Files
CosmicClash/server/api/service_test.go
T
Josh Creek 801fca7cb0 fix(matchmaking): make regional RTT evidence obtainable end to end
domain.validCandidate hard-requires a non-empty PredictedRTT map, but
CreateQueueTicket persisted an empty one and the only endpoint that
could fill it returned 503 in every real binary, because Service.Probe
was assigned nowhere outside api tests. No client-created ticket could
ever be selected by the matcher. The Godot client had no probe method at
all, so even a wired backend was unreachable from the game.

Four distinct defects had to be fixed for this path to work:

Nothing issued the nonce ProbeProvider was meant to compare against, so
the contract could not be satisfied even in principle. Add
POST /v1/probes/{region}/challenge, backed by a durable single-use
challenge -- durable because any replica may serve the answer for a
challenge another replica issued. RTT is the interval between issuing
and receiving, so no client-reported latency reaches placement.

CreateQueueTicket marshalled a nil map to JSON `null`, a JSONB scalar
rather than an object, and jsonb_set rejects that with "cannot set path
in scalar". RecordProbe would have failed at runtime even once wired.
Persist an object, and normalise non-object values in the update for
rows already written.

A nil ProbeRecorder made the handler report success while persisting
nothing, which silently leaves the ticket unmatchable. That is a
misconfiguration, not a successful probe; it now returns 503.

A successful probe updated PostgreSQL only. The candidate inserted at
enqueue time carries an empty RTT map, and the Redis keyspace has its
TTL continually refreshed, so the stale entry need never repair itself.
Refresh that player's projection after the probe commits.

Client side: add the challenge/answer round trip and have the
matchmaking screen collect evidence before creating a ticket, since
queueing first produces a search that can never match. Probing every
region fully is not required -- placement uses whichever regions
answered -- but queueing with none is refused rather than silently
stalling.

New integration test drives the real enqueue and probe paths and then
asks the actual matcher predicate, rather than hand-building a candidate
the way the unit tests do -- which is exactly why they missed this.

Also make the integration schema reset drop the whole public schema: the
enumerated table list silently broke with each new migration.
2026-09-05 10:49:28 +01:00

1958 lines
83 KiB
Go

