feat(multiplayer): refresh allocator ready servers

This commit is contained in:
Josh Creek
2026-09-01 10:44:43 +01:00
parent 013eb0778b
commit 9dc1cc2d6f
8 changed files with 111 additions and 5 deletions
+62
View File
@@ -52,6 +52,68 @@ type allocationResponse struct {
} `json:"status"`
}
type gameServerListResponse struct {
Items []struct {
Metadata struct {
Name string `json:"name"`
Labels map[string]string `json:"labels"`
} `json:"metadata"`
Status struct {
State string `json:"state"`
} `json:"status"`
} `json:"items"`
}
// ListReadyServers projects only Agones Ready GameServers into the durable
// allocator registry. Compatibility fields must be present as Fleet labels;
// malformed Ready objects fail closed instead of creating selectable capacity.
func (c Client) ListReadyServers(ctx context.Context) ([]domain.ReadyServer, error) {
if c.HTTP == nil {
c.HTTP = http.DefaultClient
}
base, err := c.endpoint()
if err != nil {
return nil, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, base+"/apis/agones.dev/v1/namespaces/"+url.PathEscape(c.Namespace)+"/gameservers", nil)
if err != nil {
return nil, err
}
response, err := c.HTTP.Do(req)
if err != nil {
return nil, err
}
defer response.Body.Close()
if response.StatusCode < 200 || response.StatusCode >= 300 {
return nil, fmt.Errorf("Agones GameServer list returned %s", response.Status)
}
var decoded gameServerListResponse
if err := json.NewDecoder(io.LimitReader(response.Body, 1<<20)).Decode(&decoded); err != nil {
return nil, fmt.Errorf("decode Agones GameServer list: %w", err)
}
ready := make([]domain.ReadyServer, 0, len(decoded.Items))
for _, item := range decoded.Items {
if item.Status.State != "Ready" {
continue
}
server, err := readyServerFromGameServer(item.Metadata.Name, item.Metadata.Labels)
if err != nil {
return nil, err
}
ready = append(ready, server)
}
return ready, nil
}
func readyServerFromGameServer(name string, labels map[string]string) (domain.ReadyServer, error) {
protocol, err := strconv.Atoi(labels["cosmic-clash.io/protocol"])
server := domain.ReadyServer{ServerID: name, Region: labels["cosmic-clash.io/region"], Build: labels["cosmic-clash.io/build"], Protocol: protocol, Transport: labels["cosmic-clash.io/transport"], State: domain.ServerReady}
if err != nil || server.ServerID == "" || (server.Region != "EU" && server.Region != "NA") || server.Build == "" || server.Protocol < 1 || (server.Transport != "enet" && server.Transport != "steam_sdr") {
return domain.ReadyServer{}, fmt.Errorf("invalid Ready GameServer compatibility labels")
}
return server, nil
}
func (c Client) Allocate(ctx context.Context, request domain.AllocationRequest, labels map[string]string, now time.Time) (AllocatedServer, error) {
if c.HTTP == nil {
c.HTTP = http.DefaultClient
+24
View File
@@ -70,3 +70,27 @@ func TestAllocateRejectsUnsafeConfigurationAndProviderFailure(t *testing.T) {
t.Fatalf("provider failure err=%v", err)
}
}
func TestListReadyServersProjectsOnlyStrictReadyFleetMembers(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet || r.URL.Path != "/apis/agones.dev/v1/namespaces/games/gameservers" {
t.Fatalf("request=%s %s", r.Method, r.URL.Path)
}
_, _ = w.Write([]byte(`{"items":[{"metadata":{"name":"ready-a","labels":{"cosmic-clash.io/region":"EU","cosmic-clash.io/build":"build-1","cosmic-clash.io/protocol":"1","cosmic-clash.io/transport":"enet"}},"status":{"state":"Ready"}},{"metadata":{"name":"allocated-a","labels":{}},"status":{"state":"Allocated"}}]}`))
}))
defer server.Close()
ready, err := (Client{BaseURL: server.URL, Namespace: "games", HTTP: server.Client()}).ListReadyServers(context.Background())
if err != nil || len(ready) != 1 || ready[0] != (domain.ReadyServer{ServerID: "ready-a", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet", State: domain.ServerReady}) {
t.Fatalf("ready=%+v err=%v", ready, err)
}
}
func TestListReadyServersFailsClosedOnInvalidReadyCompatibility(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(`{"items":[{"metadata":{"name":"ready-a","labels":{"cosmic-clash.io/region":"EU","cosmic-clash.io/build":"build-1","cosmic-clash.io/protocol":"bad","cosmic-clash.io/transport":"enet"}},"status":{"state":"Ready"}}]}`))
}))
defer server.Close()
if _, err := (Client{BaseURL: server.URL, Namespace: "games", HTTP: server.Client()}).ListReadyServers(context.Background()); err == nil {
t.Fatal("invalid Ready GameServer accepted")
}
}
+12 -1
View File
@@ -45,10 +45,11 @@ func main() {
fatalf("apply migrations: %v", err)
}
now := func() time.Time { return time.Now().UTC() }
client := agones.Client{BaseURL: *agonesURL, Namespace: *namespace}
worker := allocator.Worker{
Claims: store.AllocatingMatchClaims{DB: db, Transport: *transport},
Service: allocator.Service{
Provider: agones.Client{BaseURL: *agonesURL, Namespace: *namespace},
Provider: client,
Durable: store.AllocationRegistry{DB: db},
Now: now,
},
@@ -59,6 +60,16 @@ func main() {
ticker := time.NewTicker(*interval)
defer ticker.Stop()
for {
servers, err := client.ListReadyServers(ctx)
if err != nil && ctx.Err() == nil {
log.Printf("allocator: list Ready GameServers: %v", err)
} else {
for _, server := range servers {
if err := store.RegisterReadyServer(ctx, db, server, now()); err != nil && ctx.Err() == nil {
log.Printf("allocator: register Ready GameServer %s: %v", server.ServerID, err)
}
}
}
if _, err := worker.RunOnce(ctx); err != nil && ctx.Err() == nil {
log.Printf("allocator: run once: %v", err)
}
+2 -1
View File
@@ -16,7 +16,8 @@ const RegisterReadyServerSQL = `INSERT INTO game_servers
VALUES ($1, $2, $3, $4, $5, 'READY', $6)
ON CONFLICT (server_id) DO UPDATE SET region = EXCLUDED.region,
build = EXCLUDED.build, protocol_version = EXCLUDED.protocol_version,
transport = EXCLUDED.transport, state = 'READY', updated_at = EXCLUDED.updated_at`
transport = EXCLUDED.transport, updated_at = EXCLUDED.updated_at
WHERE game_servers.state = 'READY'`
const ClaimReadyServerSQL = `UPDATE game_servers SET state = 'ALLOCATED', updated_at = $5
WHERE server_id = (
+1 -1
View File
@@ -9,7 +9,7 @@ import (
func TestAllocatorSQLClaimsAndAuditsCompatibleReadyServers(t *testing.T) {
for query, fragments := range map[string][]string{
RegisterReadyServerSQL: {"game_servers", "ON CONFLICT", "state = 'READY'"},
RegisterReadyServerSQL: {"game_servers", "ON CONFLICT", "WHERE game_servers.state = 'READY'"},
ClaimReadyServerSQL: {"state = 'READY'", "region = $1", "protocol_version = $3", "FOR UPDATE SKIP LOCKED", "ORDER BY server_id"},
InsertAllocationSQL: {"allocations", "request_digest", "state", "ALLOCATED"},
ProviderServerClaimSQL: {"state = 'READY'", "region = $2", "protocol_version = $4", "RETURNING"},
@@ -71,6 +71,13 @@ func TestPostgreSQLAllocatorClaimReplayAndCapacityFence(t *testing.T) {
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)