Files
CosmicClash/server/cmd/control-plane/main.go
T
Josh Creek 129b0c7ef0 fix(server): fan outbox events out to every control-plane replica
The Deployment runs two replicas, but WebSocket subscribers live only in
each process's in-memory hub. Every replica races to read the same
global unpublished outbox rows, and publishing succeeded even when the
winning replica held no matching local subscriber -- that replica then
set the single global published_at. A client connected to the other
replica never received the event, and delivery degraded further with
each replica added. REST recovery eventually converged, but short-lived
proposal transitions could be observed late or not at all.

Publish committed events through PostgreSQL LISTEN/NOTIFY so the replica
that owns the subscriber's connection delivers it, regardless of which
replica drained the row. The listener holds its own pgx connection --
LISTEN is session state, so a pooled database/sql connection cannot
carry it -- and reconnects with backoff, since losing it would silently
downgrade that replica's subscribers to REST-only recovery.

The fan-out is optional: without EventFanout configured, behaviour is
unchanged local-hub publication, which stays correct for a single
replica and for tests. Only outbox-sourced events are routed through it;
the in-request-path publishes remain local, as those are a latency
optimisation for the caller's own connection.

Fan-out needs a wire shape of its own because ControlPlaneEvent hides
PlayerID from clients, and the recipient is exactly what a peer replica
needs to route on.
2026-09-05 10:29:23 +01:00

204 lines
8.4 KiB
Go