package api
import (
"bufio"
"bytes"
"context"
"encoding/binary"
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/cosmic-clash/cosmic-clash/server/domain"
"github.com/cosmic-clash/cosmic-clash/server/observability"
)
type queueBackendSpy struct{ createCalls, heartbeatCalls, cancelCalls, getCalls int }
func TestQueueResponseCarriesRecoveredMatchIdentity(t *testing.T) {
response := toQueueResponse(domain.QueueTicket{TicketID: "ticket-1234567890", PlayerID: "player-1234567890", ProposalID: "proposal-1234567890", MatchID: "match-1234567890", State: domain.AssignmentReady})
if response.ProposalID != "proposal-1234567890" {
t.Fatalf("queue response proposal ID = %q", response.ProposalID)
}
if response.MatchID != "match-1234567890" {
t.Fatalf("queue response match ID = %q", response.MatchID)
}
}
type candidateIndexSpy struct {
upsertCalls, removeCalls int
upsertErr, removeErr error
last domain.Candidate
}
type probeRecorderSpy struct {
calls int
err error
last struct {
player, region string
rtt time.Duration
}
}
type resultSubmitterSpy struct {
calls int
err error
key string
result domain.MatchResult
}
type serverRegistrarSpy struct {
calls int
binding domain.WorkloadBinding
protocol int
assignmentReady bool
err error
}
type serverShutdownerSpy struct {
calls int
binding domain.WorkloadBinding
reason string
key string
err error
}
type serverConnectionSpy struct {
connectCalls int
disconnectCalls int
binding domain.WorkloadBinding
playerID string
key string
expectedGeneration uint64
generation uint64
err error
}
func (s *serverConnectionSpy) ClaimPlayerConnection(_ context.Context, binding domain.WorkloadBinding, playerID string, expectedGeneration uint64, key string, _ time.Time) (uint64, error) {
s.connectCalls++
s.binding, s.playerID, s.expectedGeneration, s.key = binding, playerID, expectedGeneration, key
return expectedGeneration + 1, s.err
}
func (s *serverConnectionSpy) RecordPlayerDisconnected(_ context.Context, binding domain.WorkloadBinding, playerID string, generation uint64, key string, _ time.Time) error {
s.disconnectCalls++
s.binding, s.playerID, s.generation, s.key = binding, playerID, generation, key
return s.err
}
func (s *serverShutdownerSpy) ShutdownServer(_ context.Context, binding domain.WorkloadBinding, reason, key string, _ time.Time) error {
s.calls++
s.binding, s.reason, s.key = binding, reason, key
return s.err
}
func (s *serverRegistrarSpy) RegisterServer(_ context.Context, binding domain.WorkloadBinding, protocol int, assignmentReady bool, _ string, _ time.Time) error {
s.calls++
s.binding, s.protocol, s.assignmentReady = binding, protocol, assignmentReady
return s.err
}
type proposalPromoterSpy struct {
calls int
proposal domain.Proposal
err error
}
func (p *proposalPromoterSpy) Promote(_ context.Context, proposal domain.Proposal, _ time.Time) error {
p.calls++
p.proposal = proposal
return p.err
}
func (r *resultSubmitterSpy) SubmitResult(_ context.Context, key string, result domain.MatchResult, _ domain.WorkloadBinding, _ []byte, _ time.Time) error {
r.calls++
r.key, r.result = key, result
return r.err
}
func (p *probeRecorderSpy) RecordProbe(_ context.Context, player, region string, rtt time.Duration, _ time.Time) error {
p.calls++
p.last.player, p.last.region, p.last.rtt = player, region, rtt
return p.err
}
func (i *candidateIndexSpy) Upsert(_ context.Context, candidate domain.Candidate) error {
i.upsertCalls++
i.last = candidate
return i.upsertErr
}
func (i *candidateIndexSpy) Remove(_ context.Context, _ domain.Playlist, _ string) error {
i.removeCalls++
return i.removeErr
}
type sessionBackendSpy struct{ calls int }
type proposalBackendSpy struct {
proposal domain.Proposal
calls int
mutations int
}
func (b *proposalBackendSpy) Respond(_ context.Context, playerID, _ string, key string, accept bool, revision uint64, now time.Time) (domain.Proposal, error) {
b.mutations++
updated, err := b.proposal.Respond(playerID, key, accept, revision, now)
if err == nil {
b.proposal = updated
}
return updated, err
}
func (b *proposalBackendSpy) Get(_ context.Context, playerID, _ string, _ time.Time) (domain.Proposal, error) {
b.calls++
if !b.proposal.HasParticipant(playerID) {
return domain.Proposal{}, domain.ErrNotParticipant
}
return b.proposal, nil
}
func (s *sessionBackendSpy) Authenticate(_ context.Context, sessionID, _ string, _ time.Time) (domain.Session, error) {
s.calls++
return domain.Session{SessionID: sessionID, PlayerID: "player-1"}, nil
}
type steamLoginSpy struct{ calls int }
func (s *steamLoginSpy) Authenticate(_ context.Context, ticket string, _ time.Time) (domain.VerifiedIdentity, error) {
s.calls++
if ticket != "valid-web-ticket" {
return domain.VerifiedIdentity{}, domain.ErrTicketRejected
}
return domain.VerifiedIdentity{PlayerID: "player-1", SteamID: "steam-1"}, nil
}
func (b *queueBackendSpy) Create(_ context.Context, playerID, ticketID, _ string, spec domain.QueueSpec, now time.Time) (domain.QueueTicket, error) {
b.createCalls++
return domain.QueueTicket{TicketID: ticketID, PlayerID: playerID, Playlist: spec.Playlist, State: domain.Queued, EnqueuedAt: now, ExpiresAt: now.Add(domain.QueueExpiryWindow)}, nil
}
func (b *queueBackendSpy) Heartbeat(_ context.Context, playerID, ticketID, _ string, revision uint64, now time.Time) (domain.QueueTicket, error) {
b.heartbeatCalls++
return domain.QueueTicket{TicketID: ticketID, PlayerID: playerID, State: domain.Queued, Revision: revision + 1, EnqueuedAt: now, ExpiresAt: now.Add(domain.QueueExpiryWindow)}, nil
}
func (b *queueBackendSpy) Cancel(_ context.Context, playerID, ticketID, _ string, revision uint64, now time.Time) (domain.QueueTicket, error) {
b.cancelCalls++
return domain.QueueTicket{TicketID: ticketID, PlayerID: playerID, State: domain.Cancelled, Revision: revision + 1, EnqueuedAt: now, ExpiresAt: now}, nil
}
func (b *queueBackendSpy) Get(_ context.Context, playerID, ticketID string, now time.Time) (domain.QueueTicket, error) {
b.getCalls++
return domain.QueueTicket{TicketID: ticketID, PlayerID: playerID, State: domain.Queued, EnqueuedAt: now, ExpiresAt: now.Add(domain.QueueExpiryWindow)}, nil
}
func TestAuthenticatedQueueAPIUsesServerCandidateAndRevisionedMutations(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)
}
queue := domain.NewQueue()
service := &Service{Sessions: sessions, Queue: queue, Now: func() time.Time { return now }, Candidate: func(playerID, ticketID string) (domain.Candidate, error) {
return domain.Candidate{PlayerID: playerID, TicketID: ticketID, EnqueuedAt: now, PredictedRTT: map[string]float64{"EU": 20}}, nil
}}
server := httptest.NewServer(service.Handler())
defer server.Close()
request := func(method, path, body string, headers map[string]string) *http.Response {
req, _ := http.NewRequest(method, server.URL+path, strings.NewReader(body))
for key, value := range headers {
req.Header.Set(key, value)
}
response, requestErr := http.DefaultClient.Do(req)
if requestErr != nil {
t.Fatal(requestErr)
}
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","playlist":"casual","client_build":"build-1","protocol_version":1}`, headers)
if response.StatusCode != http.StatusCreated {
t.Fatalf("create status = %d", response.StatusCode)
}
var created queueResponse
if err := json.NewDecoder(response.Body).Decode(&created); err != nil {
t.Fatal(err)
}
_ = response.Body.Close()
if created.PlayerID != "player-1" || created.State != "QUEUED" || created.Revision != 0 {
t.Fatalf("created = %+v", created)
}
response = request(http.MethodPost, "/v1/queue/ticket-1/heartbeat", `{}`, map[string]string{"Authorization": headers["Authorization"], "Idempotency-Key": "heartbeat-key-123456", "If-Match-Revision": "0"})
if response.StatusCode != http.StatusOK {
t.Fatalf("heartbeat status = %d", response.StatusCode)
}
_ = response.Body.Close()
response = request(http.MethodPost, "/v1/queue/ticket-1/cancel", `{}`, map[string]string{"Authorization": headers["Authorization"], "Idempotency-Key": "cancel-key-123456", "If-Match-Revision": "0"})
if response.StatusCode != http.StatusConflict {
t.Fatalf("stale cancel status = %d", response.StatusCode)
}
_ = response.Body.Close()
}
func TestDocumentedContractRoutesAdaptToServiceAPI(t *testing.T) {
now := time.Unix(1000, 0).UTC()
backend := &queueBackendSpy{}
service := &Service{
SessionBackend: &sessionBackendSpy{},
QueueBackend: backend,
Now: func() time.Time { return now },
}
server := httptest.NewServer(service.Handler())
defer server.Close()
auth := "Bearer session-1:token-1"
create, err := http.NewRequest(http.MethodPost, server.URL+"/api/v1/queue/tickets", strings.NewReader(`{"playlist":"casual","client_build":"build-1","protocol_version":1}`))
if err != nil {
t.Fatal(err)
}
create.Header.Set("Authorization", auth)
create.Header.Set("Idempotency-Key", "contract-create-key-123456")
response, err := http.DefaultClient.Do(create)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusCreated || backend.createCalls != 1 {
t.Fatalf("create status = %d, calls = %d", response.StatusCode, backend.createCalls)
}
var ticket queueResponse
if err := json.NewDecoder(response.Body).Decode(&ticket); err != nil {
t.Fatal(err)
}
_ = response.Body.Close()
if ticket.TicketID == "" {
t.Fatal("contract adapter did not assign a ticket id")
}
heartbeat, err := http.NewRequest(http.MethodPost, server.URL+"/api/v1/queue/tickets/"+ticket.TicketID+"/heartbeat", nil)
if err != nil {
t.Fatal(err)
}
heartbeat.Header.Set("Authorization", auth)
heartbeat.Header.Set("Idempotency-Key", "contract-heartbeat-key-123")
heartbeat.Header.Set("If-Match-Revision", "0")
response, err = http.DefaultClient.Do(heartbeat)
if err != nil {
t.Fatal(err)
}
_ = response.Body.Close()
if response.StatusCode != http.StatusOK || backend.heartbeatCalls != 1 {
t.Fatalf("heartbeat status = %d, calls = %d", response.StatusCode, backend.heartbeatCalls)
}
cancel, err := http.NewRequest(http.MethodDelete, server.URL+"/api/v1/queue/tickets/"+ticket.TicketID, nil)
if err != nil {
t.Fatal(err)
}
cancel.Header.Set("Authorization", auth)
cancel.Header.Set("Idempotency-Key", "contract-cancel-key-123456")
cancel.Header.Set("If-Match-Revision", "0")
response, err = http.DefaultClient.Do(cancel)
if err != nil {
t.Fatal(err)
}
_ = response.Body.Close()
if response.StatusCode != http.StatusNoContent || backend.cancelCalls != 1 {
t.Fatalf("cancel status = %d, calls = %d", response.StatusCode, backend.cancelCalls)
}
}
func TestDocumentedContractRoutesRejectNonOpaqueResourceIDs(t *testing.T) {
service := &Service{}
server := httptest.NewServer(service.Handler())
defer server.Close()
request, err := http.NewRequest(http.MethodPost, server.URL+"/api/v1/queue/tickets", strings.NewReader(`{"ticket_id":"short","playlist":"casual","client_build":"build-1","protocol_version":1}`))
if err != nil {
t.Fatal(err)
}
if response, requestErr := http.DefaultClient.Do(request); requestErr != nil {
t.Fatal(requestErr)
} else {
if response.StatusCode != http.StatusBadRequest {
t.Fatalf("short supplied ticket id status = %d, want 400", response.StatusCode)
}
response.Body.Close()
}
paths := []string{
"/api/v1/queue/tickets/short/heartbeat",
"/api/v1/proposals/proposal/unsafe/accept",
"/api/v1/assignments/match/unsafe",
"/api/v1/servers/server/unsafe/result",
}
for _, path := range paths {
request, err := http.NewRequest(http.MethodGet, server.URL+path, nil)
if err != nil {
t.Fatal(err)
}
response, err := http.DefaultClient.Do(request)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusNotFound {
t.Fatalf("%s status = %d, want 404", path, response.StatusCode)
}
response.Body.Close()
}
}
func TestAuthenticatedWebSocketDeliversOnlyTargetedRevisionedEvents(t *testing.T) {
service := &Service{SessionBackend: &sessionBackendSpy{}}
server := httptest.NewServer(service.Handler())
defer server.Close()
invalid, err := http.NewRequest(http.MethodGet, server.URL+"/v1/events", nil)
if err != nil {
t.Fatal(err)
}
invalid.Header.Set("Upgrade", "websocket")
invalid.Header.Set("Connection", "Upgrade")
invalid.Header.Set("Sec-WebSocket-Key", "not-a-websocket-key")
invalid.Header.Set("Authorization", "Bearer session-1:token-1")
invalidResponse, err := server.Client().Do(invalid)
if err != nil {
t.Fatal(err)
}
_ = invalidResponse.Body.Close()
if invalidResponse.StatusCode != http.StatusBadRequest {
t.Fatalf("invalid handshake status = %d", invalidResponse.StatusCode)
}
connection, err := net.Dial("tcp", strings.TrimPrefix(server.URL, "http://"))
if err != nil {
t.Fatal(err)
}
defer connection.Close()
_, err = io.WriteString(connection, "GET /v1/events HTTP/1.1\r\nHost: localhost\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nAuthorization: Bearer session-1:token-1\r\n\r\n")
if err != nil {
t.Fatal(err)
}
reader := bufio.NewReader(connection)
status, err := reader.ReadString('\n')
if err != nil {
t.Fatal(err)
}
if !strings.Contains(status, "101 Switching Protocols") {
t.Fatalf("handshake status = %q", status)
}
for {
line, err := reader.ReadString('\n')
if err != nil {
t.Fatal(err)
}
if line == "\r\n" {
break
}
}
time.Sleep(10 * time.Millisecond)
if err := service.PublishControlPlaneEvent(ControlPlaneEvent{Event: "state_changed", Revision: 1, ResourceID: "ticket-1234567890123456", OccurredAt: time.Unix(1000, 0).UTC(), State: "QUEUED", PlayerID: "player-1"}); err != nil {
t.Fatal(err)
}
if err := service.PublishControlPlaneEvent(ControlPlaneEvent{Event: "state_changed", Revision: 2, ResourceID: "ticket-1234567890123456", OccurredAt: time.Unix(1001, 0).UTC(), State: "PROPOSED", PlayerID: "player-2"}); err != nil {
t.Fatal(err)
}
first, err := readServerWebSocketFrame(reader)
if err != nil {
t.Fatal(err)
}
var event ControlPlaneEvent
if err := json.Unmarshal(first, &event); err != nil {
t.Fatal(err)
}
if event.PlayerID != "" || event.Revision != 1 || event.ResourceID != "ticket-1234567890123456" || event.State != "QUEUED" {
t.Fatalf("event = %+v", event)
}
}
func TestServerRosterRequiresWorkloadBindingAndReturnsRawSignedEnvelopes(t *testing.T) {
now := time.Unix(1000, 0).UTC()
service := &Service{
Now: func() time.Time { return now },
WorkloadVerify: func(token string, at time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" || !at.Equal(now) {
t.Fatal("unexpected workload verification input")
}
return domain.WorkloadBinding{ServerID: "server-1", MatchID: "match-1", AllocationID: "allocation-1"}, nil
},
Roster: func(_ context.Context, binding domain.WorkloadBinding, at time.Time) ([][]byte, error) {
if binding.ServerID != "server-1" || binding.MatchID != "match-1" || !at.Equal(now) {
t.Fatal("unexpected roster binding")
}
return [][]byte{[]byte(`{"authorisation":{"player_id":"player-1","expires_at":"1970-01-01T00:33:20Z"},"signature":"sig"}`)}, nil
},
}
server := httptest.NewServer(service.Handler())
defer server.Close()
request, _ := http.NewRequest(http.MethodGet, server.URL+"/v1/servers/server-1/roster", nil)
request.Header.Set("Authorization", "Bearer workload-token")
response, err := server.Client().Do(request)
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
if response.StatusCode != http.StatusOK {
t.Fatalf("roster status=%d", response.StatusCode)
}
var roster []json.RawMessage
if err := json.NewDecoder(response.Body).Decode(&roster); err != nil {
t.Fatal(err)
}
if len(roster) != 1 || !bytes.Contains(roster[0], []byte(`"player_id":"player-1"`)) {
t.Fatalf("roster=%s", roster[0])
}
}
func readServerWebSocketFrame(reader *bufio.Reader) ([]byte, error) {
first, err := reader.ReadByte()
if err != nil {
return nil, err
}
second, err := reader.ReadByte()
if err != nil {
return nil, err
}
if first&0x0f != 0x1 || second&0x80 != 0 {
return nil, errors.New("unexpected server websocket frame")
}
length := int(second & 0x7f)
if length == 126 {
var extended uint16
if err := binary.Read(reader, binary.BigEndian, &extended); err != nil {
return nil, err
}
length = int(extended)
}
payload := make([]byte, length)
_, err = io.ReadFull(reader, payload)
return payload, err
}
func TestEventHubClosesSlowSubscribersExactlyOnce(t *testing.T) {
hub := newEventHub()
subscriber := hub.subscribe("player-1")
event := ControlPlaneEvent{Event: "state_changed", Revision: 1, ResourceID: "ticket-1234567890123456", OccurredAt: time.Unix(1000, 0).UTC(), State: "QUEUED", PlayerID: "player-1"}
for i := 0; i < eventQueueCapacity; i++ {
if err := hub.publish(event); err != nil {
t.Fatal(err)
}
}
if err := hub.publish(event); err != nil {
t.Fatal(err)
}
for {
_, open := <-subscriber.queue
if !open {
break
}
}
hub.unsubscribe(subscriber)
}
func TestEventHubRejectsEventsOutsideTheV1Vocabulary(t *testing.T) {
hub := newEventHub()
base := ControlPlaneEvent{Revision: 1, ResourceID: "ticket-1234567890123456", OccurredAt: time.Unix(1000, 0).UTC(), PlayerID: "player-1"}
invalid := []ControlPlaneEvent{
{Event: "unknown", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "state_changed", State: "QUEUED", Revision: base.Revision, ResourceID: "short", OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "state_changed", State: "QUEUED", Revision: base.Revision, ResourceID: "ticket-1234567890/unsafe", OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "state_changed", State: "NOT_A_STATE", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "proposal_changed", State: "LIVE", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "assignment_changed", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "assignment_changed", MatchID: "match-1", ServerID: "server-1", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "assignment_changed", MatchID: "match_1234567890", ServerID: "server/unsafe", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
{Event: "error", Code: "SECRET_LEAK", Revision: base.Revision, ResourceID: base.ResourceID, OccurredAt: base.OccurredAt, PlayerID: base.PlayerID},
}
for _, event := range invalid {
if err := hub.publish(event); err == nil {
t.Fatalf("invalid event was accepted: %+v", event)
}
}
}
func TestStateChangingAPIActionsPublishTargetedEvents(t *testing.T) {
now := time.Unix(1000, 0).UTC()
backend := &queueBackendSpy{}
service := &Service{SessionBackend: &sessionBackendSpy{}, QueueBackend: backend, Now: func() time.Time { return now }, Proposals: make(map[string]*domain.Proposal)}
subscriber := service.getEventHub().subscribe("player-1")
defer service.getEventHub().unsubscribe(subscriber)
create := httptest.NewRequest(http.MethodPost, "/v1/queue", strings.NewReader(`{"ticket_id":"ticket-1234567890123456","playlist":"casual","client_build":"build-1","protocol_version":1}`))
create.Header.Set("Authorization", "Bearer session-1:token-1")
create.Header.Set("Idempotency-Key", "create-event-key-123456")
createRecorder := httptest.NewRecorder()
service.queueCreate(createRecorder, create)
if createRecorder.Code != http.StatusCreated {
t.Fatalf("create status = %d", createRecorder.Code)
}
var queueEvent ControlPlaneEvent
if err := json.Unmarshal(<-subscriber.queue, &queueEvent); err != nil {
t.Fatal(err)
}
if queueEvent.Event != "state_changed" || queueEvent.ResourceID != "ticket-1234567890123456" || queueEvent.PlayerID != "" {
t.Fatalf("queue event = %+v", queueEvent)
}
proposal, err := domain.NewProposal("proposal-1234567890123456", domain.Casual, []string{"player-1", "player-2"}, now)
if err != nil {
t.Fatal(err)
}
service.Proposals[proposal.ProposalID] = &proposal
respond := httptest.NewRequest(http.MethodPost, "/v1/proposals/"+proposal.ProposalID+"/accept", nil)
respond.Header.Set("Authorization", "Bearer session-1:token-1")
respond.Header.Set("Idempotency-Key", "proposal-event-key-123456")
respond.Header.Set("If-Match-Revision", "0")
respondRecorder := httptest.NewRecorder()
service.proposalMutation(respondRecorder, respond)
if respondRecorder.Code != http.StatusOK {
t.Fatalf("proposal status = %d", respondRecorder.Code)
}
var proposalEvent ControlPlaneEvent
if err := json.Unmarshal(<-subscriber.queue, &proposalEvent); err != nil {
t.Fatal(err)
}
if proposalEvent.Event != "proposal_changed" || proposalEvent.ResourceID != proposal.ProposalID || proposalEvent.State != "OPEN" {
t.Fatalf("proposal event = %+v", proposalEvent)
}
}
func TestFinalProposalAcceptancePromotesDurableMatchAndFailsRetryably(t *testing.T) {
now := time.Unix(1000, 0).UTC()
proposal, err := domain.NewProposal("proposal-promote-123456", domain.Casual, []string{"player-1", "player-2"}, now)
if err != nil {
t.Fatal(err)
}
backend := &proposalBackendSpy{proposal: proposal}
promoter := &proposalPromoterSpy{}
sessions := domain.NewSessionStore()
session1, token1, err := sessions.Issue("player-1", time.Hour, now)
if err != nil {
t.Fatal(err)
}
session2, token2, err := sessions.Issue("player-2", time.Hour, now)
if err != nil {
t.Fatal(err)
}
service := &Service{Sessions: sessions, ProposalBackend: backend, ProposalPromoter: promoter, Now: func() time.Time { return now }}
respond := func(credential, key, revision string) int {
req := httptest.NewRequest(http.MethodPost, "/v1/proposals/"+proposal.ProposalID+"/accept", nil)
req.Header.Set("Authorization", "Bearer "+credential)
req.Header.Set("Idempotency-Key", key)
req.Header.Set("If-Match-Revision", revision)
recorder := httptest.NewRecorder()
service.proposalMutation(recorder, req)
return recorder.Code
}
credential1 := session1.SessionID + ":" + token1
credential2 := session2.SessionID + ":" + token2
if status := respond(credential1, "proposal-promote-first", "0"); status != http.StatusOK || promoter.calls != 0 {
t.Fatalf("first acceptance status/calls = %d/%d", status, promoter.calls)
}
if status := respond(credential2, "proposal-promote-final", "1"); status != http.StatusOK || promoter.calls != 1 || promoter.proposal.State != domain.Accepted {
t.Fatalf("final acceptance status/promoter = %d/%+v", status, promoter)
}
promoter.err = errors.New("database unavailable")
// A duplicate response is replayed by the durable proposal backend and
// retries promotion instead of asking the player to accept again.
if status := respond(credential2, "proposal-promote-final", "1"); status != http.StatusServiceUnavailable || promoter.calls != 2 {
t.Fatalf("promotion retry status/calls = %d/%d", status, promoter.calls)
}
}
func TestProposalRecoveryUsesDurableBackendAndRemainsParticipantScoped(t *testing.T) {
now := time.Unix(1000, 0).UTC()
proposal, err := domain.NewProposal("proposal-1234567890123456", domain.Casual, []string{"player-1", "player-2"}, now)
if err != nil {
t.Fatal(err)
}
backend := &proposalBackendSpy{proposal: proposal}
sessions := domain.NewSessionStore()
participantSession, participantToken, err := sessions.Issue("player-1", time.Hour, now)
if err != nil {
t.Fatal(err)
}
outsiderSession, outsiderToken, err := sessions.Issue("outsider", time.Hour, now)
if err != nil {
t.Fatal(err)
}
service := &Service{Sessions: sessions, Proposals: map[string]*domain.Proposal{}, ProposalBackend: backend, Now: func() time.Time { return now }}
server := httptest.NewServer(service.Handler())
defer server.Close()
get := func(session domain.Session, token string) int {
request, err := http.NewRequest(http.MethodGet, server.URL+"/v1/proposals/"+proposal.ProposalID, nil)
if err != nil {
t.Fatal(err)
}
request.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, err := server.Client().Do(request)
if err != nil {
t.Fatal(err)
}
_ = response.Body.Close()
return response.StatusCode
}
if status := get(participantSession, participantToken); status != http.StatusOK {
t.Fatalf("participant recovery status = %d", status)
}
respond, err := http.NewRequest(http.MethodPost, server.URL+"/v1/proposals/"+proposal.ProposalID+"/accept", nil)
if err != nil {
t.Fatal(err)
}
respond.Header.Set("Authorization", "Bearer "+participantSession.SessionID+":"+participantToken)
respond.Header.Set("Idempotency-Key", "proposal-durable-response-123456")
respond.Header.Set("If-Match-Revision", "0")
response, err := server.Client().Do(respond)
if err != nil {
t.Fatal(err)
}
_ = response.Body.Close()
if response.StatusCode != http.StatusOK || backend.mutations != 1 {
t.Fatalf("durable response status = %d, mutations = %d", response.StatusCode, backend.mutations)
}
if status := get(outsiderSession, outsiderToken); status != http.StatusNotFound {
t.Fatalf("outsider recovery status = %d", status)
}
if backend.calls != 2 {
t.Fatalf("durable backend calls = %d", backend.calls)
}
}
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","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 {
t.Fatal(err)
}
if response.StatusCode != http.StatusUnauthorized {
t.Fatalf("unauthenticated status = %d", response.StatusCode)
}
_ = response.Body.Close()
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","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)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusBadRequest {
t.Fatalf("unknown field status = %d", response.StatusCode)
}
_ = response.Body.Close()
request, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/queue", strings.NewReader(`{"ticket_id":`))
request.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
request.Header.Set("Idempotency-Key", "create-key-654321")
response, err = http.DefaultClient.Do(request)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusBadRequest {
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","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)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusBadRequest {
t.Fatalf("trailing JSON status = %d", response.StatusCode)
}
_ = 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)
}
}
// TestQueueCreateEnforcesMinProtocolVersion covers the gap multiplayer-next.md
// 8.43 named "version-mismatch-specific client messaging": before this,
// queue_create accepted any protocol_version >= 1 unconditionally, so an
// outdated client below every other queued player's version would simply
// queue forever with no error at all -- the matcher's own compatibility
// check requires every formed player to share an identical protocol_version,
// so it could never be paired, and nothing ever told it why.
func TestQueueCreateEnforcesMinProtocolVersion(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)
}
calls := 0
service := &Service{
Sessions: sessions,
Queue: domain.NewQueue(),
Now: func() time.Time { return now },
MinProtocolVersion: 5,
CandidateV2: func(_ string, ticketID string, spec domain.QueueSpec) (domain.Candidate, error) {
calls++
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, string) {
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)
}
decoded, _ := io.ReadAll(response.Body)
response.Body.Close()
return response, string(decoded)
}
response, body := request(`{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":4}`)
if response.StatusCode != http.StatusUpgradeRequired {
t.Fatalf("below-floor status = %d, want 426 Upgrade Required; body=%s", response.StatusCode, body)
}
if !strings.Contains(body, "client_outdated") {
t.Fatalf("below-floor body does not name the outdated-client error: %s", body)
}
if calls != 0 {
t.Fatalf("candidate provider must not be reached for a rejected below-floor request, calls=%d", calls)
}
response, _ = request(`{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":5}`)
if response.StatusCode != http.StatusCreated {
t.Fatalf("exactly-at-floor status = %d, want 201", response.StatusCode)
}
if calls != 1 {
t.Fatalf("exactly-at-floor request should reach the provider once, calls=%d", calls)
}
}
// TestQueueCreateMinProtocolVersionZeroIsDisabled proves the floor is opt-in:
// every existing Service literal across the codebase that never sets
// MinProtocolVersion must keep accepting protocol_version 1 exactly as
// before, unconditionally.
func TestQueueCreateMinProtocolVersionZeroIsDisabled(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)
}
service := &Service{
Sessions: sessions,
Queue: domain.NewQueue(),
Now: func() time.Time { return now },
CandidateV2: func(_ string, ticketID string, spec domain.QueueSpec) (domain.Candidate, error) {
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":"casual","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.StatusCreated {
t.Fatalf("status = %d, want 201 with MinProtocolVersion left at its zero default", response.StatusCode)
}
}
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 TestQueueCreateAPIRetriesIdenticallyAndRejectsKeyReuseWithChangedPayload(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)
}
service := &Service{Sessions: sessions, Queue: domain.NewQueue(), Now: func() time.Time { return now }, Candidate: func(playerID, ticketID string) (domain.Candidate, error) {
return domain.Candidate{PlayerID: playerID, TicketID: ticketID, EnqueuedAt: now}, nil
}}
server := httptest.NewServer(service.Handler())
defer server.Close()
request := func(body, key string) (int, queueResponse) {
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", key)
response, requestErr := http.DefaultClient.Do(req)
if requestErr != nil {
t.Fatal(requestErr)
}
defer response.Body.Close()
var decoded queueResponse
if response.StatusCode == http.StatusCreated {
if err := json.NewDecoder(response.Body).Decode(&decoded); err != nil {
t.Fatal(err)
}
}
return response.StatusCode, decoded
}
body := `{"ticket_id":"ticket-idempotent","playlist":"casual","client_build":"build-1","protocol_version":1}`
status, first := request(body, "idempotency-key-123456")
if status != http.StatusCreated {
t.Fatalf("first create status=%d", status)
}
status, replay := request(body, "idempotency-key-123456")
if status != http.StatusCreated || replay != first {
t.Fatalf("identical replay status=%d first=%+v replay=%+v", status, first, replay)
}
changed := `{"ticket_id":"ticket-idempotent","playlist":"casual","client_build":"build-2","protocol_version":1}`
status, _ = request(changed, "idempotency-key-123456")
if status != http.StatusConflict {
t.Fatalf("changed-payload replay status=%d, want conflict", status)
}
}
func TestQueueAPIUsesInjectedPersistentBackendWithoutCandidateProvider(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, _ := sessions.Issue("player-1", time.Hour, now)
backend := &queueBackendSpy{}
service := &Service{Sessions: sessions, QueueBackend: backend, Now: func() time.Time { return now }}
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.StatusCreated || backend.createCalls != 1 {
t.Fatalf("status=%d backend_calls=%d", response.StatusCode, backend.createCalls)
}
}
func TestQueueAPIProjectsSuccessfulMutationsWithoutMakingRedisRequired(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, _ := sessions.Issue("player-1", time.Hour, now)
backend := &queueBackendSpy{}
index := &candidateIndexSpy{upsertErr: errors.New("redis unavailable"), removeErr: errors.New("redis unavailable")}
service := &Service{Sessions: sessions, QueueBackend: backend, CandidateIndex: index, Now: func() time.Time { return now }}
server := httptest.NewServer(service.Handler())
defer server.Close()
auth := "Bearer " + session.SessionID + ":" + token
request := func(method, path, key, revision string) *http.Response {
req, _ := http.NewRequest(method, server.URL+path, strings.NewReader(`{"ticket_id":"ticket-1","playlist":"ranked","client_build":"build-1","protocol_version":1}`))
req.Header.Set("Authorization", auth)
if key != "" {
req.Header.Set("Idempotency-Key", key)
}
if revision != "" {
req.Header.Set("If-Match-Revision", revision)
}
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
return response
}
response := request(http.MethodPost, "/v1/queue", "create-key-123456", "")
if response.StatusCode != http.StatusCreated {
t.Fatalf("create status=%d", response.StatusCode)
}
response.Body.Close()
if index.upsertCalls != 1 {
t.Fatalf("upsert calls=%d", index.upsertCalls)
}
response = request(http.MethodPost, "/v1/queue/ticket-1/cancel", "cancel-key-123456", "0")
if response.StatusCode != http.StatusOK {
t.Fatalf("cancel status=%d", response.StatusCode)
}
response.Body.Close()
if index.removeCalls != 1 {
t.Fatalf("remove calls=%d", index.removeCalls)
}
}
func TestQueueAPIDelegatesAllMutationsAndRecoveryToBackend(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, _ := sessions.Issue("player-1", time.Hour, now)
backend := &queueBackendSpy{}
service := &Service{Sessions: sessions, QueueBackend: backend, Now: func() time.Time { return now }}
server := httptest.NewServer(service.Handler())
defer server.Close()
auth := "Bearer " + session.SessionID + ":" + token
request := func(method, path, body, key, revision string) *http.Response {
req, _ := http.NewRequest(method, server.URL+path, strings.NewReader(body))
req.Header.Set("Authorization", auth)
if key != "" {
req.Header.Set("Idempotency-Key", key)
}
if revision != "" {
req.Header.Set("If-Match-Revision", revision)
}
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
return response
}
response := request(http.MethodGet, "/v1/queue/ticket-1", "", "", "")
if response.StatusCode != http.StatusOK {
t.Fatalf("get status = %d", response.StatusCode)
}
response.Body.Close()
response = request(http.MethodPost, "/v1/queue/ticket-1/heartbeat", `{}`, "heartbeat-key-123456", "0")
if response.StatusCode != http.StatusOK {
t.Fatalf("heartbeat status = %d", response.StatusCode)
}
response.Body.Close()
response = request(http.MethodPost, "/v1/queue/ticket-1/cancel", `{}`, "cancel-key-123456", "1")
if response.StatusCode != http.StatusOK {
t.Fatalf("cancel status = %d", response.StatusCode)
}
response.Body.Close()
if backend.getCalls != 1 || backend.heartbeatCalls != 1 || backend.cancelCalls != 1 {
t.Fatalf("backend calls = %+v", backend)
}
}
func TestQueueAPIUsesInjectedSessionBackend(t *testing.T) {
backend := &sessionBackendSpy{}
queue := &queueBackendSpy{}
service := &Service{SessionBackend: backend, QueueBackend: queue, Now: func() time.Time { return time.Unix(1000, 0).UTC() }}
server := httptest.NewServer(service.Handler())
defer server.Close()
req, _ := http.NewRequest(http.MethodGet, server.URL+"/v1/queue/ticket-1", nil)
req.Header.Set("Authorization", "Bearer durable-session:durable-token")
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
if response.StatusCode != http.StatusOK || backend.calls != 1 || queue.getCalls != 1 {
t.Fatalf("status=%d session_calls=%d queue_calls=%d", response.StatusCode, backend.calls, queue.getCalls)
}
}
func TestSteamSessionAPIRequiresBackendVerificationAndIssuesOpaqueSession(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
provider := &steamLoginSpy{}
service := &Service{Sessions: sessions, SteamLogin: provider, Now: func() time.Time { return now }}
server := httptest.NewServer(service.Handler())
defer server.Close()
request := func(body string) *http.Response {
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/session/steam", strings.NewReader(body))
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
return response
}
response := request(`{"web_api_ticket":"valid-web-ticket","steam_id":"spoofed"}`)
if response.StatusCode != http.StatusBadRequest {
t.Fatalf("extra field status = %d", response.StatusCode)
}
response.Body.Close()
response = request(`{"web_api_ticket":"invalid"}`)
if response.StatusCode != http.StatusUnauthorized {
t.Fatalf("invalid ticket status = %d", response.StatusCode)
}
response.Body.Close()
response = request(`{"web_api_ticket":"valid-web-ticket"}`)
if response.StatusCode != http.StatusOK {
t.Fatalf("valid ticket status = %d", response.StatusCode)
}
var result struct {
PlayerID string `json:"player_id"`
AccessToken string `json:"access_token"`
}
if err := json.NewDecoder(response.Body).Decode(&result); err != nil {
t.Fatal(err)
}
response.Body.Close()
if result.PlayerID != "player-1" || !strings.Contains(result.AccessToken, ":") || provider.calls != 2 {
t.Fatalf("session result=%+v provider_calls=%d", result, provider.calls)
}
}
func TestQueueRecoveryAPIIsAuthenticatedOwnerOnlyAndExpiresStaleTickets(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
ownerSession, ownerToken, err := sessions.Issue("player-1", time.Hour, now)
if err != nil {
t.Fatal(err)
}
otherSession, otherToken, err := sessions.Issue("player-2", time.Hour, now)
if err != nil {
t.Fatal(err)
}
queue := domain.NewQueue()
service := &Service{Sessions: sessions, Queue: queue, Now: func() time.Time { return now }, Candidate: func(playerID, ticketID string) (domain.Candidate, error) {
return domain.Candidate{PlayerID: playerID, TicketID: ticketID, EnqueuedAt: now}, nil
}}
server := httptest.NewServer(service.Handler())
defer server.Close()
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)
if err != nil || response.StatusCode != http.StatusCreated {
t.Fatalf("create status=%v err=%v", response.StatusCode, err)
}
_ = response.Body.Close()
get, _ := http.NewRequest(http.MethodGet, server.URL+"/v1/queue/ticket-recovery-123456", nil)
get.Header.Set("Authorization", "Bearer "+ownerSession.SessionID+":"+ownerToken)
response, err = http.DefaultClient.Do(get)
if err != nil || response.StatusCode != http.StatusOK {
t.Fatalf("owner recovery status=%v err=%v", response.StatusCode, err)
}
_ = response.Body.Close()
get.Header.Set("Authorization", "Bearer "+otherSession.SessionID+":"+otherToken)
response, err = http.DefaultClient.Do(get)
if err != nil || response.StatusCode != http.StatusForbidden {
t.Fatalf("cross-player recovery status=%v err=%v", response.StatusCode, err)
}
_ = response.Body.Close()
service.Now = func() time.Time { return now.Add(domain.QueueExpiryWindow) }
get.Header.Set("Authorization", "Bearer "+ownerSession.SessionID+":"+ownerToken)
response, err = http.DefaultClient.Do(get)
if err != nil || response.StatusCode != http.StatusGone {
t.Fatalf("expired recovery status=%v err=%v", response.StatusCode, err)
}
_ = response.Body.Close()
}
func TestReadOnlyAPIsEmitLifecycleEvents(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, err := sessions.Issue("player-a", time.Hour, now)
if err != nil {
t.Fatal(err)
}
queue := domain.NewQueue()
if _, err := queue.Create("player-a", "ticket-read-123456", "create-read-123456", domain.Candidate{PlayerID: "player-a", TicketID: "ticket-read-123456", Playlist: domain.Casual, EnqueuedAt: now}, now); err != nil {
t.Fatal(err)
}
proposal, err := domain.NewProposal("proposal-read-123456", domain.Casual, []string{"player-a", "player-b"}, now)
if err != nil {
t.Fatal(err)
}
events := make([]observability.Event, 0)
service := &Service{
Sessions: sessions, Queue: queue, Proposals: map[string]*domain.Proposal{proposal.ProposalID: &proposal},
RankedProfiles: map[string]domain.RankedProfile{"player-a": {Rating: domain.Rating{Value: 1500, RD: 100, Volatility: 0.06}}},
TierPolicy: domain.DefaultTierPolicy(), Now: func() time.Time { return now },
Assignment: func(_ context.Context, playerID, matchID string, _ time.Time) (AssignmentView, error) {
return AssignmentView{MatchID: matchID, PlayerID: playerID, ServerID: "server-read", Slot: 0, ProtocolVersion: 1, Transport: "enet", Endpoint: "127.0.0.1:30001", JoinAuthorisation: "join-token", ExpiresAt: now.Add(time.Minute)}, nil
},
Log: func(event observability.Event) { events = append(events, event) },
}
server := httptest.NewServer(service.Handler())
defer server.Close()
auth := "Bearer " + session.SessionID + ":" + token
for _, path := range []string{"/v1/queue/ticket-read-123456", "/v1/proposals/proposal-read-123456", "/v1/assignments/match-read-123456", "/api/v1/profile", "/v1/profile/ranked"} {
req, _ := http.NewRequest(http.MethodGet, server.URL+path, nil)
req.Header.Set("Authorization", auth)
response, requestErr := http.DefaultClient.Do(req)
if requestErr != nil {
t.Fatal(requestErr)
}
if response.StatusCode != http.StatusOK {
t.Fatalf("GET %s status = %d", path, response.StatusCode)
}
response.Body.Close()
}
seen := map[string]bool{}
for _, event := range events {
seen[event.Event] = true
}
for _, eventName := range []string{"queue_get", "proposal_get", "assignment_get", "profile_get", "ranked_profile_get"} {
if !seen[eventName] {
t.Fatalf("read event %q missing from %+v", eventName, events)
}
}
}
func TestAuthenticatedProposalAPIUsesRevisionAndIdempotencyPolicy(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, err := sessions.Issue("player-a", time.Hour, now)
if err != nil {
t.Fatal(err)
}
proposal, err := domain.NewProposal("proposal-123456789", domain.Casual, []string{"player-a", "player-b"}, now)
if err != nil {
t.Fatal(err)
}
service := &Service{Sessions: sessions, Proposals: map[string]*domain.Proposal{proposal.ProposalID: &proposal}, Now: func() time.Time { return now }}
server := httptest.NewServer(service.Handler())
defer server.Close()
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/proposals/"+proposal.ProposalID+"/accept", nil)
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
req.Header.Set("Idempotency-Key", "proposal-response-123456")
req.Header.Set("If-Match-Revision", "0")
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusOK {
t.Fatalf("proposal accept status = %d", response.StatusCode)
}
_ = response.Body.Close()
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/proposals/"+proposal.ProposalID+"/accept", nil)
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
req.Header.Set("Idempotency-Key", "proposal-response-654321")
req.Header.Set("If-Match-Revision", "0")
response, err = http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusConflict {
t.Fatalf("stale proposal response status = %d", response.StatusCode)
}
_ = response.Body.Close()
}
func TestRankedProfileAPIReturnsBackendTierAndHidesCasualData(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, err := sessions.Issue("player-a", time.Hour, now)
if err != nil {
t.Fatal(err)
}
policy, err := domain.NewTierPolicy([]domain.TierBand{{Tier: domain.RankTierBronze, MinRating: 0}, {Tier: domain.RankTierGold, MinRating: 1500}})
if err != nil {
t.Fatal(err)
}
service := &Service{
Sessions: sessions,
RankedProfiles: map[string]domain.RankedProfile{"player-a": {Rating: domain.Rating{Value: 1600, RD: 200, Volatility: 0.06}, RankedGames: 10, CurrentSeasonID: "season-current", CurrentSeasonEndsAt: now.Add(48 * time.Hour), LastSeasonID: "season-1"}},
TierPolicy: policy,
Now: func() time.Time { return now },
}
server := httptest.NewServer(service.Handler())
defer server.Close()
req, _ := http.NewRequest(http.MethodGet, server.URL+"/v1/profile/ranked", nil)
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
if response.StatusCode != http.StatusOK {
t.Fatalf("ranked profile status = %d", response.StatusCode)
}
var body rankedProfileResponse
if err := json.NewDecoder(response.Body).Decode(&body); err != nil {
t.Fatal(err)
}
if body.Tier != string(domain.RankTierGold) || body.Provisional || body.RankedGames != 10 || body.SeasonID != "season-current" || body.SeasonEndsAt != "1970-01-03T00:16:40Z" {
t.Fatalf("ranked profile response = %+v", body)
}
}
func TestServerResultAPIRequiresBoundWorkloadAndDelegatesDurableSubmission(t *testing.T) {
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{MatchID: "match-1", ServerID: "server-1"}
submitter := &resultSubmitterSpy{}
service := &Service{Now: func() time.Time { return now }, WorkloadVerify: func(token string, at time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" || !at.Equal(now) {
t.Fatalf("verifier input=%q %v", token, at)
}
return binding, nil
}, ResultSubmitter: submitter}
server := httptest.NewServer(service.Handler())
defer server.Close()
body := `{"match_id":"match-1","result_nonce":"nonce-1234567890","score":{"team_0":3,"team_1":2},"integrity_state":"CERTIFIED"}`
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/result", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "result-key-123456")
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusAccepted {
t.Fatalf("status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
if submitter.calls != 1 || submitter.key != "result-key-123456" || submitter.result.Team0Score != 3 {
t.Fatalf("submission=%+v calls=%d", submitter, submitter.calls)
}
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-2/result", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "result-key-123456")
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusUnauthorized {
t.Fatalf("wrong server status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
}
func TestContractServerRoutesAdaptTwoSegmentPaths(t *testing.T) {
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{AllocationID: "allocation-1", MatchID: "match_1234567890", ServerID: "server_123456789"}
submitter := &resultSubmitterSpy{}
registrar := &serverRegistrarSpy{}
service := &Service{Now: func() time.Time { return now }, WorkloadVerify: func(token string, _ time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" {
return domain.WorkloadBinding{}, errors.New("bad token")
}
return binding, nil
}, ResultSubmitter: submitter, ServerRegistrar: registrar}
server := httptest.NewServer(service.Handler())
defer server.Close()
registerBody := `{"match_id":"match_1234567890","protocol_version":1,"image_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","assignment_ready":false}`
req, _ := http.NewRequest(http.MethodPost, server.URL+"/api/v1/servers/server_123456789/register", strings.NewReader(registerBody))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "contract-register-key-1")
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusNoContent || registrar.calls != 1 {
t.Fatalf("register status=%v err=%v calls=%d", response.StatusCode, err, registrar.calls)
}
response.Body.Close()
resultBody := `{"match_id":"match_1234567890","result_nonce":"nonce-1234567890","score":{"team_0":3,"team_1":2},"integrity_state":"CERTIFIED"}`
req, _ = http.NewRequest(http.MethodPost, server.URL+"/api/v1/servers/server_123456789/result", strings.NewReader(resultBody))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "contract-result-key-123")
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusAccepted || submitter.calls != 1 {
t.Fatalf("result status=%v err=%v calls=%d", response.StatusCode, err, submitter.calls)
}
response.Body.Close()
}
// TestServerMutationLoggingNeverLeaksRequestSecrets is a secret canary: it
// drives the register and result routes end to end with realistic-looking
// bearer tokens and a result nonce, captures every event actually emitted
// through Service.Log during those real requests, and asserts the literal
// secret values never appear anywhere in the encoded output -- not just that
// observability.redact() strips a synthetic value under a known key name (see
// TestEncodeCorrelatesStagesAndRedactsNestedCredentials in the observability
// package for that narrower unit test).
func TestServerMutationLoggingNeverLeaksRequestSecrets(t *testing.T) {
const bearerToken = "wl-canary-secret-do-not-log-9f8e7d6c5b4a"
const resultNonce = "nonce-canary-secret-value-1a2b3c4d5e6f"
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{AllocationID: "allocation-1", MatchID: "match-1", ServerID: "server-1"}
registrar := &serverRegistrarSpy{}
submitter := &resultSubmitterSpy{}
var captured [][]byte
service := &Service{
Now: func() time.Time { return now },
WorkloadVerify: func(token string, _ time.Time) (domain.WorkloadBinding, error) {
if token != bearerToken {
return domain.WorkloadBinding{}, errors.New("bad token")
}
return binding, nil
},
ServerRegistrar: registrar,
ResultSubmitter: submitter,
Log: func(event observability.Event) {
payload, err := observability.Encode(event)
if err != nil {
t.Fatalf("encode event: %v", err)
}
captured = append(captured, payload)
},
}
server := httptest.NewServer(service.Handler())
defer server.Close()
registerBody := `{"match_id":"match-1","protocol_version":1,"image_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","assignment_ready":true}`
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/register", strings.NewReader(registerBody))
req.Header.Set("Authorization", "Bearer "+bearerToken)
req.Header.Set("Idempotency-Key", "canary-register-key-1")
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusNoContent {
t.Fatalf("register status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
resultBody := fmt.Sprintf(`{"match_id":"match-1","result_nonce":%q,"score":{"team_0":3,"team_1":2},"integrity_state":"CERTIFIED"}`, resultNonce)
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/result", strings.NewReader(resultBody))
req.Header.Set("Authorization", "Bearer "+bearerToken)
req.Header.Set("Idempotency-Key", "canary-result-key-1234")
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusAccepted {
t.Fatalf("result status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
// An unauthorized attempt must also log nothing sensitive -- it's the one
// call site handling a token that never even verified successfully.
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/register", strings.NewReader(registerBody))
req.Header.Set("Authorization", "Bearer wrong-"+bearerToken)
req.Header.Set("Idempotency-Key", "canary-register-key-2")
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusUnauthorized {
t.Fatalf("unauthorized register status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
if len(captured) == 0 {
t.Fatal("no events were logged; the canary can't prove anything")
}
all := string(bytes.Join(captured, []byte("\n")))
if strings.Contains(all, bearerToken) {
t.Fatalf("bearer token leaked into logged events: %s", all)
}
if strings.Contains(all, resultNonce) {
t.Fatalf("result nonce leaked into logged events: %s", all)
}
}
func TestQueueAndProposalMutationsLogLifecycleEvents(t *testing.T) {
now := time.Unix(1000, 0).UTC()
queueBackend := &queueBackendSpy{}
proposal, err := domain.NewProposal("proposal-1", domain.Casual, []string{"player-1", "player-2"}, now)
if err != nil {
t.Fatal(err)
}
proposalBackend := &proposalBackendSpy{proposal: proposal}
var captured []observability.Event
service := &Service{
SessionBackend: &sessionBackendSpy{},
QueueBackend: queueBackend,
ProposalBackend: proposalBackend,
Now: func() time.Time { return now },
Log: func(event observability.Event) { captured = append(captured, event) },
}
server := httptest.NewServer(service.Handler())
defer server.Close()
auth := "Bearer session-1:token-1"
request := func(method, path, body string, headers map[string]string) *http.Response {
req, _ := http.NewRequest(method, server.URL+path, strings.NewReader(body))
req.Header.Set("Authorization", auth)
for key, value := range headers {
req.Header.Set(key, value)
}
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
return response
}
create := request(http.MethodPost, "/v1/queue", `{"ticket_id":"ticket-1","playlist":"casual","client_build":"build-1","protocol_version":1}`, map[string]string{"Idempotency-Key": "log-create-key-123456"})
if create.StatusCode != http.StatusCreated {
t.Fatalf("create status = %d", create.StatusCode)
}
create.Body.Close()
heartbeat := request(http.MethodPost, "/v1/queue/ticket-1/heartbeat", `{}`, map[string]string{"Idempotency-Key": "log-heartbeat-key-123456", "If-Match-Revision": "0"})
if heartbeat.StatusCode != http.StatusOK {
t.Fatalf("heartbeat status = %d", heartbeat.StatusCode)
}
heartbeat.Body.Close()
cancel := request(http.MethodPost, "/v1/queue/ticket-1/cancel", `{}`, map[string]string{"Idempotency-Key": "log-cancel-key-123456", "If-Match-Revision": "0"})
if cancel.StatusCode != http.StatusOK {
t.Fatalf("cancel status = %d", cancel.StatusCode)
}
cancel.Body.Close()
respond := request(http.MethodPost, "/v1/proposals/proposal-1/accept", `{}`, map[string]string{"Idempotency-Key": "log-respond-key-123456", "If-Match-Revision": "0"})
if respond.StatusCode != http.StatusOK {
t.Fatalf("proposal accept status = %d", respond.StatusCode)
}
respond.Body.Close()
// Same stale revision again -- the real domain.Proposal.Respond behind
// proposalBackendSpy fences this for real, unlike the dumb queue spy
// above, so this proves the rejection path logs too.
staleRespond := request(http.MethodPost, "/v1/proposals/proposal-1/accept", `{}`, map[string]string{"Idempotency-Key": "log-respond-key-234567", "If-Match-Revision": "0"})
if staleRespond.StatusCode != http.StatusConflict {
t.Fatalf("stale proposal accept status = %d", staleRespond.StatusCode)
}
staleRespond.Body.Close()
want := []struct{ event, id, stage string }{
{"queue_create", "ticket-1", "queued"},
{"queue_heartbeat", "ticket-1", "queued"},
{"queue_cancel", "ticket-1", "cancelled"},
{"proposal_response", "proposal-1", "open"},
{"proposal_response", "proposal-1", "rejected"},
}
if len(captured) != len(want) {
t.Fatalf("captured %d events, want %d: %+v", len(captured), len(want), captured)
}
for i, w := range want {
got := captured[i]
gotID := got.QueueID
if got.Event == "proposal_response" {
gotID = got.ProposalID
}
if got.Event != w.event || gotID != w.id || got.Stage != w.stage {
t.Fatalf("event[%d] = %+v, want {%s %s %s}", i, got, w.event, w.id, w.stage)
}
}
}
func TestServerRegistrationAPIRequiresBoundWorkloadAndValidDigest(t *testing.T) {
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{AllocationID: "allocation-1", MatchID: "match-1", ServerID: "server-1"}
registrar := &serverRegistrarSpy{}
service := &Service{Now: func() time.Time { return now }, WorkloadVerify: func(token string, _ time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" {
return domain.WorkloadBinding{}, errors.New("bad token")
}
return binding, nil
}, ServerRegistrar: registrar}
server := httptest.NewServer(service.Handler())
defer server.Close()
body := `{"match_id":"match-1","protocol_version":1,"image_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","assignment_ready":false}`
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/register", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "register-key-123456")
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusNoContent || registrar.calls != 1 || registrar.binding != binding || registrar.protocol != 1 || registrar.assignmentReady {
t.Fatalf("status=%v err=%v registrar=%+v", response.StatusCode, err, registrar)
}
response.Body.Close()
body = `{"match_id":"match-1","protocol_version":1,"image_digest":"bad","assignment_ready":false}`
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/register", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "register-key-123456")
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusUnprocessableEntity || registrar.calls != 1 {
t.Fatalf("invalid registration status=%v err=%v calls=%d", response.StatusCode, err, registrar.calls)
}
response.Body.Close()
}
func TestServerMutationConflictsAreExportedAsADistinctPrometheusCounter(t *testing.T) {
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{AllocationID: "allocation-1", MatchID: "match-1", ServerID: "server-1"}
registrar := &serverRegistrarSpy{err: domain.ErrConflict}
metrics := observability.NewMetrics()
service := &Service{Now: func() time.Time { return now }, Metrics: metrics, WorkloadVerify: func(token string, _ time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" {
return domain.WorkloadBinding{}, errors.New("bad token")
}
return binding, nil
}, ServerRegistrar: registrar}
server := httptest.NewServer(service.Handler())
defer server.Close()
body := `{"match_id":"match-1","protocol_version":1,"image_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","assignment_ready":false}`
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/register", strings.NewReader(body))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "register-key-123456")
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusConflict {
t.Fatalf("status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
metricsResponse, err := http.Get(server.URL + "/metrics")
if err != nil {
t.Fatal(err)
}
defer metricsResponse.Body.Close()
exported, err := io.ReadAll(metricsResponse.Body)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(string(exported), `cosmic_clash_api_server_conflicts_total{kind="register"} 1`) {
t.Fatalf("register conflict was not exported: %s", exported)
}
}
func TestServerShutdownAPIRequiresBoundWorkloadAndDelegatesAcknowledgement(t *testing.T) {
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{AllocationID: "allocation-1", MatchID: "match-1", ServerID: "server-1"}
shutdowner := &serverShutdownerSpy{}
service := &Service{Now: func() time.Time { return now }, WorkloadVerify: func(token string, at time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" || !at.Equal(now) {
t.Fatalf("verifier input=%q %v", token, at)
}
return binding, nil
}, ServerShutdowner: shutdowner}
server := httptest.NewServer(service.Handler())
defer server.Close()
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/shutdown", strings.NewReader(`{"reason":"server_draining"}`))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "shutdown-key-123456")
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusNoContent {
t.Fatalf("status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
if shutdowner.calls != 1 || shutdowner.binding != binding || shutdowner.reason != "server_draining" || shutdowner.key != "shutdown-key-123456" {
t.Fatalf("shutdown=%+v", shutdowner)
}
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/servers/server-1/shutdown", strings.NewReader(`{"reason":"bad\nreason"}`))
req.Header.Set("Authorization", "Bearer workload-token")
req.Header.Set("Idempotency-Key", "shutdown-key-123456")
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusUnprocessableEntity || shutdowner.calls != 1 {
t.Fatalf("invalid shutdown status=%v err=%v calls=%d", response.StatusCode, err, shutdowner.calls)
}
response.Body.Close()
}
func TestServerConnectionAPIRequiresBoundWorkloadAndOpaqueAssignedPlayer(t *testing.T) {
now := time.Unix(1000, 0).UTC()
binding := domain.WorkloadBinding{AllocationID: "allocation-123456", MatchID: "match-1234567890", ServerID: "server-123456789"}
recorder := &serverConnectionSpy{}
service := &Service{Now: func() time.Time { return now }, WorkloadVerify: func(token string, _ time.Time) (domain.WorkloadBinding, error) {
if token != "workload-token" {
return domain.WorkloadBinding{}, errors.New("bad token")
}
return binding, nil
}, ServerConnections: recorder}
server := httptest.NewServer(service.Handler())
defer server.Close()
request := func(operation, serverID, playerID, token, key, bodySuffix string) (int, string) {
body := fmt.Sprintf(`{"player_id":%q%s}`, playerID, bodySuffix)
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/servers/"+serverID+"/"+operation, strings.NewReader(body))
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Idempotency-Key", key)
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
responseBody, _ := io.ReadAll(response.Body)
response.Body.Close()
return response.StatusCode, string(responseBody)
}
if got, body := request("connect", binding.ServerID, "player-123456789", "workload-token", "connect-player-123456789", `,"expected_generation":0`); got != http.StatusOK || !strings.Contains(body, `"generation":1`) {
t.Fatalf("connection status = %d", got)
}
if recorder.connectCalls != 1 || recorder.binding != binding || recorder.playerID != "player-123456789" || recorder.expectedGeneration != 0 || recorder.key != "connect-player-123456789" {
t.Fatalf("connection receipt = %+v", recorder)
}
if got, body := request("connect", binding.ServerID, "player-legacy-123456", "workload-token", "connect-legacy-123456", ""); got != http.StatusNoContent || body != "" {
t.Fatalf("legacy connection status=%d body=%q", got, body)
}
if got, _ := request("connect", "server-000000000", "player-123456789", "workload-token", "connect-player-123456789", ""); got != http.StatusUnauthorized {
t.Fatalf("wrong server status = %d", got)
}
if got, _ := request("connect", binding.ServerID, "short", "workload-token", "connect-player-short-123", ""); got != http.StatusUnprocessableEntity {
t.Fatalf("short player status = %d", got)
}
if recorder.connectCalls != 2 {
t.Fatalf("invalid receipts reached backend: %d", recorder.connectCalls)
}
if got, _ := request("disconnect", binding.ServerID, "player-123456789", "workload-token", "disconnect-player-123456789", `,"generation":1`); got != http.StatusNoContent {
t.Fatalf("disconnect status = %d", got)
}
if recorder.disconnectCalls != 1 || recorder.generation != 1 {
t.Fatalf("disconnect receipt = %+v", recorder)
}
if got, _ := request("disconnect", binding.ServerID, "player-123456789", "workload-token", "disconnect-zero-123456", ""); got != http.StatusUnprocessableEntity || recorder.disconnectCalls != 1 {
t.Fatalf("zero-generation disconnect status=%d calls=%d", got, recorder.disconnectCalls)
}
recorder.err = errors.New("database unavailable")
if got, _ := request("connect", binding.ServerID, "player-123456789", "workload-token", "connect-player-retry-123", `,"expected_generation":1`); got != http.StatusServiceUnavailable {
t.Fatalf("recorder outage status = %d, want retryable 503", got)
}
}
func TestMetricsEndpointExportsBoundedAPILatencyAndSkipsItsOwnScrape(t *testing.T) {
metrics := observability.NewMetrics()
service := &Service{Metrics: metrics, Now: time.Now}
server := httptest.NewServer(service.Handler())
defer server.Close()
response, err := http.Get(server.URL + "/healthz")
if err != nil {
t.Fatal(err)
}
response.Body.Close()
response, err = http.Get(server.URL + "/unknown/secret-token")
if err != nil {
t.Fatal(err)
}
response.Body.Close()
response, err = http.Get(server.URL + "/metrics")
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
body, err := io.ReadAll(response.Body)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusOK || !strings.Contains(string(body), `operation="other",status="4xx"`) || strings.Contains(string(body), "secret-token") {
t.Fatalf("metrics status=%d body=%s", response.StatusCode, body)
}
}
func TestControlPlaneLivenessAndDatastoreReadinessAreIndependent(t *testing.T) {
ready := false
checks := 0
limiter, err := NewRateLimiter(1, time.Minute, 8)
if err != nil {
t.Fatal(err)
}
service := &Service{
RateLimiter: limiter,
Now: func() time.Time { return time.Unix(1000, 0) },
ReadinessCheck: func(context.Context) error {
checks++
if !ready {
return errors.New("database unavailable")
}
return nil
},
}
handler := service.Handler()
status := func(method, path string) int {
recorder := httptest.NewRecorder()
handler.ServeHTTP(recorder, httptest.NewRequest(method, path, nil))
return recorder.Code
}
if got := status(http.MethodGet, "/healthz"); got != http.StatusOK {
t.Fatalf("liveness during datastore outage = %d", got)
}
if got := status(http.MethodGet, "/readyz"); got != http.StatusServiceUnavailable {
t.Fatalf("readiness during datastore outage = %d", got)
}
ready = true
if got := status(http.MethodGet, "/readyz"); got != http.StatusOK {
t.Fatalf("recovered readiness = %d", got)
}
if got := status(http.MethodGet, "/healthz"); got != http.StatusOK {
t.Fatalf("repeated probe was incorrectly rate limited: %d", got)
}
if checks != 2 {
t.Fatalf("readiness checks = %d", checks)
}
if got := status(http.MethodPost, "/readyz"); got != http.StatusMethodNotAllowed {
t.Fatalf("readiness mutation status = %d", got)
}
}
func TestControlPlaneReadinessFailsClosedWithoutCheck(t *testing.T) {
recorder := httptest.NewRecorder()
(&Service{}).Handler().ServeHTTP(recorder, httptest.NewRequest(http.MethodGet, "/readyz", nil))
if recorder.Code != http.StatusServiceUnavailable {
t.Fatalf("unconfigured readiness status = %d", recorder.Code)
}
}
func TestProbeAPIUsesServerEvidenceAndRejectsClientRTTField(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, err := sessions.Issue("player-a", time.Hour, now)
if err != nil {
t.Fatal(err)
}
called := false
// A ProbeRecorder is required: accepting a probe without persisting it
// reports success while leaving predicted_rtt empty, which silently keeps
// the ticket invisible to the matcher.
service := &Service{Sessions: sessions, Now: func() time.Time { return now }, ProbeRecorder: &probeRecorderSpy{}, Probe: func(_ context.Context, playerID, region string, location, nonce []byte, receivedAt time.Time) (domain.ProbeEvidence, []byte, error) {
called = true
if playerID != "player-a" || region != "EU" || string(location) != "opaque" || string(nonce) != "nonce" || !receivedAt.Equal(now) {
t.Fatalf("probe provider arguments = %q %s %q %q %v", playerID, region, location, nonce, receivedAt)
}
return domain.ProbeEvidence{OpaqueLocation: location, Nonce: nonce, IssuedAt: now, Region: region, ServerRTT: 40 * time.Millisecond}, nonce, nil
}}
server := httptest.NewServer(service.Handler())
defer server.Close()
request := `{"opaque_location":"b3BhcXVl","nonce":"bm9uY2U="}`
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/probes/EU", strings.NewReader(request))
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusAccepted || !called {
t.Fatalf("valid probe status=%v err=%v called=%v", response.StatusCode, err, called)
}
_ = response.Body.Close()
request = `{"opaque_location":"b3BhcXVl","nonce":"bm9uY2U=","server_rtt_ms":1}`
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/probes/EU", strings.NewReader(request))
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusBadRequest {
t.Fatalf("client RTT field status=%v err=%v", response.StatusCode, err)
}
_ = response.Body.Close()
}
func TestProbeAPIRecordsOnlyValidatedServerEvidence(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, _ := sessions.Issue("player-a", time.Hour, now)
recorder := &probeRecorderSpy{}
service := &Service{Sessions: sessions, Now: func() time.Time { return now }, ProbeRecorder: recorder, Probe: func(_ context.Context, _ string, region string, location, nonce []byte, _ time.Time) (domain.ProbeEvidence, []byte, error) {
return domain.ProbeEvidence{OpaqueLocation: location, Nonce: nonce, IssuedAt: now, Region: region, ServerRTT: 37 * time.Millisecond}, nonce, nil
}}
server := httptest.NewServer(service.Handler())
defer server.Close()
req, _ := http.NewRequest(http.MethodPost, server.URL+"/v1/probes/NA", strings.NewReader(`{"opaque_location":"b3BhcXVl","nonce":"bm9uY2U="}`))
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, err := http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusAccepted {
t.Fatalf("status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
if recorder.calls != 1 || recorder.last.player != "player-a" || recorder.last.region != "NA" || recorder.last.rtt != 37*time.Millisecond {
t.Fatalf("recorded probe=%+v calls=%d", recorder.last, recorder.calls)
}
recorder.err = errors.New("database unavailable")
req, _ = http.NewRequest(http.MethodPost, server.URL+"/v1/probes/NA", strings.NewReader(`{"opaque_location":"b3BhcXVl","nonce":"bm9uY2U="}`))
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, err = http.DefaultClient.Do(req)
if err != nil || response.StatusCode != http.StatusServiceUnavailable {
t.Fatalf("persistence status=%v err=%v", response.StatusCode, err)
}
response.Body.Close()
}
func TestProposalRecoveryIsParticipantScopedAndExpiresAtReadBoundary(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
owner, ownerToken, err := sessions.Issue("player-a", time.Hour, now)
if err != nil {
t.Fatal(err)
}
other, otherToken, err := sessions.Issue("player-z", time.Hour, now)
if err != nil {
t.Fatal(err)
}
proposal, err := domain.NewProposal("proposal-recovery", domain.Casual, []string{"player-a", "player-b"}, now)
if err != nil {
t.Fatal(err)
}
current := now
service := &Service{Sessions: sessions, Proposals: map[string]*domain.Proposal{proposal.ProposalID: &proposal}, Now: func() time.Time { return current }}
server := httptest.NewServer(service.Handler())
defer server.Close()
get := func(session domain.Session, token string) (int, proposalResponse) {
req, _ := http.NewRequest(http.MethodGet, server.URL+"/v1/proposals/proposal-recovery", nil)
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, requestErr := http.DefaultClient.Do(req)
if requestErr != nil {
t.Fatal(requestErr)
}
defer response.Body.Close()
var body proposalResponse
if response.StatusCode == http.StatusOK {
if err := json.NewDecoder(response.Body).Decode(&body); err != nil {
t.Fatal(err)
}
}
return response.StatusCode, body
}
status, recovered := get(owner, ownerToken)
if status != http.StatusOK || recovered.State != string(domain.Open) || recovered.Revision != 0 {
t.Fatalf("owner recovery status=%d body=%+v", status, recovered)
}
status, _ = get(other, otherToken)
if status != http.StatusNotFound {
t.Fatalf("non-participant recovery status=%d, want 404", status)
}
current = now.Add(domain.ProposalWindow)
status, recovered = get(owner, ownerToken)
if status != http.StatusOK || recovered.State != string(domain.Expired) || recovered.Revision != 1 {
t.Fatalf("expired recovery status=%d body=%+v", status, recovered)
}
}
func TestAssignmentRecoveryIsPlayerScopedAndRejectsExpiredOrMismatchedViews(t *testing.T) {
now := time.Unix(1000, 0).UTC()
sessions := domain.NewSessionStore()
session, token, err := sessions.Issue("player-a", time.Hour, now)
if err != nil {
t.Fatal(err)
}
other, otherToken, err := sessions.Issue("player-z", time.Hour, now)
if err != nil {
t.Fatal(err)
}
current := now
service := &Service{Sessions: sessions, Now: func() time.Time { return current }, Assignment: func(_ context.Context, _ string, matchID string, _ time.Time) (AssignmentView, error) {
return AssignmentView{MatchID: matchID, ServerID: "server-1", PlayerID: "player-a", Slot: 2, ExpiresAt: now.Add(time.Minute), ProtocolVersion: 1, Transport: "enet", Endpoint: "127.0.0.1:30001", JoinAuthorisation: "signed-join"}, nil
}}
server := httptest.NewServer(service.Handler())
defer server.Close()
get := func(path string) (int, AssignmentView) {
req, _ := http.NewRequest(http.MethodGet, server.URL+path, nil)
req.Header.Set("Authorization", "Bearer "+session.SessionID+":"+token)
response, requestErr := http.DefaultClient.Do(req)
if requestErr != nil {
t.Fatal(requestErr)
}
defer response.Body.Close()
var view AssignmentView
if response.StatusCode == http.StatusOK {
if err := json.NewDecoder(response.Body).Decode(&view); err != nil {
t.Fatal(err)
}
}
return response.StatusCode, view
}
status, view := get("/v1/assignments/match-1")
if status != http.StatusOK || view.PlayerID != "player-a" || view.Slot != 2 {
t.Fatalf("assignment status=%d view=%+v", status, view)
}
req, _ := http.NewRequest(http.MethodGet, server.URL+"/v1/assignments/match-1", nil)
req.Header.Set("Authorization", "Bearer "+other.SessionID+":"+otherToken)
response, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
if response.StatusCode != http.StatusNotFound {
t.Fatalf("misbound assignment status=%d, want 404", response.StatusCode)
}
response.Body.Close()
status, _ = get("/v1/assignments/")
if status != http.StatusNotFound {
t.Fatalf("malformed assignment path status=%d", status)
}
current = now.Add(time.Minute)
status, _ = get("/v1/assignments/match-1")
if status != http.StatusServiceUnavailable {
t.Fatalf("expired assignment status=%d", status)
}
}
func TestAssignmentEventUsesAuthoritativeRevisionWithoutChangingResponseShape(t *testing.T) {
now := time.Unix(1000, 0).UTC()
view := AssignmentView{MatchID: "match-1", ServerID: "server-1", PlayerID: "player-a", Revision: 7}
event := assignmentChangedEvent(view, now)
if event.Revision != 7 || event.ResourceID != "match-1" || event.PlayerID != "player-a" {
t.Fatalf("assignment event = %+v", event)
}
payload, err := json.Marshal(view)
if err != nil {
t.Fatal(err)
}
if strings.Contains(string(payload), "revision") {
t.Fatalf("assignment response leaked event revision: %s", payload)
}
}
func TestAssignmentEndpointValidationRejectsAmbiguousOrUnsafeEndpoints(t *testing.T) {
for _, endpoint := range []string{"", "127.0.0.1", "127.0.0.1:0", "127.0.0.1:70000", "https://127.0.0.1:1", "127.0.0.1:1/path"} {
if validAssignmentEndpoint(endpoint) {
t.Fatalf("unsafe endpoint accepted: %q", endpoint)
}
}
for _, endpoint := range []string{"127.0.0.1:1", "example.invalid:65535", "[2001:db8::1]:31001"} {
if !validAssignmentEndpoint(endpoint) {
t.Fatalf("valid endpoint rejected: %q", endpoint)
}
}
}