From 6c0163c3ec889a1cade4d334a6e20fa6cae0a923 Mon Sep 17 00:00:00 2001 From: Josh Creek <8179928+jcreek@users.noreply.github.com> Date: Mon, 31 Aug 2026 20:20:42 +0100 Subject: [PATCH] feat: add retry-safe matchmaking queue domain --- server/domain/queue.go | 150 ++++++++++++++++++++++++++++++++++++ server/domain/queue_test.go | 47 +++++++++++ 2 files changed, 197 insertions(+) create mode 100644 server/domain/queue.go create mode 100644 server/domain/queue_test.go diff --git a/server/domain/queue.go b/server/domain/queue.go new file mode 100644 index 00000000..2598beaf --- /dev/null +++ b/server/domain/queue.go @@ -0,0 +1,150 @@ +package domain + +import ( + "crypto/sha256" + "errors" + "fmt" + "sort" + "strings" + "time" +) + +const ( + QueueHeartbeatInterval = 10 * time.Second + QueueExpiryWindow = 30 * time.Second +) + +var ( + ErrPlayerQueued = errors.New("player already owns an active queue ticket") + ErrTicketNotFound = errors.New("queue ticket not found") + ErrNotTicketOwner = errors.New("queue ticket is owned by another player") + ErrTicketExpired = errors.New("queue ticket expired") +) + +type QueueTicket struct { + TicketID string + PlayerID string + Candidate Candidate + State State + Revision uint64 + EnqueuedAt time.Time + ExpiresAt time.Time +} + +type queueMutation struct { + digest [32]byte + ticket QueueTicket +} + +type Queue struct { + tickets map[string]QueueTicket + byPlayer map[string]string + mutations map[string]queueMutation +} + +func NewQueue() *Queue { + return &Queue{tickets: make(map[string]QueueTicket), byPlayer: make(map[string]string), mutations: make(map[string]queueMutation)} +} + +// Create is the in-process equivalent of the PostgreSQL ownership fence. The +// production adapter must perform the same check in one transaction and use +// the same idempotency semantics. +func (q *Queue) Create(playerID, ticketID, idempotencyKey string, candidate Candidate, now time.Time) (QueueTicket, error) { + digest := sha256.Sum256([]byte(createPayload(playerID, ticketID, candidate))) + if prior, ok := q.mutations[idempotencyKey]; ok { + if prior.digest != digest { return QueueTicket{}, fmt.Errorf("%w: create payload changed", ErrConflict) } + return prior.ticket, nil + } + if idempotencyKey == "" || playerID == "" || ticketID == "" || candidate.TicketID != ticketID { + return QueueTicket{}, fmt.Errorf("%w: invalid queue create", ErrConflict) + } + if _, ok := q.byPlayer[playerID]; ok { return QueueTicket{}, ErrPlayerQueued } + if _, ok := q.tickets[ticketID]; ok { return QueueTicket{}, fmt.Errorf("%w: ticket ID already exists", ErrConflict) } + ticket := QueueTicket{TicketID: ticketID, PlayerID: playerID, Candidate: candidate, State: Queued, EnqueuedAt: now, ExpiresAt: now.Add(QueueExpiryWindow)} + q.tickets[ticketID] = ticket + q.byPlayer[playerID] = ticketID + q.mutations[idempotencyKey] = queueMutation{digest: digest, ticket: ticket} + return ticket, nil +} + +func (q *Queue) Heartbeat(playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (QueueTicket, error) { + digest := sha256.Sum256([]byte(fmt.Sprintf("heartbeat:%s:%d", ticketID, expectedRevision))) + if prior, ok := q.mutations[idempotencyKey]; ok { + if prior.digest != digest { return QueueTicket{}, fmt.Errorf("%w: heartbeat payload changed", ErrConflict) } + return prior.ticket, nil + } + ticket, err := q.ownedTicket(playerID, ticketID) + if err != nil { return QueueTicket{}, err } + if now.After(ticket.ExpiresAt) || now.Equal(ticket.ExpiresAt) { return QueueTicket{}, ErrTicketExpired } + if ticket.Revision != expectedRevision { return QueueTicket{}, ErrStaleRevision } + if ticket.State != Queued && ticket.State != Proposed { return QueueTicket{}, fmt.Errorf("%w: heartbeat in %s", ErrConflict, ticket.State) } + if idempotencyKey == "" { return QueueTicket{}, fmt.Errorf("%w: empty heartbeat key", ErrConflict) } + ticket.Revision++ + ticket.ExpiresAt = now.Add(QueueExpiryWindow) + q.tickets[ticketID] = ticket + q.mutations[idempotencyKey] = queueMutation{digest: digest, ticket: ticket} + return ticket, nil +} + +func (q *Queue) Cancel(playerID, ticketID, idempotencyKey string, expectedRevision uint64, now time.Time) (QueueTicket, error) { + digest := sha256.Sum256([]byte(fmt.Sprintf("cancel:%s:%d", ticketID, expectedRevision))) + if prior, ok := q.mutations[idempotencyKey]; ok { + if prior.digest != digest { return QueueTicket{}, fmt.Errorf("%w: cancel payload changed", ErrConflict) } + return prior.ticket, nil + } + ticket, err := q.ownedTicket(playerID, ticketID) + if err != nil { return QueueTicket{}, err } + if ticket.Revision != expectedRevision { return QueueTicket{}, ErrStaleRevision } + if idempotencyKey == "" { return QueueTicket{}, fmt.Errorf("%w: empty cancel key", ErrConflict) } + ticket.State = Cancelled + ticket.Revision++ + ticket.ExpiresAt = now + q.tickets[ticketID] = ticket + delete(q.byPlayer, playerID) + q.mutations[idempotencyKey] = queueMutation{digest: digest, ticket: ticket} + return ticket, nil +} + +func (q *Queue) Expire(now time.Time) []QueueTicket { + var expired []QueueTicket + for id, ticket := range q.tickets { + if (ticket.State == Queued || ticket.State == Proposed) && !now.Before(ticket.ExpiresAt) { + ticket.State = Expired + ticket.Revision++ + q.tickets[id] = ticket + delete(q.byPlayer, ticket.PlayerID) + expired = append(expired, ticket) + } + } + sort.Slice(expired, func(i, j int) bool { return expired[i].TicketID < expired[j].TicketID }) + return expired +} + +func (q *Queue) Candidates(now time.Time) []Candidate { + q.Expire(now) + result := make([]Candidate, 0) + for _, ticket := range q.tickets { + if ticket.State == Queued { result = append(result, ticket.Candidate) } + } + sort.Slice(result, func(i, j int) bool { + if !result[i].EnqueuedAt.Equal(result[j].EnqueuedAt) { return result[i].EnqueuedAt.Before(result[j].EnqueuedAt) } + return result[i].TicketID < result[j].TicketID + }) + return result +} + +func (q *Queue) ownedTicket(playerID, ticketID string) (QueueTicket, error) { + ticket, ok := q.tickets[ticketID] + if !ok { return QueueTicket{}, ErrTicketNotFound } + if ticket.PlayerID != playerID { return QueueTicket{}, ErrNotTicketOwner } + return ticket, nil +} + +func createPayload(playerID, ticketID string, candidate Candidate) string { + regions := make([]string, 0, len(candidate.PredictedRTT)) + for region := range candidate.PredictedRTT { regions = append(regions, region) } + sort.Strings(regions) + rtts := make([]string, 0, len(regions)) + for _, region := range regions { rtts = append(rtts, fmt.Sprintf("%s=%.9f", region, candidate.PredictedRTT[region])) } + return strings.Join([]string{playerID, ticketID, candidate.PlayerID, candidate.TicketID, fmt.Sprintf("%.9f", candidate.Rating), candidate.EnqueuedAt.UTC().Format(time.RFC3339Nano), strings.Join(rtts, ",")}, "\x00") +} diff --git a/server/domain/queue_test.go b/server/domain/queue_test.go new file mode 100644 index 00000000..ae27a41f --- /dev/null +++ b/server/domain/queue_test.go @@ -0,0 +1,47 @@ +package domain + +import ( + "errors" + "testing" + "time" +) + +func TestQueueFencesOneActiveTicketPerPlayerAndReplaysCreate(t *testing.T) { + q := NewQueue() + now := time.Unix(1000, 0) + c := Candidate{TicketID: "ticket-a", PlayerID: "player-a", EnqueuedAt: now, PredictedRTT: map[string]float64{"EU": 20}} + first, err := q.Create("player-a", "ticket-a", "create-key-123456", c, now) + if err != nil { t.Fatal(err) } + replay, err := q.Create("player-a", "ticket-a", "create-key-123456", c, now.Add(time.Second)) + if err != nil || replay != first { t.Fatalf("create replay = %+v, %v", replay, err) } + other := c; other.TicketID = "ticket-b" + if _, err := q.Create("player-a", "ticket-b", "create-key-654321", other, now); !errors.Is(err, ErrPlayerQueued) { t.Fatalf("second active ticket error = %v", err) } +} + +func TestQueueHeartbeatExtendsExpiryExactlyAndRejectsStaleReplay(t *testing.T) { + q := NewQueue(); now := time.Unix(1000, 0) + c := Candidate{TicketID: "ticket-a", PlayerID: "player-a", EnqueuedAt: now} + if _, err := q.Create("player-a", "ticket-a", "create-key-123456", c, now); err != nil { t.Fatal(err) } + updated, err := q.Heartbeat("player-a", "ticket-a", "heartbeat-key-123", 0, now.Add(10*time.Second)) + if err != nil { t.Fatal(err) } + if !updated.ExpiresAt.Equal(now.Add(40 * time.Second)) || updated.Revision != 1 { t.Fatalf("bad heartbeat: %+v", updated) } + replay, err := q.Heartbeat("player-a", "ticket-a", "heartbeat-key-123", 0, now.Add(50*time.Second)) + if err != nil || replay != updated { t.Fatalf("heartbeat replay = %+v, %v", replay, err) } + if _, err := q.Heartbeat("player-a", "ticket-a", "heartbeat-key-456", 0, now.Add(20*time.Second)); !errors.Is(err, ErrStaleRevision) { t.Fatalf("stale heartbeat error = %v", err) } +} + +func TestQueueExpiryReleasesOwnershipAndDoesNotReturnExpiredCandidates(t *testing.T) { + q := NewQueue(); now := time.Unix(1000, 0) + c := Candidate{TicketID: "ticket-a", PlayerID: "player-a", EnqueuedAt: now} + if _, err := q.Create("player-a", "ticket-a", "create-key-123456", c, now); err != nil { t.Fatal(err) } + if got := q.Candidates(now.Add(QueueExpiryWindow)); len(got) != 0 { t.Fatalf("expired candidate returned: %+v", got) } + if _, err := q.Create("player-a", "ticket-b", "create-key-654321", Candidate{TicketID: "ticket-b", PlayerID: "player-a"}, now.Add(QueueExpiryWindow)); err != nil { t.Fatalf("ownership was not released: %v", err) } +} + +func TestQueueCreateIdempotencyIncludesCandidatePayload(t *testing.T) { + q := NewQueue(); now := time.Unix(1000, 0) + base := Candidate{TicketID: "ticket-a", PlayerID: "player-a", Rating: 1500, EnqueuedAt: now, PredictedRTT: map[string]float64{"EU": 20}} + if _, err := q.Create("player-a", "ticket-a", "create-key-123456", base, now); err != nil { t.Fatal(err) } + changed := base; changed.Rating = 1800 + if _, err := q.Create("player-a", "ticket-a", "create-key-123456", changed, now); !errors.Is(err, ErrConflict) { t.Fatalf("changed create payload error = %v", err) } +}