diff --git a/server/internal/wsapi/codec.go b/server/internal/wsapi/codec.go index 282d5e7..40041ec 100644 --- a/server/internal/wsapi/codec.go +++ b/server/internal/wsapi/codec.go @@ -79,3 +79,9 @@ func pongMsg(clientTimeMs, serverTimeMs int64) *noituv1.ServerMessage { Pong: &noituv1.Pong{ClientTimeMs: clientTimeMs, ServerTimeMs: serverTimeMs}, }} } + +func quickMatchStatusMsg(queued bool) *noituv1.ServerMessage { + return &noituv1.ServerMessage{Payload: &noituv1.ServerMessage_QuickMatchStatus{ + QuickMatchStatus: &noituv1.QuickMatchStatus{Queued: queued}, + }} +} diff --git a/server/internal/wsapi/hub.go b/server/internal/wsapi/hub.go index 8d3d705..eb0a1a8 100644 --- a/server/internal/wsapi/hub.go +++ b/server/internal/wsapi/hub.go @@ -32,6 +32,9 @@ var ( // started shutting down, so the creator is told to come back rather than // to wait — a full room fills back up, a draining one never will. errDraining = errors.New("wsapi: server draining") + // errAlreadyQueued answers a second QuickMatch from a session already + // waiting in the pairing queue. + errAlreadyQueued = errors.New("wsapi: already queued for quick match") ) // defaultMaxRooms bounds live rooms across the whole process when nothing else @@ -59,6 +62,13 @@ type hub struct { rooms map[string]*room sessions map[string]*session // by resume token + // waiting is the FIFO of sessions queued for a quick match. No key and no + // skill: the pool this server serves is small enough that "the next + // stranger who also asked" is the whole matching policy. Guarded by mu + // rather than a lock of its own — the hub is not a bottleneck any of this + // adds meaningful contention to. + waiting []*session + joinLimiter *keyedLimiter // draining refuses every new room once set, so a creator is told to come @@ -151,6 +161,69 @@ func (h *hub) createRoom(s *session) error { return nil } +// quickMatch pairs s with the next stranger waiting, or queues it as that +// stranger for whoever asks next. +// +// A match sends both sides their QuickMatchStatus itself, before either +// input reaches the room: the room's own messages — RoomState, then +// GameStarted — are sent from its goroutine afterwards, and doing the status +// sends here first is what guarantees neither of them can arrive still +// claiming "queued". The enqueue path sends its own status for the same +// reason, symmetry, and because the caller has nobody else to hear from. +func (h *hub) quickMatch(s *session) error { + h.mu.Lock() + for _, w := range h.waiting { + if w == s { + h.mu.Unlock() + return errAlreadyQueued + } + } + if len(h.waiting) == 0 { + h.waiting = append(h.waiting, s) + h.mu.Unlock() + metrics.quickMatchQueued.Add(1) + s.send(quickMatchStatusMsg(true)) + return nil + } + waiter := h.waiting[0] + h.waiting = h.waiting[1:] + h.mu.Unlock() + + r, err := h.newRegisteredRoom(roomModePvP) + if err != nil { + // The waiter has no dispatch call site of its own to answer this + // through, being the caller of an earlier message; s is told by + // session.dispatch's own roomCreateError path instead. + waiter.send(roomCreateError(waiter.id, err)) + return err + } + + metrics.quickMatchMatched.Add(1) + waiter.send(quickMatchStatusMsg(false)) + s.send(quickMatchStatusMsg(false)) + // autoStart carries through createInput because it is a room field the + // goroutine sets for itself from handleCreate — nothing outside that + // goroutine ever touches it directly. + r.send(createInput{sess: waiter, autoStart: true}) + r.send(joinInput{sess: s}) + return nil +} + +// cancelQuickMatch drops s from the pairing queue if it is there. Idempotent: +// called on a teardown or a room-entry path that may or may not have found it +// queued, and neither is worth a special case at the call site. +func (h *hub) cancelQuickMatch(s *session) { + h.mu.Lock() + defer h.mu.Unlock() + for i, w := range h.waiting { + if w == s { + h.waiting = append(h.waiting[:i], h.waiting[i+1:]...) + metrics.quickMatchCancelled.Add(1) + return + } + } +} + // joinRoom offers a second player to a room. Whether they are seated is the // room's decision, not the hub's. func (h *hub) joinRoom(code string, s *session) error { @@ -289,6 +362,10 @@ func randomCode() string { // shutdown tells every live room to stop, so clients learn why rather than // finding the socket gone. +// +// A quick-match waiter needs nothing extra here: it registered with the hub +// the moment its Hello landed, same as any other connected session, so the +// loop below already reaches it. func (h *hub) shutdown() { h.mu.Lock() rooms := make([]*room, 0, len(h.rooms)) @@ -299,6 +376,7 @@ func (h *hub) shutdown() { for _, s := range h.sessions { sessions = append(sessions, s) } + h.waiting = nil h.mu.Unlock() for _, s := range sessions { diff --git a/server/internal/wsapi/metrics.go b/server/internal/wsapi/metrics.go index 54065a0..2c36be3 100644 --- a/server/internal/wsapi/metrics.go +++ b/server/internal/wsapi/metrics.go @@ -53,6 +53,15 @@ type metricSet struct { resumesAttempted *expvar.Int resumesSucceeded *expvar.Int + + // quickMatchQueued, quickMatchCancelled and quickMatchMatched count the + // three things that can happen to a QuickMatch: it waits, it is withdrawn + // (by CancelQuickMatch, a disconnect, or entering a room another way), or + // it is paired. Matched counts pairings, not players, so it rises by one + // per room a quick match opened. + quickMatchQueued *expvar.Int + quickMatchCancelled *expvar.Int + quickMatchMatched *expvar.Int } // metrics is the one instance every call site writes through. Built at @@ -78,5 +87,9 @@ func newMetricSet() *metricSet { botMoves: expvar.NewMap("noitu_bot_moves"), resumesAttempted: expvar.NewInt("noitu_resumes_attempted"), resumesSucceeded: expvar.NewInt("noitu_resumes_succeeded"), + + quickMatchQueued: expvar.NewInt("noitu_quick_match_queued"), + quickMatchCancelled: expvar.NewInt("noitu_quick_match_cancelled"), + quickMatchMatched: expvar.NewInt("noitu_quick_match_matched"), } } diff --git a/server/internal/wsapi/quick_match_test.go b/server/internal/wsapi/quick_match_test.go new file mode 100644 index 0000000..7ee67b9 --- /dev/null +++ b/server/internal/wsapi/quick_match_test.go @@ -0,0 +1,168 @@ +package wsapi + +import ( + "testing" + + "github.com/coder/websocket" + noituv1 "github.com/tiennm99dev/noitu/server/gen/noitu/v1" +) + +// TestQuickMatchPairsTwoWaiters is the core promise: two strangers who both +// ask are seated together and the first game begins itself, with nobody +// pressing ready or start. +func TestQuickMatchPairsTwoWaiters(t *testing.T) { + _, url := newTestServer(t, chainDict(), Config{}) + + first := dial(t, url) + first.hello("Người chờ") + first.quickMatch() + if !first.await("quick_match_status").GetQuickMatchStatus().GetQueued() { + t.Fatal("the first asker should be queued, nobody else waiting yet") + } + + second := dial(t, url) + second.hello("Người đến sau") + second.quickMatch() + + if got := second.await("quick_match_status").GetQuickMatchStatus().GetQueued(); got { + t.Error("the one who completes the pair must not be left queued") + } + if got := first.await("quick_match_status").GetQuickMatchStatus().GetQueued(); got { + t.Error("the waiter must be told the wait is over once matched") + } + + // The room begins its own first game: neither side sets ready or start. + firstStart := first.await("game_started").GetGameStarted() + secondStart := second.await("game_started").GetGameStarted() + if firstStart.GetMyTurn() == secondStart.GetMyTurn() { + t.Fatal("exactly one player should be dealt the opening turn") + } +} + +// TestQuickMatchCancelStopsTheWait covers withdrawing from the queue and +// staying out of it: a cancelled waiter must not surface later as somebody's +// match. +func TestQuickMatchCancelStopsTheWait(t *testing.T) { + _, url := newTestServer(t, chainDict(), Config{}) + + waiter := dial(t, url) + waiter.hello("Người chờ") + waiter.quickMatch() + if !waiter.await("quick_match_status").GetQuickMatchStatus().GetQueued() { + t.Fatal("expected to be queued") + } + + waiter.cancelQuickMatch() + if got := waiter.await("quick_match_status").GetQuickMatchStatus().GetQueued(); got { + t.Error("a cancel should answer queued:false") + } + + // A cancel a second time is not an error: it is idempotent. + waiter.cancelQuickMatch() + if got := waiter.await("quick_match_status").GetQuickMatchStatus().GetQueued(); got { + t.Error("cancelling twice should still answer queued:false, not an error") + } + + other := dial(t, url) + other.hello("Người khác") + other.quickMatch() + if !other.await("quick_match_status").GetQuickMatchStatus().GetQueued() { + t.Error("the cancelled waiter must not still be in line to be matched with") + } +} + +// TestQuickMatchDropsADisconnectedWaiter is the ghost case: a waiter whose +// socket ends without a CancelQuickMatch must not be handed to the next +// stranger who asks. +func TestQuickMatchDropsADisconnectedWaiter(t *testing.T) { + _, url := newTestServer(t, chainDict(), Config{}) + + ghost := dial(t, url) + ghost.hello("Ma") + ghost.quickMatch() + ghost.await("quick_match_status") + _ = ghost.conn.Close(websocket.StatusNormalClosure, "") + settle() + + a := dial(t, url) + a.hello("A") + a.quickMatch() + if !a.await("quick_match_status").GetQuickMatchStatus().GetQueued() { + t.Fatal("A should be freshly queued, not matched with a ghost") + } + + b := dial(t, url) + b.hello("B") + b.quickMatch() + + a.await("game_started") + b.await("game_started") +} + +// TestQuickMatchRefusedWhenAlreadySeated guards against a session that +// already holds a seat asking to be queued as well. +func TestQuickMatchRefusedWhenAlreadySeated(t *testing.T) { + _, url := newTestServer(t, chainDict(), Config{}) + + c := dial(t, url) + c.hello("Chủ phòng") + c.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_CreateRoom{CreateRoom: &noituv1.CreateRoom{}}}) + c.await("room_state") + + c.quickMatch() + if got := c.await("error").GetError().GetCode(); got != "already_in_a_room" { + t.Errorf("error code = %q, want already_in_a_room", got) + } +} + +// TestQuickMatchRefusesASecondRequest guards the other half of the same +// check: a session already queued must not be queued twice. +func TestQuickMatchRefusesASecondRequest(t *testing.T) { + _, url := newTestServer(t, chainDict(), Config{}) + + c := dial(t, url) + c.hello("Người chờ") + c.quickMatch() + c.await("quick_match_status") + + c.quickMatch() + if got := c.await("error").GetError().GetCode(); got != "already_queued" { + t.Errorf("error code = %q, want already_queued", got) + } +} + +// TestQuickMatchServerFullTellsBothSides is the room-cap edge: a match that +// cannot open a room must not leave either side believing it is still +// queued. +func TestQuickMatchServerFullTellsBothSides(t *testing.T) { + _, url := newTestServer(t, chainDict(), Config{MaxRooms: 1}) + + filler := dial(t, url) + filler.hello("Chiếm chỗ") + filler.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_CreateRoom{CreateRoom: &noituv1.CreateRoom{}}}) + filler.await("room_state") + + waiter := dial(t, url) + waiter.hello("Người chờ") + waiter.quickMatch() + waiter.await("quick_match_status") + + newcomer := dial(t, url) + newcomer.hello("Người đến sau") + newcomer.quickMatch() + + if got := newcomer.await("error").GetError().GetCode(); got != "server_full" { + t.Errorf("newcomer error = %q, want server_full", got) + } + if got := waiter.await("error").GetError().GetCode(); got != "server_full" { + t.Errorf("waiter error = %q, want server_full", got) + } + + // Neither is left queued: a fresh pair must not resurrect this one. + third := dial(t, url) + third.hello("Người thứ ba") + third.quickMatch() + if !third.await("quick_match_status").GetQuickMatchStatus().GetQueued() { + t.Fatal("third should be freshly queued, not immediately matched with a stale waiter") + } +} diff --git a/server/internal/wsapi/room.go b/server/internal/wsapi/room.go index 1e46573..ad22042 100644 --- a/server/internal/wsapi/room.go +++ b/server/internal/wsapi/room.go @@ -93,6 +93,10 @@ const defaultIdleWindow = 10 * time.Minute // writes happened before `go r.run()`. type createInput struct { sess *session + // autoStart marks a room a quick match opened rather than a player asking + // for a code: once both seats are filled and connected, the room begins + // its own first game instead of waiting on readiness and StartGame. + autoStart bool } type startBotInput struct { @@ -253,6 +257,12 @@ type room struct { // the room is cancelled mid-game instead of finishing normally. liveCounted atomic.Bool + // autoStart marks a room opened by a quick match. Once both seats are + // filled and connected it begins its own first game — see handleJoin — + // and is cleared right there, so every later game in the room is agreed + // with readiness and StartGame like any other. + autoStart bool + seats [maxPlayers]*seat // owner is the seat that may start a game and free the other one. It is a @@ -539,8 +549,14 @@ func (r *room) run() { func (r *room) handleCreate(m createInput) { r.seats[0] = &seat{id: "p1", nickname: m.sess.nickname(), sess: m.sess, chatFrom: r.chatSeq} r.owner = "p1" + r.autoStart = m.autoStart m.sess.attach(r, "p1") r.lobbyChanged = true + // A quick match already popped this session off the pairing queue before + // sending it here, but a plain CreateRoom might still be seating somebody + // who was also waiting in it from another attempt — one dequeue serves + // both room-entry paths. + r.hub.cancelQuickMatch(m.sess) // Deliberately sent to a brand-new room's creator, where it is always // empty: it is what replaces the conversation a client may still be // holding from a room it was in before this one. @@ -587,6 +603,7 @@ func (r *room) handleStartBot(m startBotInput) { r.seats[1] = &seat{id: botPlayerID, nickname: "Máy"} r.owner = "p1" m.sess.attach(r, "p1") + r.hub.cancelQuickMatch(m.sess) if err := r.beginGame(); err != nil { slog.Error("could not start bot game", "room", r.code, "err", err) @@ -635,7 +652,24 @@ func (r *room) handleJoin(m joinInput) { } m.sess.attach(r, string(id)) r.lobbyChanged = true + r.hub.cancelQuickMatch(m.sess) r.sendChatHistory(r.seats[free]) + + // A quick match seats both players itself rather than waiting on + // readiness and StartGame — there is no owner here to press it, only two + // strangers who both already asked to be matched. The lobby is shown + // first, with both seats filled, so the wait ends on an ordinary room a + // beat before GameStarted rather than jumping straight into one with no + // seating frame behind it. + if r.autoStart && r.seatedCount() >= minPlayers && r.allConnected() { + r.autoStart = false + r.lobbyChanged = false + r.broadcastRoomState() + if err := r.beginGame(); err != nil { + slog.Error("could not start quick-matched game", "room", r.code, "err", err) + r.broadcastError("game_start_failed") + } + } } // handleLobby applies one lobby action. diff --git a/server/internal/wsapi/session.go b/server/internal/wsapi/session.go index 38502c6..1dbc1e2 100644 --- a/server/internal/wsapi/session.go +++ b/server/internal/wsapi/session.go @@ -262,6 +262,10 @@ func (s *session) trySend(m *noituv1.ServerMessage) bool { func (s *session) run() { defer s.close() defer s.leaveRoom() + // A connection that ends while queued must not leave a ghost in line: the + // next two strangers to ask are paired with each other, not with a socket + // that is already gone. + defer s.hub.cancelQuickMatch(s) s.conn.SetReadLimit(maxFrameBytes) @@ -462,6 +466,31 @@ func (s *session) dispatch(msg *noituv1.ClientMessage) error { s.send(errorMsg("room_not_found")) } + case *noituv1.ClientMessage_QuickMatch: + if r, _ := s.currentRoom(); r != nil { + s.send(errorMsg("already_in_a_room")) + return nil + } + // A match mints a room exactly as CreateRoom does, so it is charged + // the same way and for the same reason. + if !s.roomLimiter.allow(time.Now()) { + s.send(errorMsg("too_many_rooms")) + return nil + } + if err := s.hub.quickMatch(s); err != nil { + if errors.Is(err, errAlreadyQueued) { + s.send(errorMsg("already_queued")) + } else { + s.send(roomCreateError(s.id, err)) + } + } + + case *noituv1.ClientMessage_CancelQuickMatch: + // Idempotent by design: a cancel that finds nothing queued is not an + // error, it is the answer the player wanted. + s.hub.cancelQuickMatch(s) + s.send(quickMatchStatusMsg(false)) + case *noituv1.ClientMessage_SubmitWord: s.handleSubmit(p.SubmitWord) diff --git a/server/internal/wsapi/wsapi_test.go b/server/internal/wsapi/wsapi_test.go index dd61b77..39d57d0 100644 --- a/server/internal/wsapi/wsapi_test.go +++ b/server/internal/wsapi/wsapi_test.go @@ -277,6 +277,8 @@ func payloadCase(m *noituv1.ServerMessage) string { return "chat_message" case *noituv1.ServerMessage_ChatHistory: return "chat_history" + case *noituv1.ServerMessage_QuickMatchStatus: + return "quick_match_status" } // Named rather than empty: a missing arm here makes every await for that // message time out with nothing to say about why. @@ -1151,6 +1153,18 @@ func (c *testClient) leaveRoom() { c.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_LeaveRoom{LeaveRoom: &noituv1.LeaveRoom{}}}) } +func (c *testClient) quickMatch() { + c.t.Helper() + c.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_QuickMatch{QuickMatch: &noituv1.QuickMatch{}}}) +} + +func (c *testClient) cancelQuickMatch() { + c.t.Helper() + c.send(&noituv1.ClientMessage{ + Payload: &noituv1.ClientMessage_CancelQuickMatch{CancelQuickMatch: &noituv1.CancelQuickMatch{}}, + }) +} + func (c *testClient) resign() { c.t.Helper() c.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_Resign{Resign: &noituv1.Resign{}}})