From b2ee9ec92d1667dc1f813331ac85cc9656687e90 Mon Sep 17 00:00:00 2001 From: Josh Creek <8179928+jcreek@users.noreply.github.com> Date: Mon, 31 Aug 2026 21:57:40 +0100 Subject: [PATCH] feat: validate queue compatibility metadata --- multiplayer-next.md | 3 +- multiplayer-todo.md | 2 +- server/api/service.go | 32 ++++++++++++--- server/api/service_test.go | 82 +++++++++++++++++++++++++++++++++++--- server/domain/matcher.go | 22 +++++++--- server/domain/queue.go | 3 +- 6 files changed, 126 insertions(+), 18 deletions(-) diff --git a/multiplayer-next.md b/multiplayer-next.md index 435054e2..f289b2bf 100644 --- a/multiplayer-next.md +++ b/multiplayer-next.md @@ -75,7 +75,8 @@ product policy are in [`docs/MATCHMAKING.md`](docs/MATCHMAKING.md). - [ ] **IN PROGRESS:** Add one PostgreSQL-owned queue ticket/player with 10 s heartbeat, 30 s expiry, Redis candidate cache and restart/failover repair. - Authenticated owner-only recovery reads now return terminal expiry correctly; + Authenticated queue creation now requires playlist, client build and + protocol version and passes them to the server-owned candidate provider; PostgreSQL/Redis wiring remains. - [ ] **IN PROGRESS:** Validate opaque Steam ping locations and nonce-bound probes server-side; require <=100 ms, enforce discrepancy quarantine and the locked widening/ diff --git a/multiplayer-todo.md b/multiplayer-todo.md index 627cf8e4..382b2988 100644 --- a/multiplayer-todo.md +++ b/multiplayer-todo.md @@ -1191,7 +1191,7 @@ the local/CI/community transport, not a silent production fallback. | # | Task | Acceptance | |---|---|---| -| 8.14 `[D:8.4,8.5,8.8]` | **IN PROGRESS.** Pure Go queue domain enforces one active ticket per verified player under concurrent mutation, 10 s heartbeat/30 s expiry, retry-safe create/heartbeat/cancel, owner-only recovery reads and deterministic candidate projection; store layer adds a rebuildable candidate-cache boundary and authenticated HTTP queue adapter | `server/domain/queue.go`, `server/store/candidates.go` and `server/api/service.go` cover ownership/expiry/idempotency, concurrent create fencing, expired recovery as a terminal error, server-owned candidate resolution, bounded/strict JSON input and cache loss/atomic rebuild; PostgreSQL row adapter, real Redis index/TTLs and restart/failover integration remain | +| 8.14 `[D:8.4,8.5,8.8]` | **IN PROGRESS.** Pure Go queue domain enforces one active ticket per verified player under concurrent mutation, 10 s heartbeat/30 s expiry, retry-safe create/heartbeat/cancel, owner-only recovery reads and deterministic candidate projection; store layer adds a rebuildable candidate-cache boundary and authenticated HTTP queue adapter with playlist/build/protocol compatibility metadata | `server/domain/queue.go`, `server/store/candidates.go` and `server/api/service.go` cover ownership/expiry/idempotency, concurrent create fencing, expired recovery as a terminal error, server-owned candidate resolution, strict compatibility metadata and cache loss/atomic rebuild; PostgreSQL row adapter, real Redis index/TTLs and restart/failover integration remain | | 8.15 `[D:7.8,8.3]` | **IN PROGRESS.** Pure Go probe validation treats Steam location as opaque, requires nonce/freshness/region and server-computed RTT, and implements discrepancy quarantine/release; authenticated HTTP now accepts only opaque location/nonce input through a server-owned probe provider | `server/domain/probes.go`, adversarial fixtures and `server/api/service.go` cover stale/wrong/forged evidence, the 25 ms/30% threshold, three-sample quarantine, five-clean release, authenticated provider arguments and rejection of client RTT fields; Steam coordinator and regional probe adapters remain | | 8.16 `[D:8.14,8.15]` | **IN PROGRESS.** Pure Go candidate/team selection implements the <=100 ms ceiling, pairwise widening tolerance, anchor inclusion, deterministic set/region scoring and balanced team partitioning; queue-backed formation now consumes the server-owned projection and fences duplicate player identities | `server/domain/matcher.go`, `teams.go` and adversarial fixtures cover no-common-region, tolerance boundaries, lexical ties, mean-rating balance, malformed candidates, duplicate identities and queue-backed oldest-anchor formation; full population fixtures and durable matcher claim integration remain | | 8.17 `[D:8.14,8.16]` | **IN PROGRESS.** Pure Go proposal policy sends a 10-second response window to every selected human, requires unanimous acceptance, applies exact decline/timeout cooldowns and ranked escalation; authenticated API exposes revisioned accept/decline mutations; formed matches now pass through a playlist-aware proposal boundary | `server/domain/proposal.go`, `formation.go` and `server/api/service.go` plus adversarial fixtures cover partial/unanimous response, expiry, replay/conflict, stale API revision, casual lineup preparation and ranked metadata validation; queue precedence and allocation integration remain | diff --git a/server/api/service.go b/server/api/service.go index baeb0588..3ec40048 100644 --- a/server/api/service.go +++ b/server/api/service.go @@ -19,12 +19,14 @@ import ( const maxBodyBytes = 8 << 10 type CandidateProvider func(playerID, ticketID string) (domain.Candidate, error) +type CandidateProviderV2 func(playerID, ticketID string, spec domain.QueueSpec) (domain.Candidate, error) type ProbeProvider func(playerID, region string, opaqueLocation, nonce []byte, receivedAt time.Time) (domain.ProbeEvidence, []byte, error) type Service struct { Sessions *domain.SessionStore Queue *domain.Queue Candidate CandidateProvider + CandidateV2 CandidateProviderV2 Probe ProbeProvider Now func() time.Time Proposals map[string]*domain.Proposal @@ -49,7 +51,10 @@ func (s *Service) health(w http.ResponseWriter, _ *http.Request) { } type queueCreateRequest struct { - TicketID string `json:"ticket_id"` + TicketID string `json:"ticket_id"` + Playlist string `json:"playlist"` + ClientBuild string `json:"client_build"` + ProtocolVersion int `json:"protocol_version"` } type queueResponse struct { TicketID string `json:"ticket_id"` @@ -58,6 +63,7 @@ type queueResponse struct { Revision uint64 `json:"revision"` EnqueuedAt time.Time `json:"enqueued_at"` ExpiresAt time.Time `json:"expires_at"` + Playlist string `json:"playlist"` } func (s *Service) queueCreate(w http.ResponseWriter, r *http.Request) { @@ -69,7 +75,7 @@ func (s *Service) queueCreate(w http.ResponseWriter, r *http.Request) { if !ok { return } - if s.Queue == nil || s.Candidate == nil { + if s.Queue == nil || (s.Candidate == nil && s.CandidateV2 == nil) { writeError(w, http.StatusServiceUnavailable, "queue_unavailable") return } @@ -77,7 +83,7 @@ func (s *Service) queueCreate(w http.ResponseWriter, r *http.Request) { if !decodeBody(w, r, &input) { return } - if input.TicketID == "" { + if input.TicketID == "" || (input.Playlist != string(domain.Casual) && input.Playlist != string(domain.Ranked)) || input.ClientBuild == "" || len(input.ClientBuild) > 128 || input.ProtocolVersion < 1 { writeError(w, http.StatusBadRequest, "invalid_request") return } @@ -87,11 +93,27 @@ func (s *Service) queueCreate(w http.ResponseWriter, r *http.Request) { return } now := s.now() - candidate, err := s.Candidate(playerID, input.TicketID) + spec := domain.QueueSpec{Playlist: domain.Playlist(input.Playlist), ClientBuild: input.ClientBuild, ProtocolVersion: input.ProtocolVersion} + var candidate domain.Candidate + var err error + if s.CandidateV2 != nil { + candidate, err = s.CandidateV2(playerID, input.TicketID, spec) + } else { + candidate, err = s.Candidate(playerID, input.TicketID) + // Legacy providers predate queue compatibility metadata. The API has + // validated the request; keep the resulting projection self-describing. + candidate.Playlist = spec.Playlist + candidate.ClientBuild = spec.ClientBuild + candidate.ProtocolVersion = spec.ProtocolVersion + } if err != nil { writeError(w, http.StatusUnprocessableEntity, "candidate_unavailable") return } + if candidate.PlayerID != playerID || candidate.TicketID != input.TicketID || candidate.Playlist != spec.Playlist || candidate.ClientBuild != spec.ClientBuild || candidate.ProtocolVersion != spec.ProtocolVersion { + writeError(w, http.StatusUnprocessableEntity, "candidate_mismatch") + return + } ticket, err := s.Queue.Create(playerID, input.TicketID, key, candidate, now) if err != nil { writeDomainError(w, err) @@ -322,7 +344,7 @@ func decodeBody(w http.ResponseWriter, r *http.Request, target any) bool { } func toQueueResponse(ticket domain.QueueTicket) queueResponse { - return queueResponse{TicketID: ticket.TicketID, PlayerID: ticket.PlayerID, State: string(ticket.State), Revision: ticket.Revision, EnqueuedAt: ticket.EnqueuedAt, ExpiresAt: ticket.ExpiresAt} + return queueResponse{TicketID: ticket.TicketID, PlayerID: ticket.PlayerID, Playlist: string(ticket.Playlist), State: string(ticket.State), Revision: ticket.Revision, EnqueuedAt: ticket.EnqueuedAt, ExpiresAt: ticket.ExpiresAt} } func toProposalResponse(proposal domain.Proposal) proposalResponse { diff --git a/server/api/service_test.go b/server/api/service_test.go index d84fd332..68afce88 100644 --- a/server/api/service_test.go +++ b/server/api/service_test.go @@ -36,7 +36,7 @@ func TestAuthenticatedQueueAPIUsesServerCandidateAndRevisionedMutations(t *testi return response } headers := map[string]string{"Authorization": "Bearer " + session.SessionID + ":" + token, "Idempotency-Key": "create-key-123456"} - response := request(http.MethodPost, "/v1/queue", `{"ticket_id":"ticket-1"}`, headers) + response := request(http.MethodPost, "/v1/queue", `{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":1}`, headers) if response.StatusCode != http.StatusCreated { t.Fatalf("create status = %d", response.StatusCode) } @@ -64,7 +64,7 @@ func TestQueueAPIRejectsUnauthenticatedUnknownAndOversizedInput(t *testing.T) { service := &Service{Sessions: domain.NewSessionStore(), Queue: domain.NewQueue(), Candidate: func(string, string) (domain.Candidate, error) { return domain.Candidate{}, nil }} server := httptest.NewServer(service.Handler()) defer server.Close() - request, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1","player_id":"attacker"}`)) + request, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":1,"player_id":"attacker"}`)) request.Header.Set("Idempotency-Key", "create-key-123456") response, err := http.DefaultClient.Do(request) if err != nil { @@ -77,7 +77,7 @@ func TestQueueAPIRejectsUnauthenticatedUnknownAndOversizedInput(t *testing.T) { sessionStore := domain.NewSessionStore() session, token, _ := sessionStore.Issue("player-1", time.Hour, time.Now()) service.Sessions = sessionStore - request, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1","unknown":true}`)) + request, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":1,"unknown":true}`)) request.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token) request.Header.Set("Idempotency-Key", "create-key-123456") response, err = http.DefaultClient.Do(request) @@ -99,7 +99,7 @@ func TestQueueAPIRejectsUnauthenticatedUnknownAndOversizedInput(t *testing.T) { t.Fatalf("malformed body status = %d", response.StatusCode) } _ = response.Body.Close() - request, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1"}{"ticket_id":"ticket-2"}`)) + request, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":1}{"ticket_id":"ticket-2"}`)) request.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token) request.Header.Set("Idempotency-Key", "create-key-789012") response, err = http.DefaultClient.Do(request) @@ -112,6 +112,78 @@ func TestQueueAPIRejectsUnauthenticatedUnknownAndOversizedInput(t *testing.T) { _ = response.Body.Close() } +func TestQueueCreateRequiresCompatibilityMetadataAndPassesItToProvider(t *testing.T) { + now := time.Unix(1000, 0).UTC() + sessions := domain.NewSessionStore() + session, token, err := sessions.Issue("player-1", time.Hour, now) + if err != nil { + t.Fatal(err) + } + var got domain.QueueSpec + service := &Service{ + Sessions: sessions, + Queue: domain.NewQueue(), + Now: func() time.Time { return now }, + CandidateV2: func(_ string, ticketID string, spec domain.QueueSpec) (domain.Candidate, error) { + got = spec + return domain.Candidate{PlayerID: "player-1", TicketID: ticketID, Playlist: spec.Playlist, ClientBuild: spec.ClientBuild, ProtocolVersion: spec.ProtocolVersion, EnqueuedAt: now}, nil + }, + } + server := httptest.NewServer(service.Handler()) + defer server.Close() + request := func(body string) *http.Response { + req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(body)) + req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token) + req.Header.Set("Idempotency-Key", "create-key-123456") + response, requestErr := http.DefaultClient.Do(req) + if requestErr != nil { + t.Fatal(requestErr) + } + return response + } + response := request(`{"ticket_id":"ticket-1"}`) + if response.StatusCode != http.StatusBadRequest { + t.Fatalf("missing metadata status = %d", response.StatusCode) + } + _ = response.Body.Close() + response = request(`{"ticket_id":"ticket-1","playlist":"invalid","client_build":"build-1","protocol_version":1}`) + if response.StatusCode != http.StatusBadRequest { + t.Fatalf("invalid playlist status = %d", response.StatusCode) + } + _ = response.Body.Close() + response = request(`{"ticket_id":"ticket-1","playlist":"ranked","client_build":"build-1","protocol_version":7}`) + if response.StatusCode != http.StatusCreated { + t.Fatalf("valid metadata status = %d", response.StatusCode) + } + _ = response.Body.Close() + if got.Playlist != domain.Ranked || got.ClientBuild != "build-1" || got.ProtocolVersion != 7 { + t.Fatalf("provider received %+v", got) + } +} + +func TestQueueCreateRejectsCandidateMetadataMismatch(t *testing.T) { + now := time.Unix(1000, 0).UTC() + sessions := domain.NewSessionStore() + session, token, _ := sessions.Issue("player-1", time.Hour, now) + service := &Service{Sessions: sessions, Queue: domain.NewQueue(), Now: func() time.Time { return now }, CandidateV2: func(_ string, ticketID string, spec domain.QueueSpec) (domain.Candidate, error) { + spec.ClientBuild = "tampered" + return domain.Candidate{PlayerID: "player-1", TicketID: ticketID, Playlist: spec.Playlist, ClientBuild: spec.ClientBuild, ProtocolVersion: spec.ProtocolVersion, EnqueuedAt: now}, nil + }} + server := httptest.NewServer(service.Handler()) + defer server.Close() + req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1","playlist":"ranked","client_build":"build-1","protocol_version":1}`)) + req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token) + req.Header.Set("Idempotency-Key", "create-key-123456") + response, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + defer response.Body.Close() + if response.StatusCode != http.StatusUnprocessableEntity { + t.Fatalf("mismatch status = %d", response.StatusCode) + } +} + func TestQueueRecoveryAPIIsAuthenticatedOwnerOnlyAndExpiresStaleTickets(t *testing.T) { now := time.Unix(1000, 0).UTC() sessions := domain.NewSessionStore() @@ -129,7 +201,7 @@ func TestQueueRecoveryAPIIsAuthenticatedOwnerOnlyAndExpiresStaleTickets(t *testi }} server := httptest.NewServer(service.Handler()) defer server.Close() - create, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-recovery-123456"}`)) + create, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":"ticket-recovery-123456","playlist":"casual","client_build":"build-1","protocol_version":1}`)) create.Header.Set("Authorization", "Bearer "+ownerSession.SessionID+":"+ownerToken) create.Header.Set("Idempotency-Key", "queue-create-recovery-123456") response, err := http.DefaultClient.Do(create) diff --git a/server/domain/matcher.go b/server/domain/matcher.go index 571a82eb..40574138 100644 --- a/server/domain/matcher.go +++ b/server/domain/matcher.go @@ -15,14 +15,26 @@ const ( RatingWidenPeriod = 30.0 ) +// QueueSpec is the compatibility contract selected by the authenticated +// client. CandidateProviderV2 may use it to resolve a server-owned projection +// from the verified account and current deployment configuration. +type QueueSpec struct { + Playlist Playlist + ClientBuild string + ProtocolVersion int +} + // Candidate is the server-side projection of a verified, live queue ticket. // RTT values come from backend probes, never from the client request body. type Candidate struct { - TicketID string - PlayerID string - Rating float64 - EnqueuedAt time.Time - PredictedRTT map[string]float64 + TicketID string + PlayerID string + Playlist Playlist + ClientBuild string + ProtocolVersion int + Rating float64 + EnqueuedAt time.Time + PredictedRTT map[string]float64 } type Selection struct { diff --git a/server/domain/queue.go b/server/domain/queue.go index 9e0f87a1..71b97c64 100644 --- a/server/domain/queue.go +++ b/server/domain/queue.go @@ -26,6 +26,7 @@ type QueueTicket struct { TicketID string PlayerID string Candidate Candidate + Playlist Playlist State State Revision uint64 EnqueuedAt time.Time @@ -70,7 +71,7 @@ func (q *Queue) Create(playerID, ticketID, idempotencyKey string, candidate Cand if _, ok := q.tickets[ticketID]; ok { return QueueTicket{}, fmt.Errorf("%w: ticket ID already exists", ErrConflict) } - ticket := QueueTicket{TicketID: ticketID, PlayerID: playerID, Candidate: candidate, State: Queued, EnqueuedAt: now, ExpiresAt: now.Add(QueueExpiryWindow)} + ticket := QueueTicket{TicketID: ticketID, PlayerID: playerID, Candidate: candidate, Playlist: candidate.Playlist, State: Queued, EnqueuedAt: now, ExpiresAt: now.Add(QueueExpiryWindow)} q.tickets[ticketID] = ticket q.byPlayer[playerID] = ticketID q.mutations[idempotencyKey] = queueMutation{digest: digest, ticket: ticket}