mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 00:14:00 +00:00
4eaa3304c3
Wire the same Service.Log hook added for the server register/result routes into queue create/heartbeat/cancel and proposal accept/decline: log the resulting state on success (queue_create, queue_heartbeat, queue_cancel, proposal_response) or 'rejected' on a domain error, using only the ticket/proposal ID and outcome -- never the domain error text itself, which isn't documented as credential-free. Read-only routes (queue GET, proposal GET, assignment fetch) and the early availability/not-found rejections that return before reaching the domain call are deliberately not logged in this pass. Covered by a new end-to-end test driving real create/heartbeat/cancel and an accept followed by a stale-revision accept (fenced for real by the domain layer behind proposalBackendSpy, unlike the dumb queue spy), asserting the exact sequence of events logged.
1473 lines
62 KiB
Go
1473 lines
62 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 }
|
|
|
|
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
|
|
}
|
|
|
|
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, _ 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 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 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: "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: "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)
|
|
}
|
|
}
|
|
|
|
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 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, 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-1" {
|
|
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-1", ServerID: "server-1"}
|
|
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-1","protocol_version":1,"image_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","assignment_ready":false}`
|
|
req, _ := http.NewRequest(http.MethodPost, server.URL+"/api/v1/servers/server-1/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-1","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-1/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 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
|
|
service := &Service{Sessions: sessions, Now: func() time.Time { return now }, Probe: func(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(_ 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)
|
|
}
|
|
}
|
|
}
|