feat: validate queue compatibility metadata

This commit is contained in:
Josh Creek
2026-08-31 21:57:40 +01:00
parent 0551fb0b1f
commit b2ee9ec92d
6 changed files with 126 additions and 18 deletions
+2 -1
View File
@@ -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/
+1 -1
View File
@@ -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 |
+27 -5
View File
@@ -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 {
+77 -5
View File
@@ -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)
+17 -5
View File
@@ -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 {
+2 -1
View File
@@ -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}