package domain import ( "errors" "reflect" "sync" "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 || !reflect.DeepEqual(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 || !reflect.DeepEqual(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) } } func TestQueueConcurrentCreateKeepsOneActiveTicketPerPlayer(t *testing.T) { q := NewQueue() now := time.Unix(1000, 0) var wg sync.WaitGroup results := make(chan error, 2) for i := 0; i < 2; i++ { wg.Add(1) go func(i int) { defer wg.Done() id := string(rune('a' + i)) _, err := q.Create("same-player", "ticket-"+id, "create-"+id+"-123456", Candidate{TicketID: "ticket-" + id, PlayerID: "same-player", EnqueuedAt: now}, now) results <- err }(i) } wg.Wait() close(results) succeeded := 0 for err := range results { if err == nil { succeeded++ } else if !errors.Is(err, ErrPlayerQueued) { t.Fatalf("unexpected concurrent create error: %v", err) } } if succeeded != 1 { t.Fatalf("concurrent creates succeeded %d times", succeeded) } }