Files
CosmicClash/server/supervisor/supervisor.go
T
Josh Creek 0bad9e07db feat(multiplayer): report process-ready to the control plane from the supervisor
Add opt-in control-plane registration to server/supervisor: once Agones
Ready succeeds, POST /v1/servers/{id}/register (assignment_ready=false)
using a workload token read fresh from disk each call -- matching how a
Kubernetes projected service account token is rotated in place by
kubelet, unlike a cached/env-var secret. ControlPlaneURL empty (the
default) is a total no-op, so direct/Compose mode and allocated-without-
control-plane mode are both byte-for-byte unaffected; New() rejects a
half-configured registration (URL set without token path/server/match/
digest) rather than silently skipping it.

ServerID/MatchID/ImageDigest are read from env vars named by CLI flags
(--server-id-env, --match-id-env, --image-digest-env), matching the
existing --drain-token-env convention in this same binary, rather than
parsed out of the Agones SDK's own GameServer JSON -- that shape isn't
independently verifiable from here, whereas the Kubernetes Downward API
(fieldRef: metadata.name) populating an env var is a standard, safe
pattern already used elsewhere in this codebase for exactly this class
of secret.

A registration failure now kills the child (matching the existing
waitReady failure path) rather than leaving Agones-Ready-but-
control-plane-unregistered process running -- a real gap the second new
test (TestControlPlaneRegistrationFailureKillsChildRatherThanRunningUnregistered)
had to be corrected to actually exercise: its first draft omitted
ReadyURL and was failing at waitReady, before ever reaching the code
path it claimed to test.
2026-09-01 13:27:43 +01:00

380 lines
12 KiB
Go