package main
import (
"context"
"database/sql"
"flag"
"fmt"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/cosmic-clash/cosmic-clash/server/api"
"github.com/cosmic-clash/cosmic-clash/server/domain"
"github.com/cosmic-clash/cosmic-clash/server/migrations"
"github.com/cosmic-clash/cosmic-clash/server/observability"
"github.com/cosmic-clash/cosmic-clash/server/store"
_ "github.com/jackc/pgx/v5/stdlib"
"github.com/redis/go-redis/v9"
)
func main() {
listen := flag.String("listen", ":8080", "HTTP listen address")
role := flag.String("role", "api", "control-plane role; currently api")
dsn := flag.String("dsn", os.Getenv("COSMIC_CLASH_POSTGRES_DSN"), "PostgreSQL connection string")
migrationDir := flag.String("migrations", "migrations", "directory containing numbered SQL migrations")
redisAddr := flag.String("redis-addr", os.Getenv("COSMIC_CLASH_REDIS_ADDR"), "optional Redis address for the candidate projection")
redisPrefix := flag.String("redis-prefix", envOrDefault("COSMIC_CLASH_REDIS_PREFIX", "cosmic-clash"), "Redis key prefix")
redisTTL := flag.Duration("redis-ttl", 60*time.Second, "TTL for transient candidate projection entries")
workloadSecret := flag.String("workload-secret", os.Getenv("COSMIC_CLASH_WORKLOAD_SECRET"), "HMAC secret for control-plane-issued workload tokens (see workload/signed_token.go); server registration/result submission return 503 until this is set")
degraded := flag.Bool("degraded", false, "start with new login, queue, and proposal mutations rejected; SIGUSR1 enables and SIGUSR2 disables this mode")
rateLimit := flag.Int("rate-limit", 120, "maximum requests per per-credential/IP fixed window")
rateWindow := flag.Duration("rate-limit-window", time.Minute, "fixed window for the per-replica request limiter")
rateMaxKeys := flag.Int("rate-limit-max-keys", 10000, "maximum credential/IP keys retained by the per-replica request limiter")
trustedProxyCIDRs := flag.String("trusted-proxy-cidrs", os.Getenv("COSMIC_CLASH_TRUSTED_PROXY_CIDRS"), "comma-separated immediate proxy CIDRs allowed to supply X-Forwarded-For")
minProtocolVersion := flag.Int("min-protocol-version", 0, "reject queue_create below this protocol_version with 426 Upgrade Required instead of queueing a client the matcher can never pair with anyone; zero disables the floor")
flag.Parse()
if *role != "api" {
fatalf("unsupported role %q (only api is implemented)", *role)
}
if *dsn == "" {
fatalf("--dsn or COSMIC_CLASH_POSTGRES_DSN is required")
}
if *redisTTL <= 0 {
fatalf("--redis-ttl must be positive")
}
if *minProtocolVersion < 0 {
fatalf("--min-protocol-version must be non-negative")
}
rateLimiter, err := api.NewRateLimiter(*rateLimit, *rateWindow, *rateMaxKeys)
if err != nil {
fatalf("invalid request limiter configuration: %v", err)
}
clientIPs, err := api.NewClientIPResolver(*trustedProxyCIDRs)
if err != nil {
fatalf("invalid trusted proxy configuration: %v", err)
}
db, err := sql.Open("pgx", *dsn)
if err != nil {
fatalf("open PostgreSQL: %v", err)
}
defer db.Close()
startupCtx, startupCancel := context.WithTimeout(context.Background(), 30*time.Second)
defer startupCancel()
if err := db.PingContext(startupCtx); err != nil {
fatalf("ping PostgreSQL: %v", err)
}
if err := migrations.Apply(startupCtx, db, *migrationDir); err != nil {
fatalf("apply migrations: %v", err)
}
var candidateIndex api.CandidateIndex
var redisClient *redis.Client
if *redisAddr != "" {
redisClient = redis.NewClient(&redis.Options{Addr: *redisAddr})
defer redisClient.Close()
candidateIndex = store.RedisCandidateIndex{Client: redisClient, Prefix: *redisPrefix, TTL: *redisTTL}
}
if *workloadSecret == "" {
fmt.Fprintln(os.Stderr, "control-plane: warning: --workload-secret / COSMIC_CLASH_WORKLOAD_SECRET is unset; server registration and result submission will return 503")
}
service := newAPIService(db, *workloadSecret, candidateIndex)
service.RateLimiter = rateLimiter
service.ClientIPs = clientIPs
service.MinProtocolVersion = *minProtocolVersion
admission := api.NewAdmissionGate(*degraded)
service.Admission = admission
server := &http.Server{Addr: *listen, Handler: service.Handler(), ReadHeaderTimeout: 5 * time.Second}
serveErr := make(chan error, 1)
go func() { serveErr <- server.ListenAndServe() }()
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
operatorSignals := make(chan os.Signal, 2)
signal.Notify(operatorSignals, syscall.SIGUSR1, syscall.SIGUSR2)
defer signal.Stop(operatorSignals)
go func() {
for sig := range operatorSignals {
switch sig {
case syscall.SIGUSR1:
admission.SetDegraded(true)
fmt.Fprintln(os.Stderr, "control-plane: degraded admission enabled")
case syscall.SIGUSR2:
admission.SetDegraded(false)
fmt.Fprintln(os.Stderr, "control-plane: degraded admission disabled")
}
}
}()
// Fan committed outbox events out to every replica. Subscribers live in
// each process's in-memory hub, but any replica may drain a given outbox
// row, so without this a client connected elsewhere never sees the event
// and delivery degrades as replicas are added.
service.EventFanout = func(event api.ControlPlaneEvent) error {
payload, err := api.EncodeFannedOutEvent(event)
if err != nil {
return err
}
return store.NotifyControlPlaneEvent(ctx, db, payload)
}
go store.ListenControlPlaneEvents(ctx, *dsn, func(payload []byte) {
event, err := api.DecodeFannedOutEvent(payload)
if err != nil {
return
}
// Publishing to a player with no local subscriber is a no-op, so every
// replica can handle every notification.
_ = service.PublishControlPlaneEvent(event)
}, func(err error) {
fmt.Fprintf(os.Stderr, "control-plane: event fan-out listener: %v\n", err)
})
go api.RunProposalOutboxDispatcher(ctx, db, service)
go api.RunResultOutboxDispatcher(ctx, db, service)
go api.RunStateOutboxDispatcher(ctx, db, service)
select {
case err := <-serveErr:
if err != nil && err != http.ErrServerClosed {
fatalf("serve API: %v", err)
}
case <-ctx.Done():
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second)
defer shutdownCancel()
if err := server.Shutdown(shutdownCtx); err != nil {
fatalf("shutdown API: %v", err)
}
}
}
func newAPIHandler(db *sql.DB, workloadSecret string, indexes ...api.CandidateIndex) http.Handler {
return newAPIService(db, workloadSecret, indexes...).Handler()
}
func newAPIService(db *sql.DB, workloadSecret string, indexes ...api.CandidateIndex) *api.Service {
var candidateIndex api.CandidateIndex
if len(indexes) > 0 {
candidateIndex = indexes[0]
}
return &api.Service{
SessionBackend: store.PostgresSessions{DB: db},
SessionIssuer: store.PostgresSessions{DB: db},
QueueBackend: store.PostgresQueue{DB: db},
ProposalBackend: api.ProposalProviderFromStore(db),
ProposalPromoter: api.ProposalPromoterFromStore(db),
ServerRegistrar: api.ServerRegistrarFromStore(db),
ServerShutdowner: api.ServerShutdownerFromStore(db),
ServerConnections: api.ServerConnectionsFromStore(db),
ResultSubmitter: store.PostgresResults{DB: db},
RankedProfileProvider: store.PostgresRankedProfiles{DB: db},
TierPolicy: domain.DefaultTierPolicy(),
Assignment: api.AssignmentProviderFromStore(db),
Roster: func(ctx context.Context, binding domain.WorkloadBinding, now time.Time) ([][]byte, error) {
return store.GetAssignmentRoster(ctx, db, binding.MatchID, binding.ServerID, now)
},
CandidateIndex: candidateIndex,
ProbeRecorder: store.PostgresQueue{DB: db},
WorkloadVerify: api.WorkloadVerifierFromSignedToken([]byte(workloadSecret), db),
ReadinessCheck: db.PingContext,
Now: func() time.Time { return time.Now().UTC() },
Log: logEvent,
Metrics: observability.NewMetrics(),
}
}
// logEvent writes one credential-safe structured event per line to stderr.
// Best-effort: a logging failure must never fail or block the request it
// describes, so encode errors are swallowed rather than surfaced.
func logEvent(event observability.Event) {
payload, err := observability.Encode(event)
if err != nil {
return
}
fmt.Fprintln(os.Stderr, string(payload))
}
func envOrDefault(name, fallback string) string {
if value := os.Getenv(name); value != "" {
return value
}
return fallback
}
func fatalf(format string, args ...any) {
fmt.Fprintf(os.Stderr, "control-plane: "+format+"\n", args...)
os.Exit(1)
}