From 0023bdab6eba711905606117bc645128f8715888 Mon Sep 17 00:00:00 2001 From: Josh Creek <8179928+jcreek@users.noreply.github.com> Date: Tue, 1 Sep 2026 09:51:52 +0100 Subject: [PATCH] feat: reconcile Agones allocations durably --- multiplayer-next.md | 4 ++- multiplayer-todo.md | 2 +- server/allocator/service.go | 46 ++++++++++++++++++++++++ server/allocator/service_test.go | 54 ++++++++++++++++++++++++++++ server/store/allocator_sql.go | 61 +++++++++++++++++++++++++++++++- server/store/queue_sql.go | 4 +++ 6 files changed, 168 insertions(+), 3 deletions(-) create mode 100644 server/allocator/service.go create mode 100644 server/allocator/service_test.go diff --git a/multiplayer-next.md b/multiplayer-next.md index 84b0dd4c..6c54fcba 100644 --- a/multiplayer-next.md +++ b/multiplayer-next.md @@ -57,7 +57,9 @@ product policy are in [`docs/MATCHMAKING.md`](docs/MATCHMAKING.md). projections and atomically claims compatible capacity with replay/conflict fencing; `server/agones` now submits and validates namespaced `GameServerAllocation` responses, including dynamic address/port data; - provider-to-durable claim reconciliation and assignment publication remain. + `server/allocator` now requires provider allocation reconciliation into the + durable registry before returning an endpoint; signed roster publication and + live Agones integration remain. - [ ] **IN PROGRESS:** Run the Go control plane against PostgreSQL/Redis with independently runnable API, matcher, allocator and maintenance roles. The `cmd/control-plane` API role now opens PostgreSQL, applies migrations, wires diff --git a/multiplayer-todo.md b/multiplayer-todo.md index 9b604343..e6504753 100644 --- a/multiplayer-todo.md +++ b/multiplayer-todo.md @@ -1213,7 +1213,7 @@ the local/CI/community transport, not a silent production fallback. | 8.27 `[D:8.26]` | **IN PROGRESS.** Go supervisor package provides local-safe Agones REST discovery, validates assigned address/port data, injects dynamic `SDR_LISTEN_PORT`/`SDR_IP`, performs explicit process-ready probing and Ready transition; direct mode bypasses Agones | `server/supervisor/` covers allocated/direct startup, invalid endpoint rejection, dynamic endpoint/Ready ordering and authenticated drain; allocated Godot now supplies a loopback readiness/drain control surface and `agones_sdk.gd` supplies sidecar Health/Ready/Shutdown/annotation REST operations; metadata watch, real Agones annotation/shutdown confirmation and emulator integration remain | | 8.28 `[D:8.6,8.27]` | **IN PROGRESS.** Supervisor separates explicit process-ready from Agones Ready and never scrapes stdout; allocated mode refuses to mark Ready without a configured readiness probe | `server/supervisor/` tests prove Ready follows the probe and direct mode remains functional; `server_control.gd`, `agones_sdk.gd` and process-level smokes prove loopback `/ready`, `/health`, bearer-protected `/drain`, sidecar-shaped Health/Ready calls and drain admission fencing; detached-container and Health-reclaim integration remain | | 8.29 `[D:8.26,8.27]` | **IN PROGRESS.** Supervisor discovers and validates the Agones endpoint, propagates the actual dynamic `--port`, and exports `SDR_LISTEN_PORT`/`SDR_IP` only for Hosted-SDR while preserving an isolated ENet path | `server/supervisor/` tests cover invalid address/port rejection, dynamic port argument/env propagation and SDR-vs-ENet separation; real Agones dynamic/passthrough mapping, POP/cert/firewall/NAT and multi-match fixture remain | -| 8.30 `[D:8.18,8.26,8.28,8.29]` | **IN PROGRESS.** Pure Go allocator filters Ready servers by region/build/protocol/transport, atomically claims one with idempotent allocation replay, and now owns the assignment-publication boundary; PostgreSQL adds durable GameServer registration and compatible `SKIP LOCKED` claims with request-digest fencing; `server/agones` submits and validates namespaced `GameServerAllocation` responses and dynamic endpoints | `server/domain/allocator.go`, `server/store/allocator_sql.go`, `server/agones/allocation.go`, `server/migrations/0004_allocator_registry.sql` and tests cover deterministic compatible selection, exhaustion, conflicting/identical allocation replay, unknown allocations, SQL claim ordering, invalid input, provider error/malformed response/IPv6 endpoint handling and assignment replay/conflict; `TestPostgreSQLAllocatorClaimReplayAndCapacityFence` now covers live registration/selection/replay/conflict/no-capacity when the disposable database gate is run; provider-to-durable claim reconciliation, signed roster metadata, bounded cross-replica retry and live integration remain | +| 8.30 `[D:8.18,8.26,8.28,8.29]` | **IN PROGRESS.** Pure Go allocator filters Ready servers by region/build/protocol/transport, atomically claims one with idempotent allocation replay, and now owns the assignment-publication boundary; PostgreSQL adds durable GameServer registration and compatible `SKIP LOCKED` claims with request-digest fencing; `server/agones` submits and validates namespaced `GameServerAllocation` responses and dynamic endpoints; `server/allocator` reconciles provider success into durable state before exposing the endpoint | `server/domain/allocator.go`, `server/store/allocator_sql.go`, `server/agones/allocation.go`, `server/allocator/service.go`, `server/migrations/0004_allocator_registry.sql` and tests cover deterministic compatible selection, exhaustion, conflicting/identical allocation replay, unknown allocations, SQL claim ordering, invalid input, provider error/malformed response/IPv6 endpoint handling, durable-reconciliation failure isolation and assignment replay/conflict; `TestPostgreSQLAllocatorClaimReplayAndCapacityFence` now covers live registration/selection/replay/conflict/no-capacity when the disposable database gate is run; signed roster metadata, bounded cross-replica retry and live integration remain | | 8.31 `[D:8.9,8.30]` | **IN PROGRESS.** Pure Go assignment gate requires Allocated state, exact allocation ID/match/server/region/build/protocol/transport compatibility, non-empty hosted endpoint and verified manifest signature before exposure; allocator publication cannot expose Ready state | `server/domain/assignment.go` and `allocator.go` plus adversarial fixtures cover early-connect, tampered signature/manifest, wrong compatibility, empty endpoint, unknown allocation and post-publication mutation rejection; Agones metadata watch, hosted-address registration, production signer and client-ticket publication remain | | 8.32 `[D:8.2,8.26,8.30]` | **IN PROGRESS.** Provider-neutral FleetAutoscaler baseline preserves a two-process Ready buffer, caps warm capacity, and leaves Allocated scale-down independent of the Ready floor; Fleet image references remain digest-pinned for current/rollback pre-pull | `deploy/k8s/base/fleet-autoscaler.yaml` and manifest tests cover Fleet ownership, Buffer policy and floor/cap invariants; regional on-demand node pools/failure domains, pre-pull rollout, warm-allocation p95/p99 and N+1 certification remain | | 8.33 `[D:8.26,8.32]` | **IN PROGRESS.** Fleet scheduling now requires on-demand capacity and spreads Ready processes across zones with skew 1; the autoscaler preserves the two-process Ready floor | `deploy/k8s/base/fleet.yaml` and manifest tests reject interruptible placement and single-zone concentration structurally; regional node pools, forced node-loss testing and measured N+1 headroom remain | diff --git a/server/allocator/service.go b/server/allocator/service.go new file mode 100644 index 00000000..fe96a5c1 --- /dev/null +++ b/server/allocator/service.go @@ -0,0 +1,46 @@ +// Package allocator coordinates provider allocation with durable control-plane +// state. It does not expose an endpoint until both boundaries succeed. +package allocator + +import ( + "context" + "time" + + "github.com/cosmic-clash/cosmic-clash/server/agones" + "github.com/cosmic-clash/cosmic-clash/server/domain" +) + +type Provider interface { + Allocate(context.Context, domain.AllocationRequest, map[string]string, time.Time) (agones.AllocatedServer, error) +} + +type Durable interface { + RecordProviderAllocation(context.Context, domain.Allocation, time.Time) (domain.Allocation, error) +} + +type Service struct { + Provider Provider + Durable Durable + Now func() time.Time +} + +func (s Service) Allocate(ctx context.Context, request domain.AllocationRequest, labels map[string]string) (agones.AllocatedServer, error) { + if s.Provider == nil || s.Durable == nil || s.Now == nil { + return agones.AllocatedServer{}, errNotConfigured + } + now := s.Now() + result, err := s.Provider.Allocate(ctx, request, labels, now) + if err != nil { + return agones.AllocatedServer{}, err + } + if _, err := s.Durable.RecordProviderAllocation(ctx, result.Allocation, now); err != nil { + return agones.AllocatedServer{}, err + } + return result, nil +} + +var errNotConfigured = &configurationError{} + +type configurationError struct{} + +func (*configurationError) Error() string { return "allocator service is not configured" } diff --git a/server/allocator/service_test.go b/server/allocator/service_test.go new file mode 100644 index 00000000..1e1f4aa9 --- /dev/null +++ b/server/allocator/service_test.go @@ -0,0 +1,54 @@ +package allocator + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/cosmic-clash/cosmic-clash/server/agones" + "github.com/cosmic-clash/cosmic-clash/server/domain" +) + +type providerSpy struct { + calls int + result agones.AllocatedServer + err error +} + +func (p *providerSpy) Allocate(_ context.Context, _ domain.AllocationRequest, _ map[string]string, _ time.Time) (agones.AllocatedServer, error) { + p.calls++ + return p.result, p.err +} + +type durableSpy struct { + calls int + allocation domain.Allocation + err error +} + +func (d *durableSpy) RecordProviderAllocation(_ context.Context, allocation domain.Allocation, _ time.Time) (domain.Allocation, error) { + d.calls++ + d.allocation = allocation + return allocation, d.err +} + +func TestServiceDurablyRecordsProviderAllocationBeforeReturning(t *testing.T) { + provider := &providerSpy{result: agones.AllocatedServer{Allocation: domain.Allocation{AllocationID: "a", MatchID: "m", ServerID: "gs", State: domain.ServerAllocated}, Endpoint: "127.0.0.1:7777"}} + durable := &durableSpy{} + service := Service{Provider: provider, Durable: durable, Now: func() time.Time { return time.Unix(1000, 0) }} + result, err := service.Allocate(context.Background(), domain.AllocationRequest{AllocationID: "a", MatchID: "m", Region: "EU", Build: "b", Protocol: 1, Transport: "enet"}, map[string]string{"region": "EU"}) + if err != nil || result.Endpoint == "" || durable.calls != 1 || durable.allocation.ServerID != "gs" { + t.Fatalf("result=%+v err=%v durable=%+v", result, err, durable) + } +} + +func TestServiceDoesNotReturnProviderResultAfterDurableFailure(t *testing.T) { + provider := &providerSpy{result: agones.AllocatedServer{Allocation: domain.Allocation{AllocationID: "a", State: domain.ServerAllocated}, Endpoint: "127.0.0.1:7777"}} + durable := &durableSpy{err: errors.New("database unavailable")} + service := Service{Provider: provider, Durable: durable, Now: func() time.Time { return time.Unix(1000, 0) }} + result, err := service.Allocate(context.Background(), domain.AllocationRequest{AllocationID: "a", MatchID: "m", Region: "EU", Build: "b", Protocol: 1, Transport: "enet"}, map[string]string{"region": "EU"}) + if err == nil || result.Endpoint != "" || durable.calls != 1 { + t.Fatalf("result=%+v err=%v calls=%d", result, err, durable.calls) + } +} diff --git a/server/store/allocator_sql.go b/server/store/allocator_sql.go index 5df43249..a4a93696 100644 --- a/server/store/allocator_sql.go +++ b/server/store/allocator_sql.go @@ -37,6 +37,14 @@ const SelectAllocationSQL = `SELECT allocation_id, match_id, server_id, region, protocol_version, transport, allocated_at, request_digest FROM allocations WHERE allocation_id = $1` +const ProviderServerClaimSQL = `UPDATE game_servers SET state = 'ALLOCATED', updated_at = $6 +WHERE server_id = $1 AND state = 'READY' AND region = $2 AND build = $3 + AND protocol_version = $4 AND transport = $5 +RETURNING server_id` + +const ServerAllocationConflictSQL = `SELECT allocation_id FROM allocations +WHERE server_id = $1 FOR UPDATE` + func RegisterReadyServer(ctx context.Context, db *sql.DB, server domain.ReadyServer, now time.Time) error { if db == nil || server.ServerID == "" || (server.Region != "EU" && server.Region != "NA") || server.Build == "" || server.Protocol <= 0 || (server.Transport != "enet" && server.Transport != "steam_sdr") || server.State != domain.ServerReady || now.IsZero() { return fmt.Errorf("invalid ready server registration") @@ -46,7 +54,7 @@ func RegisterReadyServer(ctx context.Context, db *sql.DB, server domain.ReadySer } func ClaimAllocation(ctx context.Context, db *sql.DB, request domain.AllocationRequest, now time.Time) (domain.Allocation, error) { - if db == nil || request.AllocationID == "" || request.MatchID == "" || (request.Region != "EU" && request.Region != "NA") || request.Build == "" || request.Protocol <= 0 || (request.Transport != "enet" && request.Transport != "steam_sdr") || now.IsZero() { + if !validAllocationInput(db, request, now) { return domain.Allocation{}, domain.ErrAllocationInput } digest := allocationRequestDigest(request) @@ -80,6 +88,57 @@ func ClaimAllocation(ctx context.Context, db *sql.DB, request domain.AllocationR return allocation, err } +// RecordProviderAllocation reconciles a provider-side Agones claim with the +// durable registry. It is deliberately separate from ClaimAllocation because +// Agones has already selected the server; no client-facing assignment may use +// the result until this exact tuple is durably recorded. +func RecordProviderAllocation(ctx context.Context, db *sql.DB, allocation domain.Allocation, now time.Time) (domain.Allocation, error) { + request := domain.AllocationRequest{AllocationID: allocation.AllocationID, MatchID: allocation.MatchID, Region: allocation.Region, Build: allocation.Build, Protocol: allocation.Protocol, Transport: allocation.Transport} + if !validAllocationInput(db, request, now) || allocation.State != domain.ServerAllocated || allocation.ServerID == "" { + return domain.Allocation{}, domain.ErrAllocationInput + } + digest := allocationRequestDigest(request) + var recorded domain.Allocation + err := RunSerializable(ctx, db, DefaultSerializableAttempts, func(ctx context.Context, tx *sql.Tx) error { + var prior domain.Allocation + var priorDigest []byte + err := tx.QueryRowContext(ctx, SelectAllocationSQL, allocation.AllocationID).Scan(&prior.AllocationID, &prior.MatchID, &prior.ServerID, &prior.Region, &prior.Build, &prior.Protocol, &prior.Transport, &prior.AllocatedAt, &priorDigest) + if err == nil { + if !bytes.Equal(priorDigest, digest[:]) || prior.ServerID != allocation.ServerID { + return domain.ErrConflict + } + recorded = prior + recorded.State = domain.ServerAllocated + return nil + } + if err != sql.ErrNoRows { + return err + } + var existing string + if err := tx.QueryRowContext(ctx, ServerAllocationConflictSQL, allocation.ServerID).Scan(&existing); err == nil { + return domain.ErrConflict + } else if err != sql.ErrNoRows { + return err + } + var serverID string + if err := tx.QueryRowContext(ctx, ProviderServerClaimSQL, allocation.ServerID, allocation.Region, allocation.Build, allocation.Protocol, allocation.Transport, now).Scan(&serverID); err != nil { + if err == sql.ErrNoRows { + return domain.ErrNoCapacity + } + return err + } + recorded = allocation + recorded.AllocatedAt = now + _, err = tx.ExecContext(ctx, InsertAllocationSQL, allocation.AllocationID, allocation.MatchID, serverID, allocation.Region, allocation.Build, allocation.Protocol, allocation.Transport, digest[:], now) + return err + }) + return recorded, err +} + +func validAllocationInput(db *sql.DB, request domain.AllocationRequest, now time.Time) bool { + return db != nil && request.AllocationID != "" && request.MatchID != "" && (request.Region == "EU" || request.Region == "NA") && request.Build != "" && request.Protocol > 0 && (request.Transport == "enet" || request.Transport == "steam_sdr") && !now.IsZero() +} + func allocationRequestDigest(request domain.AllocationRequest) [32]byte { return sha256.Sum256([]byte(fmt.Sprintf("%s\x00%s\x00%s\x00%s\x00%d\x00%s", request.AllocationID, request.MatchID, request.Region, request.Build, request.Protocol, request.Transport))) } diff --git a/server/store/queue_sql.go b/server/store/queue_sql.go index 4d0b51c3..f0d3a67d 100644 --- a/server/store/queue_sql.go +++ b/server/store/queue_sql.go @@ -137,6 +137,10 @@ type queueTicketRecord struct { type PostgresQueue struct{ DB *sql.DB } +func (q PostgresQueue) RecordProviderAllocation(ctx context.Context, allocation domain.Allocation, now time.Time) (domain.Allocation, error) { + return RecordProviderAllocation(ctx, q.DB, allocation, now) +} + const QueueProbeRecordSQL = `UPDATE queue_tickets SET predicted_rtt = jsonb_set(COALESCE(predicted_rtt, '{}'::jsonb), ARRAY[$2], to_jsonb($3::double precision), true) WHERE player_id = $1 AND state IN ('QUEUED', 'PROPOSED') AND expires_at > $4`