// Package supervisor contains the small PID-1 lifecycle boundary around an
// allocated Godot process. The Agones client is HTTP-only so local/Compose
// execution remains independent of the cloud SDK.
package supervisor
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net"
"net/http"
"net/url"
"os"
"os/exec"
"strconv"
"strings"
"time"
)
type GameServer struct {
Status struct {
Address string `json:"address"`
Ports []struct {
Name string `json:"name"`
Port int `json:"port"`
} `json:"ports"`
} `json:"status"`
}
type Config struct {
Command []string
Environment []string
SDKBaseURL string
ReadyURL string
Transport string
DrainURL string
DrainToken string
ReadyTimeout time.Duration
PollInterval time.Duration
HTTPClient *http.Client
// ControlPlaneURL, when set, opts into reporting process-ready to the
// matchmaking control plane (multiplayer-next.md task 8.28) once Agones
// Ready succeeds. Leaving it empty preserves every existing behavior
// exactly -- direct/Compose mode and allocated-without-control-plane mode
// are both unaffected. WorkloadTokenPath is read fresh on every call
// rather than cached, matching how a Kubernetes projected service account
// token is rotated in place by kubelet before it expires; ServerID,
// MatchID and ImageDigest are expected to be populated from the pod spec
// (Downward API / mounted build metadata), not guessed at from the
// Agones SDK's own GameServer response.
ControlPlaneURL string
WorkloadTokenPath string
ServerID string
MatchID string
ProtocolVersion int
ImageDigest string
}
type Supervisor struct {
config Config
client *http.Client
cmd *exec.Cmd
}
const DefaultDrainGrace = 285 * time.Second
func New(config Config) (*Supervisor, error) {
if len(config.Command) == 0 || config.Command[0] == "" {
return nil, fmt.Errorf("supervisor command is required")
}
if config.ReadyTimeout <= 0 {
config.ReadyTimeout = 30 * time.Second
}
if config.PollInterval <= 0 {
config.PollInterval = 100 * time.Millisecond
}
if config.Transport == "" {
config.Transport = "enet"
}
if config.Transport != "enet" && config.Transport != "steam_sdr" {
return nil, fmt.Errorf("unsupported transport %q", config.Transport)
}
if config.HTTPClient == nil {
config.HTTPClient = http.DefaultClient
}
if (config.DrainURL == "") != (config.DrainToken == "") {
return nil, fmt.Errorf("drain URL and token must be configured together")
}
if config.DrainURL != "" {
if err := validateLocalDrainURL(config.DrainURL); err != nil {
return nil, err
}
}
if config.ControlPlaneURL != "" && (config.WorkloadTokenPath == "" || config.ServerID == "" || config.MatchID == "" || config.ProtocolVersion < 1 || config.ImageDigest == "") {
return nil, fmt.Errorf("control-plane registration requires a workload token path, server ID, match ID, protocol version and image digest")
}
return &Supervisor{config: config, client: config.HTTPClient}, nil
}
func validateLocalDrainURL(raw string) error {
parsed, err := url.Parse(raw)
if err != nil || (parsed.Scheme != "http" && parsed.Scheme != "https") || parsed.Host == "" || parsed.Path == "" || parsed.User != nil || parsed.RawQuery != "" || parsed.Fragment != "" {
return fmt.Errorf("drain URL must be a loopback HTTP endpoint")
}
host := parsed.Hostname()
if host != "localhost" {
ip := net.ParseIP(host)
if ip == nil || !ip.IsLoopback() {
return fmt.Errorf("drain URL must be a loopback HTTP endpoint")
}
}
return nil
}
// Start launches the process and marks Agones Ready only after the explicit
// readiness probe succeeds. No stdout/log scraping is used. With no SDK URL,
// this is direct/Compose mode and the command is simply started.
func (s *Supervisor) Start(ctx context.Context) error {
env := append([]string(nil), os.Environ()...)
env = append(env, s.config.Environment...)
if s.config.SDKBaseURL != "" {
port, address, err := s.assignedEndpoint(ctx)
if err != nil {
return err
}
if s.config.Transport == "steam_sdr" {
env = append(env, "SDR_LISTEN_PORT="+strconv.Itoa(port), "SDR_IP="+address+":"+strconv.Itoa(port))
}
command := withPort(s.config.Command, port)
s.cmd = exec.CommandContext(ctx, command[0], command[1:]...)
} else {
s.cmd = exec.CommandContext(ctx, s.config.Command[0], s.config.Command[1:]...)
}
s.cmd.Env = env
if err := s.cmd.Start(); err != nil {
return err
}
if s.config.SDKBaseURL == "" {
return nil
}
if err := s.waitReady(ctx); err != nil {
_ = s.cmd.Process.Kill()
return err
}
if err := s.sdkPost(ctx, "/ready"); err != nil {
return err
}
if err := s.registerControlPlane(ctx, false); err != nil {
// Unlike a bare Agones Ready, this failure leaves the match's durable
// control-plane record stuck at ALLOCATING with no way for the
// matcher/allocator to learn this process is actually listening --
// players would wait indefinitely for a server that Agones considers
// healthy. Kill the child so Kubernetes reschedules rather than
// leaving that silent split-brain running.
_ = s.cmd.Process.Kill()
return err
}
return nil
}
// registerControlPlane reports the allocated process's readiness to the
// matchmaking control plane (POST /v1/servers/{id}/register). It is a no-op
// whenever ControlPlaneURL is unset, which is the default and preserves
// every existing direct/Compose/allocated-only behavior exactly. The
// workload token is read fresh from disk on every call rather than cached --
// a Kubernetes projected service account token is rotated in place by
// kubelet before it expires, so caching it risks presenting a stale one on a
// long-lived process.
func (s *Supervisor) registerControlPlane(ctx context.Context, assignmentReady bool) error {
if s.config.ControlPlaneURL == "" {
return nil
}
tokenBytes, err := os.ReadFile(s.config.WorkloadTokenPath)
if err != nil {
return fmt.Errorf("read workload token: %w", err)
}
token := strings.TrimSpace(string(tokenBytes))
if token == "" {
return fmt.Errorf("workload token file %q is empty", s.config.WorkloadTokenPath)
}
body, err := json.Marshal(struct {
MatchID string `json:"match_id"`
ProtocolVersion int `json:"protocol_version"`
ImageDigest string `json:"image_digest"`
AssignmentReady bool `json:"assignment_ready"`
}{s.config.MatchID, s.config.ProtocolVersion, s.config.ImageDigest, assignmentReady})
if err != nil {
return err
}
endpoint := strings.TrimRight(s.config.ControlPlaneURL, "/") + "/v1/servers/" + url.PathEscape(s.config.ServerID) + "/register"
request, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(body))
if err != nil {
return err
}
request.Header.Set("Content-Type", "application/json")
request.Header.Set("Authorization", "Bearer "+token)
// Idempotent per (server, readiness stage): a supervisor restart or a
// dropped response retrying this exact call must replay, not conflict.
request.Header.Set("Idempotency-Key", "supervisor-register-"+s.config.ServerID+"-"+strconv.FormatBool(assignmentReady))
response, err := s.client.Do(request)
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode/100 != 2 {
return fmt.Errorf("control-plane register returned %s", response.Status)
}
return nil
}
func withPort(command []string, port int) []string {
result := append([]string(nil), command...)
for i, arg := range result {
if strings.HasPrefix(arg, "--port=") {
result[i] = "--port=" + strconv.Itoa(port)
return result
}
}
return append(result, "--port="+strconv.Itoa(port))
}
func (s *Supervisor) Wait() error {
if s.cmd == nil {
return fmt.Errorf("supervisor has not started")
}
return s.cmd.Wait()
}
// Run owns the PID-1 termination sequence. The child gets its own context so
// cancellation of the supervisor does not kill it before the authenticated
// drain request has had a chance to stop new admissions. A non-responsive
// child is force-killed after drainGrace; a drain failure is recorded only by
// the returned error if the child exits cleanly, while the deadline still
// prevents a stuck process from hanging termination forever.
func (s *Supervisor) Run(ctx context.Context, drainGrace time.Duration) error {
if s == nil || ctx == nil {
return fmt.Errorf("supervisor context is required")
}
if drainGrace <= 0 {
drainGrace = DefaultDrainGrace
}
processCtx, cancel := context.WithCancel(context.Background())
defer cancel()
if err := s.Start(processCtx); err != nil {
return err
}
wait := make(chan error, 1)
go func() { wait <- s.Wait() }()
select {
case err := <-wait:
return err
case <-ctx.Done():
}
var drainErr error
if s.config.DrainURL != "" {
drainCtx, drainCancel := context.WithTimeout(context.Background(), 5*time.Second)
drainErr = s.Drain(drainCtx)
drainCancel()
}
timer := time.NewTimer(drainGrace)
defer timer.Stop()
select {
case err := <-wait:
if drainErr != nil {
return fmt.Errorf("child exited after drain failure: %w", drainErr)
}
return err
case <-timer.C:
if s.cmd != nil && s.cmd.Process != nil {
_ = s.cmd.Process.Kill()
}
<-wait
if drainErr != nil {
return fmt.Errorf("drain failed and child was force-killed: %w", drainErr)
}
return fmt.Errorf("child force-killed after drain deadline")
}
}
// Drain asks the allocated Godot process to stop accepting new work. The
// token is sent only over the configured localhost control endpoint and is
// never placed in command arguments or logs.
func (s *Supervisor) Drain(ctx context.Context) error {
if s.config.DrainURL == "" || s.config.DrainToken == "" {
return fmt.Errorf("authenticated drain endpoint is required")
}
request, err := http.NewRequestWithContext(ctx, http.MethodPost, s.config.DrainURL, nil)
if err != nil {
return err
}
request.Header.Set("Authorization", "Bearer "+s.config.DrainToken)
response, err := s.client.Do(request)
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode/100 != 2 {
return fmt.Errorf("drain endpoint returned %s", response.Status)
}
return nil
}
func (s *Supervisor) assignedEndpoint(ctx context.Context) (int, string, error) {
var server GameServer
if err := s.sdkGet(ctx, "/gameserver", &server); err != nil {
return 0, "", err
}
if len(server.Status.Ports) == 0 || strings.TrimSpace(server.Status.Address) == "" || strings.ContainsAny(server.Status.Address, " \t\r\n") {
return 0, "", fmt.Errorf("Agones returned no assigned endpoint")
}
for _, port := range server.Status.Ports {
if port.Port > 0 && port.Port <= 65535 && (port.Name == "game" || len(server.Status.Ports) == 1) {
return port.Port, server.Status.Address, nil
}
}
return 0, "", fmt.Errorf("Agones returned no usable game port")
}
func (s *Supervisor) waitReady(ctx context.Context) error {
if s.config.ReadyURL == "" {
return fmt.Errorf("allocated mode requires an explicit readiness URL")
}
deadline := time.NewTimer(s.config.ReadyTimeout)
defer deadline.Stop()
for {
request, err := http.NewRequestWithContext(ctx, http.MethodGet, s.config.ReadyURL, nil)
if err == nil {
response, requestErr := s.client.Do(request)
if requestErr == nil {
_ = response.Body.Close()
if response.StatusCode >= 200 && response.StatusCode < 300 {
return nil
}
}
}
select {
case <-ctx.Done():
return ctx.Err()
case <-deadline.C:
return fmt.Errorf("process-ready probe timed out")
case <-time.After(s.config.PollInterval):
}
}
}
func (s *Supervisor) sdkGet(ctx context.Context, path string, target any) error {
request, err := http.NewRequestWithContext(ctx, http.MethodGet, strings.TrimRight(s.config.SDKBaseURL, "/")+path, nil)
if err != nil {
return err
}
response, err := s.client.Do(request)
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode/100 != 2 {
return fmt.Errorf("Agones GET %s returned %s", path, response.Status)
}
return json.NewDecoder(response.Body).Decode(target)
}
func (s *Supervisor) sdkPost(ctx context.Context, path string) error {
request, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimRight(s.config.SDKBaseURL, "/")+path, nil)
if err != nil {
return err
}
response, err := s.client.Do(request)
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode/100 != 2 {
return fmt.Errorf("Agones POST %s returned %s", path, response.Status)
}
return nil
}