feat: add retry-safe matchmaking queue domain

This commit is contained in:
Josh Creek
2026-08-31 20:20:42 +01:00
parent b79d358db9
commit 6c0163c3ec
2 changed files with 197 additions and 0 deletions
+150
View File
@@ -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")
}
+47
View File
@@ -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) }
}