mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 17:53:44 +00:00
79318b56bd
The last two commits fixed the severe stranding bug in decline and timeout, but left a real responsiveness gap: cancelling a ticket directly while it's part of an OPEN proposal used to leave the OTHER participant waiting out the full response window for something the system already knew couldn't happen -- their proposal partner just abandoned the queue. ProposalExpireRequeueSQL eventually rescues them, but only after the full window elapses, not immediately. CascadeCancelToOpenProposal runs inside the same transaction as the cancel itself: if the cancelled ticket belonged to a currently-OPEN proposal, decline that proposal right now and requeue every other participant immediately via the same ProposalDeclineRequeueSQL the decline path already uses. The cancelling player's own ticket correctly stays CANCELLED, not swept back into the requeue meant for everyone else (ProposalDeclineRequeueSQL only touches tickets still at PROPOSED). Covered by a real PostgreSQL integration test: cancelling one participant's ticket mid-proposal immediately declines the proposal and requeues the other participant with a refreshed expiry, while the cancelling player's own ticket stays CANCELLED. First draft used a stale expected revision (0) for the cancel call -- CreateProposal's own QueueTicketProposeSQL already bumps a ticket's revision to 1 when forming the proposal, caught immediately by actually running the test against real Postgres rather than assuming. Clean across 5 runs after the fix, plus the full integration and unit suites.
1222 lines
57 KiB
Go
1222 lines
57 KiB
Go
//go:build integration
|
|
|
|
package store
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"database/sql"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
|
"github.com/cosmic-clash/cosmic-clash/server/migrations"
|
|
_ "github.com/jackc/pgx/v5/stdlib"
|
|
)
|
|
|
|
// This binary is deliberately opt-in. It requires a disposable PostgreSQL
|
|
// instance supplied by scripts/run_postgres_integration.sh.
|
|
func openIntegrationPostgres(t *testing.T) *sql.DB {
|
|
t.Helper()
|
|
dsn := os.Getenv("COSMIC_CLASH_POSTGRES_DSN")
|
|
if dsn == "" {
|
|
t.Skip("COSMIC_CLASH_POSTGRES_DSN is not set")
|
|
}
|
|
db, err := sql.Open("pgx", dsn)
|
|
if err != nil {
|
|
t.Fatalf("open PostgreSQL: %v", err)
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := db.PingContext(ctx); err != nil {
|
|
db.Close()
|
|
t.Fatalf("ping PostgreSQL: %v", err)
|
|
}
|
|
t.Cleanup(func() { db.Close() })
|
|
return db
|
|
}
|
|
|
|
func applyIntegrationMigrations(t *testing.T, db *sql.DB) {
|
|
t.Helper()
|
|
if _, err := db.ExecContext(context.Background(), `DROP TABLE IF EXISTS schema_migrations, assignments, audit_events, outbox, result_receipts, ranked_season_rollovers, penalties, seasons, ratings, match_participants, matches, proposal_participants, proposals, allocations, game_servers, queue_tickets, idempotency_keys, sessions, identities CASCADE`); err != nil {
|
|
t.Fatalf("reset PostgreSQL schema: %v", err)
|
|
}
|
|
if err := migrations.Apply(context.Background(), db, filepath.Join("..", "migrations")); err != nil {
|
|
t.Fatalf("apply migrations: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLAllocatorClaimReplayAndCapacityFence(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
servers := []domain.ReadyServer{
|
|
{ServerID: "allocator-server-b", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet", State: domain.ServerReady},
|
|
{ServerID: "allocator-server-a", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet", State: domain.ServerReady},
|
|
}
|
|
for _, server := range servers {
|
|
if err := RegisterReadyServer(ctx, db, server, now); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
request := domain.AllocationRequest{AllocationID: "allocation-integration-1", MatchID: "match-integration-1", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet"}
|
|
allocation, err := ClaimAllocation(ctx, db, request, now)
|
|
if err != nil {
|
|
t.Fatalf("claim: %v", err)
|
|
}
|
|
if allocation.ServerID != "allocator-server-a" || allocation.State != domain.ServerAllocated {
|
|
t.Fatalf("allocation=%+v", allocation)
|
|
}
|
|
if err := RegisterReadyServer(ctx, db, servers[1], now.Add(time.Second)); err != nil {
|
|
t.Fatalf("stale Ready projection: %v", err)
|
|
}
|
|
var lifecycle string
|
|
if err := db.QueryRowContext(ctx, `SELECT state FROM game_servers WHERE server_id = 'allocator-server-a'`).Scan(&lifecycle); err != nil || lifecycle != "ALLOCATED" {
|
|
t.Fatalf("stale Ready projection reopened allocation state=%q err=%v", lifecycle, err)
|
|
}
|
|
replay, err := ClaimAllocation(ctx, db, request, now.Add(time.Second))
|
|
if err != nil || replay.ServerID != allocation.ServerID || !replay.AllocatedAt.Equal(allocation.AllocatedAt) {
|
|
t.Fatalf("replay=%+v err=%v", replay, err)
|
|
}
|
|
conflict := request
|
|
conflict.MatchID = "match-integration-other"
|
|
if _, err := ClaimAllocation(ctx, db, conflict, now); err != domain.ErrConflict {
|
|
t.Fatalf("conflicting replay err=%v", err)
|
|
}
|
|
second := request
|
|
second.AllocationID = "allocation-integration-2"
|
|
second.MatchID = "match-integration-2"
|
|
if _, err := ClaimAllocation(ctx, db, second, now); err != nil {
|
|
t.Fatalf("second claim: %v", err)
|
|
}
|
|
third := second
|
|
third.AllocationID = "allocation-integration-3"
|
|
third.MatchID = "match-integration-3"
|
|
if _, err := ClaimAllocation(ctx, db, third, now); err != domain.ErrNoCapacity {
|
|
t.Fatalf("capacity err=%v", err)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLConcurrentAllocationClaimNeverDoubleBooksAReadyServer is the
|
|
// live counterpart to TestPostgreSQLAllocatorClaimReplayAndCapacityFence: that
|
|
// test claims strictly one request at a time, so it cannot show what happens
|
|
// when two allocator replicas race for the same compatible capacity, which is
|
|
// exactly the scenario 8.30's "bounded cross-replica retry" is about. Register
|
|
// fewer Ready servers than concurrent requests and fire them all at once;
|
|
// exactly as many must win as there was capacity, each winner must get a
|
|
// distinct server, and every loser must fail with ErrNoCapacity rather than a
|
|
// raw serialization error, a duplicate claim, or a hang.
|
|
func TestPostgreSQLConcurrentAllocationClaimNeverDoubleBooksAReadyServer(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
const capacity = 3
|
|
const contenders = 6
|
|
for i := 0; i < capacity; i++ {
|
|
server := domain.ReadyServer{ServerID: fmt.Sprintf("race-server-%d", i), Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet", State: domain.ServerReady}
|
|
if err := RegisterReadyServer(ctx, db, server, now); err != nil {
|
|
t.Fatalf("register %s: %v", server.ServerID, err)
|
|
}
|
|
}
|
|
|
|
var wg sync.WaitGroup
|
|
allocations := make([]domain.Allocation, contenders)
|
|
errs := make([]error, contenders)
|
|
wg.Add(contenders)
|
|
for i := 0; i < contenders; i++ {
|
|
go func(i int) {
|
|
defer wg.Done()
|
|
request := domain.AllocationRequest{AllocationID: fmt.Sprintf("race-allocation-%d", i), MatchID: fmt.Sprintf("race-match-%d", i), Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet"}
|
|
allocations[i], errs[i] = ClaimAllocation(ctx, db, request, now)
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
|
|
wonServers := map[string]int{}
|
|
won, lost := 0, 0
|
|
for i, err := range errs {
|
|
switch {
|
|
case err == nil:
|
|
won++
|
|
if allocations[i].ServerID == "" {
|
|
t.Fatalf("claim %d succeeded with no server", i)
|
|
}
|
|
wonServers[allocations[i].ServerID]++
|
|
case errors.Is(err, domain.ErrNoCapacity):
|
|
lost++
|
|
default:
|
|
t.Fatalf("claim %d failed with unexpected error: %v", i, err)
|
|
}
|
|
}
|
|
if won != capacity || lost != contenders-capacity {
|
|
t.Fatalf("won=%d lost=%d, want won=%d lost=%d", won, lost, capacity, contenders-capacity)
|
|
}
|
|
if len(wonServers) != capacity {
|
|
t.Fatalf("expected %d distinct servers claimed, got %d: %v", capacity, len(wonServers), wonServers)
|
|
}
|
|
for server, count := range wonServers {
|
|
if count != 1 {
|
|
t.Fatalf("server %s was claimed %d times", server, count)
|
|
}
|
|
}
|
|
var allocatedCount int
|
|
if err := db.QueryRowContext(ctx, `SELECT count(*) FROM game_servers WHERE state = 'ALLOCATED'`).Scan(&allocatedCount); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if allocatedCount != capacity {
|
|
t.Fatalf("durable ALLOCATED server count = %d, want %d", allocatedCount, capacity)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLAcceptedProposalPromotesOneAtomicAllocatingMatch(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"promote-a", "promote-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
for index, player := range []string{"promote-a", "promote-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'PROPOSED', 'build-1', 1, $3, $4)`, fmt.Sprintf("promote-ticket-%d", index), player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO proposals (proposal_id, playlist, state, expires_at, revision) VALUES ('promote-proposal', 'casual', 'ACCEPTED', $1, 2)`, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
for index, player := range []string{"promote-a", "promote-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO proposal_participants (proposal_id, player_id, ticket_id, response) VALUES ('promote-proposal', $1, $2, 'ACCEPTED')`, player, fmt.Sprintf("promote-ticket-%d", index)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
plan := AcceptedMatchPlan{MatchID: "promote-match", ProposalID: "promote-proposal", Region: "EU", Protocol: 1, Players: []MatchPlayer{{PlayerID: "promote-a", Team: 0, Slot: 0}, {PlayerID: "promote-b", Team: 1, Slot: 3}}}
|
|
if err := CreateMatchFromAcceptedProposal(ctx, db, plan, now); err != nil {
|
|
t.Fatalf("promote accepted proposal: %v", err)
|
|
}
|
|
var state string
|
|
if err := db.QueryRowContext(ctx, `SELECT state FROM matches WHERE match_id = 'promote-match'`).Scan(&state); err != nil || state != "ALLOCATING" {
|
|
t.Fatalf("match state=%q err=%v", state, err)
|
|
}
|
|
var acceptedTickets, participantCount int
|
|
if err := db.QueryRowContext(ctx, `SELECT count(*) FROM queue_tickets WHERE ticket_id LIKE 'promote-ticket-%' AND state = 'ACCEPTED'`).Scan(&acceptedTickets); err != nil || acceptedTickets != 2 {
|
|
t.Fatalf("accepted tickets=%d err=%v", acceptedTickets, err)
|
|
}
|
|
if err := db.QueryRowContext(ctx, `SELECT count(*) FROM match_participants WHERE match_id = 'promote-match'`).Scan(&participantCount); err != nil || participantCount != 2 {
|
|
t.Fatalf("participants=%d err=%v", participantCount, err)
|
|
}
|
|
if err := CreateMatchFromAcceptedProposal(ctx, db, plan, now.Add(time.Second)); err != nil {
|
|
t.Fatalf("identical match promotion replay: %v", err)
|
|
}
|
|
conflict := plan
|
|
conflict.Players = append([]MatchPlayer(nil), plan.Players...)
|
|
conflict.Players[1].Slot = 4
|
|
if err := CreateMatchFromAcceptedProposal(ctx, db, conflict, now.Add(2*time.Second)); err == nil {
|
|
t.Fatal("conflicting match promotion replay was accepted")
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLAllocationMatchClaimLeaseAndBindFence(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for index, player := range []string{"allocation-match-a", "allocation-match-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'ACCEPTED', 'build-1', 1, $3, $4)`, fmt.Sprintf("allocation-match-ticket-%d", index), player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO matches (match_id, playlist, state, region, protocol_version) VALUES ('allocation-match', 'casual', 'ALLOCATING', 'EU', 1)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
for index, player := range []string{"allocation-match-a", "allocation-match-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO match_participants (match_id, player_id, ticket_id, slot, team) VALUES ('allocation-match', $1, $2, $3, $4)`, player, fmt.Sprintf("allocation-match-ticket-%d", index), index*3, index); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
claim, found, err := ClaimAllocatingMatch(ctx, db, "enet", now)
|
|
if err != nil || !found || claim.Request != (domain.AllocationRequest{AllocationID: "allocation-allocation-match", MatchID: "allocation-match", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet"}) {
|
|
t.Fatalf("claim=%+v found=%t err=%v", claim, found, err)
|
|
}
|
|
if err := ReleaseAllocatedMatchClaim(ctx, db, claim.Request.MatchID, "different-allocation"); err != domain.ErrConflict {
|
|
t.Fatalf("wrong-claim release err=%v", err)
|
|
}
|
|
if err := ReleaseAllocatedMatchClaim(ctx, db, claim.Request.MatchID, claim.Request.AllocationID); err != nil {
|
|
t.Fatalf("release claim: %v", err)
|
|
}
|
|
reclaimed, found, err := ClaimAllocatingMatch(ctx, db, "enet", now.Add(time.Second))
|
|
if err != nil || !found || reclaimed.Request.AllocationID != claim.Request.AllocationID {
|
|
t.Fatalf("reclaimed=%+v found=%t err=%v", reclaimed, found, err)
|
|
}
|
|
server := domain.ReadyServer{ServerID: "allocation-server", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet", State: domain.ServerReady}
|
|
if err := RegisterReadyServer(ctx, db, server, now); err != nil {
|
|
t.Fatalf("register allocation server: %v", err)
|
|
}
|
|
allocation, err := ClaimAllocation(ctx, db, reclaimed.Request, now.Add(time.Second))
|
|
if err != nil {
|
|
t.Fatalf("record provider allocation: %v", err)
|
|
}
|
|
if err := BindAllocatedMatch(ctx, db, allocation); err != nil {
|
|
t.Fatalf("bind allocation: %v", err)
|
|
}
|
|
var allocatingTickets int
|
|
if err := db.QueryRowContext(ctx, `SELECT count(*) FROM queue_tickets WHERE ticket_id LIKE 'allocation-match-ticket-%' AND state = 'ALLOCATING'`).Scan(&allocatingTickets); err != nil || allocatingTickets != 2 {
|
|
t.Fatalf("allocating tickets=%d err=%v", allocatingTickets, err)
|
|
}
|
|
if _, found, err := ClaimAllocatingMatch(ctx, db, "enet", now.Add(2*time.Second)); err != nil || found {
|
|
t.Fatalf("bound match re-claimed found=%t err=%v", found, err)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLQueueAdapterAgainstRealDatabase(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ('integration-player', 'integration-steam')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
spec := domain.QueueSpec{Playlist: domain.Casual, ClientBuild: "integration-build", ProtocolVersion: 1}
|
|
ticket, err := CreateQueueTicket(ctx, db, "integration-ticket", "integration-player", "integration-create-0001", spec, now)
|
|
if err != nil {
|
|
t.Fatalf("create queue ticket: %v", err)
|
|
}
|
|
if ticket.State != domain.Queued || ticket.Revision != 0 {
|
|
t.Fatalf("unexpected ticket: %+v", ticket)
|
|
}
|
|
replay, err := CreateQueueTicket(ctx, db, "integration-ticket", "integration-player", "integration-create-0001", spec, now.Add(time.Second))
|
|
if err != nil {
|
|
t.Fatalf("idempotent queue replay: %v", err)
|
|
}
|
|
if replay.TicketID != ticket.TicketID || !replay.ExpiresAt.Equal(ticket.ExpiresAt) {
|
|
t.Fatalf("replay changed durable result: %+v vs %+v", replay, ticket)
|
|
}
|
|
if _, err := CreateQueueTicket(ctx, db, "integration-ticket-2", "integration-player", "integration-create-0002", spec, now); err == nil {
|
|
t.Fatal("second active player ticket was accepted")
|
|
}
|
|
if _, err := GetQueueTicket(ctx, db, "integration-player", "integration-ticket", now); err != nil {
|
|
t.Fatalf("owner recovery: %v", err)
|
|
}
|
|
if _, err := GetQueueTicket(ctx, db, "other-player", "integration-ticket", now); err == nil {
|
|
t.Fatal("non-owner recovered queue ticket")
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLQueueHeartbeatAndCancelAreRevisionFenced(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ('heartbeat-player', 'heartbeat-steam')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
spec := domain.QueueSpec{Playlist: domain.Ranked, ClientBuild: "integration-build", ProtocolVersion: 1}
|
|
if _, err := CreateQueueTicket(ctx, db, "heartbeat-ticket", "heartbeat-player", "heartbeat-create-0001", spec, now); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
heartbeat, err := HeartbeatQueueTicket(ctx, db, "heartbeat-player", "heartbeat-ticket", "heartbeat-op-0000001", 0, now.Add(5*time.Second))
|
|
if err != nil {
|
|
t.Fatalf("heartbeat: %v", err)
|
|
}
|
|
if heartbeat.Revision != 1 || !heartbeat.ExpiresAt.Equal(now.Add(35*time.Second)) {
|
|
t.Fatalf("unexpected heartbeat result: %+v", heartbeat)
|
|
}
|
|
if _, err := HeartbeatQueueTicket(ctx, db, "heartbeat-player", "heartbeat-ticket", "heartbeat-op-0000002", 0, now.Add(6*time.Second)); err == nil {
|
|
t.Fatal("stale heartbeat revision was accepted")
|
|
}
|
|
cancelled, err := CancelQueueTicket(ctx, db, "heartbeat-player", "heartbeat-ticket", "heartbeat-op-0000003", 1, now.Add(7*time.Second))
|
|
if err != nil {
|
|
t.Fatalf("cancel: %v", err)
|
|
}
|
|
if cancelled.State != domain.Cancelled || cancelled.Revision != 2 {
|
|
t.Fatalf("unexpected cancellation result: %+v", cancelled)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLConcurrentQueueHeartbeatIsRevisionFencedUnderRealRace is the
|
|
// live counterpart to the sequential stale-heartbeat check above: calling the
|
|
// second heartbeat only after the first has already committed proves the SQL
|
|
// predicate is correct, but not that it actually fences two requests that
|
|
// genuinely overlap at the database. A client can legitimately double-send a
|
|
// heartbeat (a slow response triggering a client-side retry, or two tabs/
|
|
// processes for the same player), and both requests can reach PostgreSQL
|
|
// truly concurrently -- this races that directly.
|
|
func TestPostgreSQLConcurrentQueueHeartbeatIsRevisionFencedUnderRealRace(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ('race-heartbeat-player', 'race-heartbeat-steam')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
spec := domain.QueueSpec{Playlist: domain.Casual, ClientBuild: "integration-build", ProtocolVersion: 1}
|
|
if _, err := CreateQueueTicket(ctx, db, "race-heartbeat-ticket", "race-heartbeat-player", "race-heartbeat-create-01", spec, now); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
const attempts = 5
|
|
var wg sync.WaitGroup
|
|
tickets := make([]domain.QueueTicket, attempts)
|
|
errs := make([]error, attempts)
|
|
wg.Add(attempts)
|
|
for i := 0; i < attempts; i++ {
|
|
go func(i int) {
|
|
defer wg.Done()
|
|
tickets[i], errs[i] = HeartbeatQueueTicket(ctx, db, "race-heartbeat-player", "race-heartbeat-ticket", fmt.Sprintf("race-heartbeat-op-%08d", i), 0, now.Add(time.Duration(i)*time.Millisecond))
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
|
|
won, lost := 0, 0
|
|
for i, err := range errs {
|
|
if err == nil {
|
|
won++
|
|
if tickets[i].Revision != 1 {
|
|
t.Fatalf("winning heartbeat %d landed at revision %d, want 1", i, tickets[i].Revision)
|
|
}
|
|
continue
|
|
}
|
|
lost++
|
|
}
|
|
if won != 1 {
|
|
t.Fatalf("won=%d, want exactly 1 of %d concurrent heartbeats at the same expected revision to win", won, attempts)
|
|
}
|
|
if lost != attempts-1 {
|
|
t.Fatalf("lost=%d, want %d", lost, attempts-1)
|
|
}
|
|
var revision uint64
|
|
if err := db.QueryRow(`SELECT revision FROM queue_tickets WHERE ticket_id = 'race-heartbeat-ticket'`).Scan(&revision); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if revision != 1 {
|
|
t.Fatalf("durable revision = %d, want exactly 1 (a stale winner re-applying would leave it higher)", revision)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLAssignmentPersistenceIsPlayerScopedAndExpiryBound(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"assignment-player", "assignment-other"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ('assignment-ticket', 'assignment-player', 'casual', 'ASSIGNED', 'integration-build', 1, $1, $2)`, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO matches (match_id, playlist, state, region, protocol_version, server_id) VALUES ('assignment-match', 'casual', 'ASSIGNED', 'EU', 1, 'assignment-server')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO match_participants (match_id, player_id, ticket_id, slot, team) VALUES ('assignment-match', 'assignment-player', 'assignment-ticket', 0, 0)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
assignment := DurableAssignment{MatchID: "assignment-match", PlayerID: "assignment-player", AllocationID: "allocation-1", ServerID: "assignment-server", Slot: 0, Region: "EU", ClientBuild: "integration-build", ProtocolVersion: 1, Transport: "enet", Endpoint: "127.0.0.1:7777", JoinAuthorisation: "join-token", ManifestDigest: []byte("manifest"), ExpiresAt: now.Add(time.Minute), Revision: 1}
|
|
if err := SaveAssignment(ctx, db, assignment); err != nil {
|
|
t.Fatalf("save assignment: %v", err)
|
|
}
|
|
got, err := GetAssignment(ctx, db, assignment.PlayerID, assignment.MatchID, now)
|
|
if err != nil {
|
|
t.Fatalf("recover assignment: %v", err)
|
|
}
|
|
if got.JoinAuthorisation != assignment.JoinAuthorisation || got.Slot != assignment.Slot {
|
|
t.Fatalf("assignment changed on round trip: %+v", got)
|
|
}
|
|
if _, err := GetAssignment(ctx, db, "assignment-other", assignment.MatchID, now); err == nil {
|
|
t.Fatal("non-owner recovered assignment")
|
|
}
|
|
if _, err := GetAssignment(ctx, db, assignment.PlayerID, assignment.MatchID, now.Add(2*time.Minute)); err == nil {
|
|
t.Fatal("expired assignment was recovered")
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLProposalClaimAndResponseAreAtomic(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"proposal-player-a", "proposal-player-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
for i, player := range []string{"proposal-player-a", "proposal-player-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'QUEUED', 'integration-build', 1, $3, $4)`, fmt.Sprintf("proposal-ticket-%d", i), player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
proposal, err := domain.NewProposal("proposal-integration", domain.Casual, []string{"proposal-player-a", "proposal-player-b"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := CreateProposal(ctx, db, proposal, map[string]string{"proposal-player-a": "proposal-ticket-0", "proposal-player-b": "proposal-ticket-1"}, now); err != nil {
|
|
t.Fatalf("create proposal: %v", err)
|
|
}
|
|
var proposed int
|
|
if err := db.QueryRow(`SELECT count(*) FROM queue_tickets WHERE state = 'PROPOSED'`).Scan(&proposed); err != nil || proposed != 2 {
|
|
t.Fatalf("proposed queue tickets = %d, err = %v", proposed, err)
|
|
}
|
|
recovered, err := GetProposal(ctx, db, "proposal-player-a", proposal.ProposalID, now)
|
|
if err != nil {
|
|
t.Fatalf("recover proposal: %v", err)
|
|
}
|
|
if len(recovered.Participants) != 2 || recovered.Revision != 0 {
|
|
t.Fatalf("unexpected recovered proposal: %+v", recovered)
|
|
}
|
|
accepted, err := RespondToProposal(ctx, db, "proposal-player-a", proposal.ProposalID, "proposal-response-a-0001", true, 0, now)
|
|
if err != nil {
|
|
t.Fatalf("first proposal acceptance: %v", err)
|
|
}
|
|
if accepted.Revision != 1 || accepted.State != domain.Open {
|
|
t.Fatalf("unexpected first acceptance: %+v", accepted)
|
|
}
|
|
accepted, err = RespondToProposal(ctx, db, "proposal-player-b", proposal.ProposalID, "proposal-response-b-0001", true, 1, now)
|
|
if err != nil {
|
|
t.Fatalf("second proposal acceptance: %v", err)
|
|
}
|
|
if accepted.State != domain.Accepted || accepted.Revision != 2 {
|
|
t.Fatalf("proposal did not close after unanimous acceptance: %+v", accepted)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLProposalDeclineRequeuesEveryParticipant is a real, severe
|
|
// bug this session found by reading the code, not by a failing test: no
|
|
// path anywhere transitioned a PROPOSED ticket back to QUEUED after a
|
|
// decline. A stranded ticket is invisible to the matcher (which only reads
|
|
// state='QUEUED'), still counts as the player's one active ticket (blocking
|
|
// a fresh queue_create), and is renewable forever by an ordinary heartbeat
|
|
// -- a player proposed a match with someone who declines had no way back
|
|
// into matchmaking without realising they had to manually cancel first.
|
|
func TestPostgreSQLProposalDeclineRequeuesEveryParticipant(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"decline-player-a", "decline-player-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
for i, player := range []string{"decline-player-a", "decline-player-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'QUEUED', 'integration-build', 1, $3, $4)`, fmt.Sprintf("decline-ticket-%d", i), player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
proposal, err := domain.NewProposal("decline-proposal", domain.Casual, []string{"decline-player-a", "decline-player-b"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := CreateProposal(ctx, db, proposal, map[string]string{"decline-player-a": "decline-ticket-0", "decline-player-b": "decline-ticket-1"}, now); err != nil {
|
|
t.Fatalf("create proposal: %v", err)
|
|
}
|
|
|
|
// player-a declines; player-b never responded at all -- the bug affects
|
|
// even a participant who was never asked to do anything wrong.
|
|
declined, err := RespondToProposal(ctx, db, "decline-player-a", proposal.ProposalID, "decline-response-a-0001", false, 0, now)
|
|
if err != nil {
|
|
t.Fatalf("decline: %v", err)
|
|
}
|
|
if declined.State != domain.Declined {
|
|
t.Fatalf("proposal did not close on decline: %+v", declined)
|
|
}
|
|
|
|
var stateA, stateB string
|
|
var expiresB time.Time
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = 'decline-ticket-0'`).Scan(&stateA); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT state, expires_at FROM queue_tickets WHERE ticket_id = 'decline-ticket-1'`).Scan(&stateB, &expiresB); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if stateA != "QUEUED" {
|
|
t.Fatalf("decliner's own ticket state = %s, want QUEUED (no cooldown mechanism exists yet to justify leaving it stuck)", stateA)
|
|
}
|
|
if stateB != "QUEUED" {
|
|
t.Fatalf("uninvolved participant's ticket state = %s, want QUEUED -- they must not be stranded by someone else's decline", stateB)
|
|
}
|
|
if !expiresB.After(now) {
|
|
t.Fatalf("requeued ticket expiry %v was not refreshed forward from %v", expiresB, now)
|
|
}
|
|
|
|
// The real, end-to-end regression: both players can be proposed a NEW
|
|
// match instead of ListQueuedCandidates silently never seeing them again.
|
|
candidates, err := ListQueuedCandidates(ctx, db, domain.Casual, now, 10)
|
|
if err != nil {
|
|
t.Fatalf("list queued candidates: %v", err)
|
|
}
|
|
found := map[string]bool{}
|
|
for _, candidate := range candidates {
|
|
found[candidate.PlayerID] = true
|
|
}
|
|
if !found["decline-player-a"] || !found["decline-player-b"] {
|
|
t.Fatalf("requeued players are not visible to the matcher: %+v", candidates)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLProposalTimeoutRequeuesEveryParticipant is the timeout
|
|
// sibling of the decline test above: a proposal that simply times out (no
|
|
// explicit decline, nobody ever responds) hits the exact same
|
|
// ProposalExpireSQL/ProposalParticipantExpireSQL path with the exact same
|
|
// gap -- neither ever touched queue_tickets, so this is the same severe
|
|
// stranding bug reached a different way. Uses GetProposal (the recovery/read
|
|
// path) rather than RespondToProposal, since a real client that just missed
|
|
// the expiry event and comes back later to check on it is exactly the
|
|
// scenario this path exists for.
|
|
func TestPostgreSQLProposalTimeoutRequeuesEveryParticipant(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"timeout-player-a", "timeout-player-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
for i, player := range []string{"timeout-player-a", "timeout-player-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'QUEUED', 'integration-build', 1, $3, $4)`, fmt.Sprintf("timeout-ticket-%d", i), player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
proposal, err := domain.NewProposal("timeout-proposal", domain.Casual, []string{"timeout-player-a", "timeout-player-b"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := CreateProposal(ctx, db, proposal, map[string]string{"timeout-player-a": "timeout-ticket-0", "timeout-player-b": "timeout-ticket-1"}, now); err != nil {
|
|
t.Fatalf("create proposal: %v", err)
|
|
}
|
|
|
|
// Nobody ever responds; recover the proposal well after its 10s window,
|
|
// exactly as a client reconnecting after missing the expiry event would.
|
|
afterExpiry := now.Add(domain.ProposalWindow + time.Second)
|
|
recovered, err := GetProposal(ctx, db, "timeout-player-a", proposal.ProposalID, afterExpiry)
|
|
if err != nil {
|
|
t.Fatalf("recover expired proposal: %v", err)
|
|
}
|
|
if recovered.State != domain.Expired {
|
|
t.Fatalf("proposal did not expire: %+v", recovered)
|
|
}
|
|
|
|
var stateA, stateB string
|
|
var expiresB time.Time
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = 'timeout-ticket-0'`).Scan(&stateA); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT state, expires_at FROM queue_tickets WHERE ticket_id = 'timeout-ticket-1'`).Scan(&stateB, &expiresB); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if stateA != "QUEUED" || stateB != "QUEUED" {
|
|
t.Fatalf("timed-out participants left stranded: a=%s b=%s", stateA, stateB)
|
|
}
|
|
if !expiresB.After(afterExpiry) {
|
|
t.Fatalf("requeued ticket expiry %v was not refreshed forward from %v", expiresB, afterExpiry)
|
|
}
|
|
candidates, err := ListQueuedCandidates(ctx, db, domain.Casual, afterExpiry, 10)
|
|
if err != nil {
|
|
t.Fatalf("list queued candidates: %v", err)
|
|
}
|
|
found := map[string]bool{}
|
|
for _, candidate := range candidates {
|
|
found[candidate.PlayerID] = true
|
|
}
|
|
if !found["timeout-player-a"] || !found["timeout-player-b"] {
|
|
t.Fatalf("requeued players are not visible to the matcher: %+v", candidates)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLCancellingAProposedTicketImmediatelyRequeuesTheOtherParticipant
|
|
// covers the responsiveness gap the decline/timeout fixes above left bounded
|
|
// but not closed: cancelling a ticket that's part of an OPEN proposal used
|
|
// to leave the OTHER participant waiting out the full 10s window for
|
|
// something the system already knew couldn't happen (their proposal partner
|
|
// just walked away). CascadeCancelToOpenProposal declines and requeues that
|
|
// proposal in the same transaction as the cancel itself.
|
|
func TestPostgreSQLCancellingAProposedTicketImmediatelyRequeuesTheOtherParticipant(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"cancel-cascade-a", "cancel-cascade-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
for i, player := range []string{"cancel-cascade-a", "cancel-cascade-b"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'QUEUED', 'integration-build', 1, $3, $4)`, fmt.Sprintf("cancel-cascade-ticket-%d", i), player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
proposal, err := domain.NewProposal("cancel-cascade-proposal", domain.Casual, []string{"cancel-cascade-a", "cancel-cascade-b"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := CreateProposal(ctx, db, proposal, map[string]string{"cancel-cascade-a": "cancel-cascade-ticket-0", "cancel-cascade-b": "cancel-cascade-ticket-1"}, now); err != nil {
|
|
t.Fatalf("create proposal: %v", err)
|
|
}
|
|
|
|
// player-a cancels their own ticket directly, well within the response
|
|
// window -- not a decline, not a timeout, just abandoning the queue.
|
|
// CreateProposal's own QueueTicketProposeSQL already bumped the ticket's
|
|
// revision from 0 to 1, so the cancel's expected revision is 1, not 0.
|
|
cancelled, err := CancelQueueTicket(ctx, db, "cancel-cascade-a", "cancel-cascade-ticket-0", "cancel-cascade-key-0001", 1, now.Add(time.Second))
|
|
if err != nil {
|
|
t.Fatalf("cancel: %v", err)
|
|
}
|
|
if cancelled.State != domain.Cancelled {
|
|
t.Fatalf("ticket did not cancel: %+v", cancelled)
|
|
}
|
|
|
|
var proposalState string
|
|
if err := db.QueryRow(`SELECT state FROM proposals WHERE proposal_id = 'cancel-cascade-proposal'`).Scan(&proposalState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if proposalState != "DECLINED" {
|
|
t.Fatalf("proposal state = %s, want DECLINED immediately, not left OPEN to time out", proposalState)
|
|
}
|
|
var stateB string
|
|
var expiresB time.Time
|
|
if err := db.QueryRow(`SELECT state, expires_at FROM queue_tickets WHERE ticket_id = 'cancel-cascade-ticket-1'`).Scan(&stateB, &expiresB); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if stateB != "QUEUED" {
|
|
t.Fatalf("other participant's ticket state = %s, want QUEUED immediately", stateB)
|
|
}
|
|
if !expiresB.After(now.Add(time.Second)) {
|
|
t.Fatalf("requeued ticket expiry %v was not refreshed forward", expiresB)
|
|
}
|
|
// The cancelling player's own ticket must stay CANCELLED, not get swept
|
|
// back up into the requeue meant for the other participant.
|
|
var stateA string
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = 'cancel-cascade-ticket-0'`).Scan(&stateA); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if stateA != "CANCELLED" {
|
|
t.Fatalf("cancelling player's own ticket state = %s, want it to stay CANCELLED", stateA)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLProposalCreationRollsBackPartialClaims(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ('rollback-player', 'rollback-steam')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ('rollback-ticket', 'rollback-player', 'casual', 'QUEUED', 'integration-build', 1, $1, $2)`, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
proposal, err := domain.NewProposal("rollback-proposal", domain.Casual, []string{"rollback-player", "missing-player"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := CreateProposal(ctx, db, proposal, map[string]string{"rollback-player": "rollback-ticket"}, now); err == nil {
|
|
t.Fatal("proposal with missing ticket mapping was accepted")
|
|
}
|
|
var proposals, participants, proposed int
|
|
if err := db.QueryRow(`SELECT count(*) FROM proposals WHERE proposal_id = 'rollback-proposal'`).Scan(&proposals); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT count(*) FROM proposal_participants WHERE proposal_id = 'rollback-proposal'`).Scan(&participants); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT count(*) FROM queue_tickets WHERE ticket_id = 'rollback-ticket' AND state = 'PROPOSED'`).Scan(&proposed); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if proposals != 0 || participants != 0 || proposed != 0 {
|
|
t.Fatalf("partial proposal claim was not rolled back: proposals=%d participants=%d proposed=%d", proposals, participants, proposed)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLConcurrentProposalCreationClaimsContestedTicketOnce is the
|
|
// live counterpart to TestPostgreSQLProposalCreationRollsBackPartialClaims:
|
|
// every other proposal test in this file (and the whole matcher/allocator
|
|
// suite) runs its transactions strictly one at a time, so none of them can
|
|
// actually exercise the SERIALIZABLE retry-and-fence path CreateProposal
|
|
// relies on -- only two goroutines racing a real connection pool can. Two
|
|
// matchers independently form a proposal that both include the same waiting
|
|
// player's ticket (a real scenario: nothing stops two matcher replicas from
|
|
// reading the same QUEUED ticket in the same poll window); exactly one
|
|
// CreateProposal must win, the other must fail with its whole transaction
|
|
// rolled back, not a database/sql panic, deadlock, or a half-inserted row.
|
|
func TestPostgreSQLConcurrentProposalCreationClaimsContestedTicketOnce(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"race-player-a", "race-player-b", "race-player-c"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
tickets := map[string]string{"race-player-a": "race-ticket-a", "race-player-b": "race-ticket-b", "race-player-c": "race-ticket-c"}
|
|
for player, ticket := range tickets {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', 'QUEUED', 'integration-build', 1, $3, $4)`, ticket, player, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
proposalA, err := domain.NewProposal("race-proposal-a", domain.Casual, []string{"race-player-a", "race-player-b"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
proposalB, err := domain.NewProposal("race-proposal-b", domain.Casual, []string{"race-player-b", "race-player-c"}, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
var wg sync.WaitGroup
|
|
errs := make([]error, 2)
|
|
wg.Add(2)
|
|
go func() {
|
|
defer wg.Done()
|
|
errs[0] = CreateProposal(ctx, db, proposalA, map[string]string{"race-player-a": tickets["race-player-a"], "race-player-b": tickets["race-player-b"]}, now)
|
|
}()
|
|
go func() {
|
|
defer wg.Done()
|
|
errs[1] = CreateProposal(ctx, db, proposalB, map[string]string{"race-player-b": tickets["race-player-b"], "race-player-c": tickets["race-player-c"]}, now)
|
|
}()
|
|
wg.Wait()
|
|
|
|
succeeded := errs[0] == nil
|
|
if succeeded == (errs[1] == nil) {
|
|
t.Fatalf("exactly one contested proposal must win, got errA=%v errB=%v", errs[0], errs[1])
|
|
}
|
|
|
|
winner, loser := "race-proposal-a", "race-proposal-b"
|
|
if !succeeded {
|
|
winner, loser = "race-proposal-b", "race-proposal-a"
|
|
}
|
|
var winnerRows, loserRows, loserParticipants int
|
|
if err := db.QueryRow(`SELECT count(*) FROM proposals WHERE proposal_id = $1`, winner).Scan(&winnerRows); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT count(*) FROM proposals WHERE proposal_id = $1`, loser).Scan(&loserRows); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT count(*) FROM proposal_participants WHERE proposal_id = $1`, loser).Scan(&loserParticipants); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if winnerRows != 1 {
|
|
t.Fatalf("winning proposal %s was not persisted", winner)
|
|
}
|
|
if loserRows != 0 || loserParticipants != 0 {
|
|
t.Fatalf("losing proposal %s was not fully rolled back: proposals=%d participants=%d", loser, loserRows, loserParticipants)
|
|
}
|
|
var contestedState string
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = $1`, tickets["race-player-b"]).Scan(&contestedState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if contestedState != "PROPOSED" {
|
|
t.Fatalf("contested ticket should be claimed by the winner, got state=%s", contestedState)
|
|
}
|
|
// The loser's OWN uncontested ticket (a or c) must have rolled back to
|
|
// QUEUED too -- CreateProposal is one transaction per proposal, so a
|
|
// contested loss on one participant must not leave another participant's
|
|
// ticket stranded as PROPOSED with no surviving proposal to reference it.
|
|
// A (player-a + contested player-b) won iff succeeded, in which case B's
|
|
// own uncontested ticket (player-c) is the one that must have rolled
|
|
// back; if A lost, it's A's own uncontested ticket (player-a) instead.
|
|
loserOnlyTicket := tickets["race-player-c"]
|
|
if !succeeded {
|
|
loserOnlyTicket = tickets["race-player-a"]
|
|
}
|
|
var loserOnlyState string
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = $1`, loserOnlyTicket).Scan(&loserOnlyState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if loserOnlyState != "QUEUED" {
|
|
t.Fatalf("loser's uncontested ticket %s should have rolled back to QUEUED, got %s", loserOnlyTicket, loserOnlyState)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLResultCompletionAndOutboxAreAtomicAndReplayable(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO matches (match_id, playlist, state, region, protocol_version, server_id) VALUES ('result-match', 'casual', 'RESULT_PENDING', 'NA', 1, 'result-server')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
payload := []byte(`{"match_id":"result-match","team0_score":2,"team1_score":1}`)
|
|
digest := sha256.Sum256(payload)
|
|
receipt := domain.ResultReceipt{ResultID: "result-receipt", MatchID: "result-match", ResultNonce: "result-nonce-123456", PayloadDigest: digest, IntegrityState: domain.IntegrityCertified, ReceivedAt: now}
|
|
if err := CompleteResult(ctx, db, receipt, "result-server", "result-event", payload, now); err != nil {
|
|
t.Fatalf("complete result: %v", err)
|
|
}
|
|
var state string
|
|
if err := db.QueryRow(`SELECT state FROM matches WHERE match_id = 'result-match'`).Scan(&state); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if state != "COMPLETED" {
|
|
t.Fatalf("result match state = %s", state)
|
|
}
|
|
events, err := ReadUnpublishedOutbox(ctx, db, 10)
|
|
if err != nil || len(events) != 1 || events[0].EventID != "result-event" {
|
|
t.Fatalf("unpublished result events = %+v, err = %v", events, err)
|
|
}
|
|
if err := MarkOutboxPublished(ctx, db, events[0].EventID, now.Add(time.Second)); err != nil {
|
|
t.Fatalf("ack result event: %v", err)
|
|
}
|
|
if remaining, err := ReadUnpublishedOutbox(ctx, db, 10); err != nil || len(remaining) != 0 {
|
|
t.Fatalf("outbox after ack = %+v, err = %v", remaining, err)
|
|
}
|
|
if err := CompleteResult(ctx, db, receipt, "result-server", "result-event-retry", payload, now.Add(time.Second)); err != nil {
|
|
t.Fatalf("identical completed result replay: %v", err)
|
|
}
|
|
conflict := receipt
|
|
conflict.ResultID = "different-result"
|
|
if err := CompleteResult(ctx, db, conflict, "result-server", "different-event", []byte(`{"conflict":true}`), now.Add(2*time.Second)); err == nil {
|
|
t.Fatal("conflicting completed result was accepted")
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLConcurrentIdenticalResultSubmissionAppliesRatingsExactlyOnce
|
|
// races real concurrent duplicate result submissions -- the scenario behind
|
|
// task 8.25's "identical duplicates idempotent" claim, which every other
|
|
// result test in this file (and the mocked-driver unit tests) only exercises
|
|
// sequentially. A game server can legitimately retry an unacknowledged
|
|
// result POST, and two such retries can land at PostgreSQL genuinely
|
|
// concurrently; every one of them must succeed (this is the identical-replay
|
|
// path, not a conflict), the match must complete exactly once, and -- the
|
|
// part that matters -- the rating update inside applyResultRatings must not
|
|
// run twice just because it raced.
|
|
func TestPostgreSQLConcurrentIdenticalResultSubmissionAppliesRatingsExactlyOnce(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
for _, player := range []string{"result-race-winner", "result-race-loser"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO ratings (player_id, rating, deviation, volatility, ranked_games) VALUES ($1, 1500, 350, 0.06, 0)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO matches (match_id, playlist, state, region, protocol_version, server_id) VALUES ('result-race-match', 'casual', 'RESULT_PENDING', 'NA', 1, 'result-race-server')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ('result-race-ticket-w', 'result-race-winner', 'casual', 'LIVE', 'build-1', 1, $1, $2), ('result-race-ticket-l', 'result-race-loser', 'casual', 'LIVE', 'build-1', 1, $1, $2)`, now, now.Add(time.Minute)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO match_participants (match_id, player_id, ticket_id, slot, team) VALUES ('result-race-match', 'result-race-winner', 'result-race-ticket-w', 0, 0), ('result-race-match', 'result-race-loser', 'result-race-ticket-l', 1, 1)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
result := domain.MatchResult{MatchID: "result-race-match", ServerID: "result-race-server", ResultNonce: "result-race-nonce-123456", Team0Score: 3, Team1Score: 1, IntegrityState: domain.IntegrityCertified}
|
|
digest := domain.ResultDigest(result)
|
|
receipt := domain.ResultReceipt{ResultID: "result-race-receipt", MatchID: result.MatchID, ResultNonce: result.ResultNonce, PayloadDigest: digest, IntegrityState: result.IntegrityState, ReceivedAt: now}
|
|
payload := []byte(`{"match_id":"result-race-match"}`)
|
|
|
|
const attempts = 5
|
|
var wg sync.WaitGroup
|
|
errs := make([]error, attempts)
|
|
wg.Add(attempts)
|
|
for i := 0; i < attempts; i++ {
|
|
go func(i int) {
|
|
defer wg.Done()
|
|
errs[i] = CompleteResultWithResult(ctx, db, receipt, result.ServerID, fmt.Sprintf("result-race-event-%d", i), payload, result, now)
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
for i, err := range errs {
|
|
if err != nil {
|
|
t.Fatalf("identical concurrent submission %d failed: %v", i, err)
|
|
}
|
|
}
|
|
|
|
var state string
|
|
if err := db.QueryRow(`SELECT state FROM matches WHERE match_id = 'result-race-match'`).Scan(&state); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if state != "COMPLETED" {
|
|
t.Fatalf("match state = %s, want COMPLETED", state)
|
|
}
|
|
var winnerGames, loserGames int
|
|
var winnerRating, loserRating float64
|
|
if err := db.QueryRow(`SELECT ranked_games, rating FROM ratings WHERE player_id = 'result-race-winner'`).Scan(&winnerGames, &winnerRating); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT ranked_games, rating FROM ratings WHERE player_id = 'result-race-loser'`).Scan(&loserGames, &loserRating); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Casual results never increment ranked_games by design (rankedIncrement
|
|
// is unconditionally 0 for domain.Casual in applyResultRatings) -- that's
|
|
// not what this test is verifying. What proves "applied exactly once, not
|
|
// N times under the race" is the rating VALUE: a second application would
|
|
// recompute from the already-updated current rating and compound further
|
|
// away from 1500, so an exact match against a single, independently
|
|
// computed application is the assertion that actually falsifies a double
|
|
// application (unlike an inequality check, which a doubled update would
|
|
// still satisfy).
|
|
if winnerGames != 0 || loserGames != 0 {
|
|
t.Fatalf("casual result should never touch ranked_games: winner=%d loser=%d", winnerGames, loserGames)
|
|
}
|
|
baseline := domain.Rating{Value: 1500, RD: 350, Volatility: 0.06}
|
|
winnerOpponents, err := domain.CasualOpponents([]domain.Opponent{{PlayerID: "result-race-loser", Rating: baseline, Score: 1}})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
wantWinner, err := domain.UpdateRating(baseline, winnerOpponents, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
loserOpponents, err := domain.CasualOpponents([]domain.Opponent{{PlayerID: "result-race-winner", Rating: baseline, Score: 0}})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
wantLoser, err := domain.UpdateRating(baseline, loserOpponents, now)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if winnerRating != wantWinner.Value {
|
|
t.Fatalf("winner rating = %v, want exactly %v (a value between these would indicate a partial/compounded update)", winnerRating, wantWinner.Value)
|
|
}
|
|
if loserRating != wantLoser.Value {
|
|
t.Fatalf("loser rating = %v, want exactly %v", loserRating, wantLoser.Value)
|
|
}
|
|
if winnerRating <= loserRating {
|
|
t.Fatalf("winner rating %v should exceed loser rating %v after a certified result", winnerRating, loserRating)
|
|
}
|
|
}
|
|
|
|
// TestPostgreSQLStalledAllocationsAreReclaimedWithoutPenalisingPlayers is the
|
|
// live counterpart to the SQL fragment test: it proves the actual data
|
|
// movement against a real database, not just that the right substrings are
|
|
// present. Two matches: one genuinely stalled (old enough to reclaim), one
|
|
// recent (must survive untouched) -- the deadline boundary and the
|
|
// no-penalty requeue are both meaningless without a real row to check.
|
|
func TestPostgreSQLStalledAllocationsAreReclaimedWithoutPenalisingPlayers(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
ctx := context.Background()
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
stalledCreatedAt := now.Add(-10 * time.Minute)
|
|
recentCreatedAt := now.Add(-5 * time.Second)
|
|
|
|
for _, player := range []string{"stall-player-a", "stall-player-b", "recent-player"} {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ($1, $1)`, player); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
insertTicket := func(ticketID, playerID, state string, expiresAt time.Time) {
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO queue_tickets (ticket_id, player_id, playlist, state, client_build, protocol_version, enqueued_at, expires_at) VALUES ($1, $2, 'casual', $3, 'integration-build', 1, $4, $5)`, ticketID, playerID, state, now, expiresAt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
insertTicket("stall-ticket-a", "stall-player-a", "PROCESS_READY", now.Add(time.Hour))
|
|
insertTicket("stall-ticket-b", "stall-player-b", "PROCESS_READY", now.Add(time.Hour))
|
|
insertTicket("recent-ticket", "recent-player", "ALLOCATING", now.Add(time.Hour))
|
|
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO matches (match_id, playlist, state, region, protocol_version, server_id, created_at) VALUES ('stalled-match', 'casual', 'PROCESS_READY', 'NA', 1, 'stalled-server', $1)`, stalledCreatedAt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO matches (match_id, playlist, state, region, protocol_version, created_at) VALUES ('recent-match', 'casual', 'ALLOCATING', 'NA', 1, $1)`, recentCreatedAt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO match_participants (match_id, player_id, ticket_id, slot, team) VALUES ('stalled-match', 'stall-player-a', 'stall-ticket-a', 0, 0), ('stalled-match', 'stall-player-b', 'stall-ticket-b', 1, 1)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO match_participants (match_id, player_id, ticket_id, slot, team) VALUES ('recent-match', 'recent-player', 'recent-ticket', 0, 0)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
reclaimed, err := ExpireStalledAllocations(ctx, db, now, 2*time.Minute, 10)
|
|
if err != nil {
|
|
t.Fatalf("expire stalled allocations: %v", err)
|
|
}
|
|
if reclaimed != 1 {
|
|
t.Fatalf("reclaimed = %d, want exactly 1 (the recent match must survive)", reclaimed)
|
|
}
|
|
|
|
var stalledState, recentState string
|
|
if err := db.QueryRow(`SELECT state FROM matches WHERE match_id = 'stalled-match'`).Scan(&stalledState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT state FROM matches WHERE match_id = 'recent-match'`).Scan(&recentState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if stalledState != "FAILED" {
|
|
t.Fatalf("stalled match state = %s, want FAILED", stalledState)
|
|
}
|
|
if recentState != "ALLOCATING" {
|
|
t.Fatalf("recent match state = %s, want untouched ALLOCATING", recentState)
|
|
}
|
|
|
|
var ticketAState, ticketBState, recentTicketState string
|
|
var ticketAExpiry time.Time
|
|
if err := db.QueryRow(`SELECT state, expires_at FROM queue_tickets WHERE ticket_id = 'stall-ticket-a'`).Scan(&ticketAState, &ticketAExpiry); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = 'stall-ticket-b'`).Scan(&ticketBState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT state FROM queue_tickets WHERE ticket_id = 'recent-ticket'`).Scan(&recentTicketState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if ticketAState != "QUEUED" || ticketBState != "QUEUED" {
|
|
t.Fatalf("stalled participants' tickets = %s, %s -- want both requeued to QUEUED, not failed/left behind", ticketAState, ticketBState)
|
|
}
|
|
if !ticketAExpiry.After(now) {
|
|
t.Fatalf("requeued ticket expiry %v was not refreshed forward from %v", ticketAExpiry, now)
|
|
}
|
|
if recentTicketState != "ALLOCATING" {
|
|
t.Fatalf("recent match's ticket state = %s, want untouched ALLOCATING", recentTicketState)
|
|
}
|
|
|
|
var activeParticipants int
|
|
if err := db.QueryRow(`SELECT count(*) FROM match_participants WHERE match_id = 'stalled-match' AND participation_active`).Scan(&activeParticipants); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if activeParticipants != 0 {
|
|
t.Fatalf("stalled match still has %d active participants, want 0 (so the player can be matched again)", activeParticipants)
|
|
}
|
|
|
|
// Idempotent: the match is now FAILED, not one of the three reclaimable
|
|
// states, so a second pass must not touch it again.
|
|
reclaimedAgain, err := ExpireStalledAllocations(ctx, db, now.Add(time.Minute), 2*time.Minute, 10)
|
|
if err != nil {
|
|
t.Fatalf("second expire pass: %v", err)
|
|
}
|
|
if reclaimedAgain != 0 {
|
|
t.Fatalf("second pass reclaimed %d matches, want 0 (already-FAILED match must not be reprocessed)", reclaimedAgain)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLRankedSeasonRolloverIsExactlyOnce(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
now := time.Now().UTC().Truncate(time.Microsecond)
|
|
ctx := context.Background()
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO identities (player_id, steam_id) VALUES ('season-player', 'season-steam')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO ratings (player_id, rating, deviation, volatility, ranked_games) VALUES ('season-player', 1900, 100, 0.12, 25)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := db.ExecContext(ctx, `INSERT INTO seasons (season_id, playlist, starts_at, ends_at) VALUES ('season-1', 'ranked', $1, $2)`, now.Add(-12*7*24*time.Hour), now); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
profile := domain.RankedProfile{Rating: domain.Rating{Value: 1900, RD: 100, Volatility: 0.12}, RankedGames: 25}
|
|
updated, applied, err := ApplyRankedSeasonRollover(ctx, db, "season-player", "season-1", profile, now)
|
|
if err != nil || !applied {
|
|
t.Fatalf("first season rollover = %+v applied=%v err=%v", updated, applied, err)
|
|
}
|
|
if updated.Value != 1800 || updated.RD != 200 {
|
|
t.Fatalf("unexpected rolled rating: %+v", updated)
|
|
}
|
|
var rating float64
|
|
var markers int
|
|
if err := db.QueryRow(`SELECT rating FROM ratings WHERE player_id = 'season-player'`).Scan(&rating); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := db.QueryRow(`SELECT count(*) FROM ranked_season_rollovers WHERE player_id = 'season-player' AND season_id = 'season-1'`).Scan(&markers); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rating != 1800 || markers != 1 {
|
|
t.Fatalf("durable rollover state rating=%v markers=%d", rating, markers)
|
|
}
|
|
_, applied, err = ApplyRankedSeasonRollover(ctx, db, "season-player", "season-1", profile, now.Add(time.Second))
|
|
if err != nil || applied {
|
|
t.Fatalf("duplicate season rollover applied=%v err=%v", applied, err)
|
|
}
|
|
if err := db.QueryRow(`SELECT rating FROM ratings WHERE player_id = 'season-player'`).Scan(&rating); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if rating != 1800 {
|
|
t.Fatalf("duplicate rollover changed rating to %v", rating)
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLMigrationsAreForwardExecutable(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
var tableCount int
|
|
if err := db.QueryRow(`SELECT count(*) FROM information_schema.tables WHERE table_schema = 'public' AND table_name = 'assignments'`).Scan(&tableCount); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if tableCount != 1 {
|
|
t.Fatal("assignments migration did not create its table")
|
|
}
|
|
}
|
|
|
|
func TestPostgreSQLMigrationsRollBackAndReapplyCleanly(t *testing.T) {
|
|
db := openIntegrationPostgres(t)
|
|
applyIntegrationMigrations(t, db)
|
|
dir := filepath.Join("..", "migrations")
|
|
tableExists := func(table string) bool {
|
|
var count int
|
|
if err := db.QueryRow(`SELECT count(*) FROM information_schema.tables WHERE table_schema = 'public' AND table_name = $1`, table).Scan(&count); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return count == 1
|
|
}
|
|
if !tableExists("assignments") || !tableExists("allocations") {
|
|
t.Fatal("expected forward-applied schema before rollback")
|
|
}
|
|
|
|
// Roll back every migration one at a time, in reverse, checking each
|
|
// down file actually undoes what its forward file created — not just
|
|
// that Rollback returns nil.
|
|
if err := migrations.Rollback(context.Background(), db, dir, 1); err != nil {
|
|
t.Fatalf("rollback 0006: %v", err)
|
|
}
|
|
var hasAllocationClaimColumn bool
|
|
if err := db.QueryRow(`SELECT count(*) > 0 FROM information_schema.columns WHERE table_name = 'matches' AND column_name = 'allocation_id'`).Scan(&hasAllocationClaimColumn); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if hasAllocationClaimColumn {
|
|
t.Fatal("0006 rollback did not drop matches.allocation_id")
|
|
}
|
|
|
|
if err := migrations.Rollback(context.Background(), db, dir, 4); err != nil {
|
|
t.Fatalf("rollback remaining down to 0001: %v", err)
|
|
}
|
|
if tableExists("assignments") || tableExists("allocations") || tableExists("game_servers") {
|
|
t.Fatal("rollback left later-migration tables behind")
|
|
}
|
|
|
|
if err := migrations.Rollback(context.Background(), db, dir, 1); err != nil {
|
|
t.Fatalf("rollback 0001: %v", err)
|
|
}
|
|
if tableExists("identities") || tableExists("matches") {
|
|
t.Fatal("0001 rollback did not drop its own tables")
|
|
}
|
|
var remaining int
|
|
if err := db.QueryRow(`SELECT count(*) FROM schema_migrations`).Scan(&remaining); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if remaining != 0 {
|
|
t.Fatalf("expected schema_migrations empty after full rollback, got %d rows", remaining)
|
|
}
|
|
|
|
// Reapplying from a fully rolled-back state must reach the same schema,
|
|
// proving down files don't leave orphaned state that trips a forward
|
|
// re-run (e.g. a constraint or index Apply then tries to recreate).
|
|
if err := migrations.Apply(context.Background(), db, dir); err != nil {
|
|
t.Fatalf("reapply after full rollback: %v", err)
|
|
}
|
|
if !tableExists("assignments") || !tableExists("allocations") {
|
|
t.Fatal("reapply after rollback did not recreate the schema")
|
|
}
|
|
}
|