diff --git a/cmd/felis/operator.go b/cmd/felis/operator.go index 7bff9f3..c1e57f9 100644 --- a/cmd/felis/operator.go +++ b/cmd/felis/operator.go @@ -107,6 +107,9 @@ func cmdOperator(args []string, _, stderr io.Writer) int { // injects into user servers. The Deployment passes it as FELIS_IMAGE (see // platform.OperatorDeployment); absent, that injection is simply skipped. FelisImage: os.Getenv("FELIS_IMAGE"), + // Uncached: the maintenance-lock check lists Jobs only when a server is + // about to start, which does not justify a namespace-wide Job informer. + Jobs: mgr.GetAPIReader(), } if err := r.SetupWithManager(mgr); err != nil { fmt.Fprintf(stderr, "felis operator: setup controller: %v\n", err) diff --git a/docs/openapi.yaml b/docs/openapi.yaml index aed73f1..30f6f42 100644 --- a/docs/openapi.yaml +++ b/docs/openapi.yaml @@ -741,6 +741,11 @@ paths: $ref: '#/components/responses/Forbidden' '404': $ref: '#/components/responses/NotFound' + '409': + description: A restore, backup or file write holds the server's world volume (maintenance_in_progress); nothing was started. + content: + application/json: + schema: { $ref: '#/components/schemas/Error' } '429': description: Wake cooldown is still active for this server. content: @@ -1193,7 +1198,7 @@ paths: application/json: schema: { $ref: '#/components/schemas/Error' } '409': - description: Server is not stopped (its world PVC is still mounted). + description: Server is not stopped (not_stopped), or a restore, backup or file write already holds its world volume (maintenance_in_progress). content: application/json: schema: { $ref: '#/components/schemas/Error' } @@ -1228,6 +1233,11 @@ paths: $ref: '#/components/responses/Forbidden' '404': $ref: '#/components/responses/NotFound' + '409': + description: A restore, backup or file write holds the server's world volume (maintenance_in_progress); nothing was started. + content: + application/json: + schema: { $ref: '#/components/schemas/Error' } '429': description: Wake cooldown is still active. content: @@ -2828,7 +2838,7 @@ paths: application/json: schema: { $ref: '#/components/schemas/Error' } '409': - description: Submission has already been reviewed. + description: Server is not stopped (not_stopped), or a restore, backup or file write already holds its world volume (maintenance_in_progress). content: application/json: schema: { $ref: '#/components/schemas/Error' } @@ -2873,7 +2883,7 @@ paths: application/json: schema: { $ref: '#/components/schemas/Error' } '409': - description: Server is not stopped (its world PVC is still mounted). + description: Server is not stopped (not_stopped), or a restore, backup or file write already holds its world volume (maintenance_in_progress). content: application/json: schema: { $ref: '#/components/schemas/Error' } @@ -3122,7 +3132,7 @@ paths: application/json: schema: { $ref: '#/components/schemas/Error' } '409': - description: Server is not stopped (its world PVC is still mounted). + description: Server is not stopped (not_stopped), or a restore, backup or file write already holds its world volume (maintenance_in_progress). content: application/json: schema: { $ref: '#/components/schemas/Error' } diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 2a46526..d65309a 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -185,6 +185,7 @@ per-server cooldown → global running cap**. Map the API result: | HTTP | Code | Cause | Fix | |---|---|---|---| | `403` | `forbidden` | `autostartPolicy=allowlist` and UUID not allowlisted, or `ownerOnly` and caller is not owner | Add the UUID / claim the server / set `autostartPolicy=public` | +| `409` | `maintenance_in_progress` | A restore, backup or file write holds the server's world volume (§3b) | Wait for the Job to finish | | `429` | (cooldown) | Wake retried within the 30s per-server `WakeCooldown` | Wait out the cooldown | | `503` | `at_capacity` | Global `MaxRunningServers` cap reached | Stop another server or raise the cap | @@ -194,11 +195,54 @@ gate; the proxy polls `GET /api/v1/internal/servers/{name}/status` every ~2s and teleports when `ready=true`. The Velocity-side consumption of these codes (`403` → "You're not allowed to -start «server»"; `429` → re-queue; other → "Couldn't start … Try again +start «server»"; `409 maintenance_in_progress` → "«server» is under +maintenance", not queued; `429` → re-queue; other → "Couldn't start … Try again shortly.") lives in the Java plugin and is **[CODE-ONLY]** — the codes it reacts to are produced by the Go-tested `authorizeWakeByUUID` / cooldown limiter, so grade the two halves separately. +### 3b. Wake, restore, backup or file save refused with `maintenance_in_progress` + +A server's world volume is ReadWriteOnce, and on a single node RWO lets a game +pod and a restore Job mount it side by side. So felis-api serialises them per +server: a restore, a backup, or a file write takes the world, and until its Job +finishes every wake (panel or join) and every other world operation on that +server gets `409 maintenance_in_progress`. File reads and listings never hold +it. The operator applies the same rule when `desiredState` is flipped to +`Running` by anything other than felis-api: the StatefulSet is not scaled up, +and the `Ready` condition reads `MaintenanceInProgress` until the Job ends. + +What holds the world, in order: + +1. An unfinished Job labelled `felis.lolicon.best/server=` with + `app.kubernetes.io/managed-by` `felis-restore`, `felis-backup`, or + `felis-files` plus `felis.lolicon.best/files-mode=write`: + + ```sh + kubectl -n minecraft get jobs -l felis.lolicon.best/server= + ``` + + A Job that is genuinely wedged is ended by its own `activeDeadlineSeconds`; + deleting it by hand releases the world at once (`kubectl -n minecraft delete + job `), at the cost of whatever it was writing. + +2. The admission lock `felis.lolicon.best/maintenance=@` on the + MinecraftServer. felis-api sets it for the milliseconds between admitting an + operation and creating its Job; it holds for at most two minutes if felis-api + died in between, and the next wake clears a stale one. To drop it by hand: + + ```sh + kubectl -n minecraft annotate minecraftserver felis.lolicon.best/maintenance- + ``` + +A restore, backup or file write refused with `409 not_stopped` although the +panel shows `Stopped` means the game pod is still terminating (its preStop save +can take a while); retry once `kubectl -n minecraft get pods -l +felis.lolicon.best/server=` shows nothing. + +[GO-TESTED: `internal/maintenance`, `k8scluster_maintenance_test.go`, +`handlers_maintenance_test.go`, operator `maintenance_test.go`.] + --- ## 4. Routing is disabled even though servers are up (online-mode coupling) @@ -947,7 +991,8 @@ installer built — only hand-built tags need a manual re-mirror. | RCON secret/auth/port errors | §1b, §1c | | Phase `Failed` | §2 | | Players land in lobby / wrong place | §3, §4 | -| Wake refused / rate-limited (403/429/503) | §3a | +| Wake refused / rate-limited (403/409/429/503) | §3a | +| `maintenance_in_progress`; server won't start after a restore | §3b | | Routing disabled, offline-mode | §4 | | Panel 401/403; fails-closed; audience error | §5 | | Local password login rejected | §5c | diff --git a/internal/api/api_test.go b/internal/api/api_test.go index fa2a7c5..f73c166 100644 --- a/internal/api/api_test.go +++ b/internal/api/api_test.go @@ -1454,12 +1454,19 @@ type fakeCluster struct { noWorld map[string]bool // server names modeled WITHOUT a world volume (never started / reaped) createErr error pingErr error + // maintErr / wakeErr: what AcquireMaintenance / SetDesiredState(Running) + // return for a server (the world-volume lock, internal/maintenance). + maintErr map[string]error + wakeErr map[string]error + acquired []string // "name:kind" per admitted AcquireMaintenance + released []string // names per ReleaseMaintenance } func newFakeCluster() *fakeCluster { return &fakeCluster{byName: map[string]*ServerInfo{}, bySub: map[string]*ServerInfo{}, desired: map[string]v1alpha1.DesiredState{}, created: map[string]CreateServerInput{}, - patched: map[string]ServerSpecPatch{}, noWorld: map[string]bool{}} + patched: map[string]ServerSpecPatch{}, noWorld: map[string]bool{}, + maintErr: map[string]error{}, wakeErr: map[string]error{}} } func (c *fakeCluster) GetServer(_ context.Context, n string) (*ServerInfo, error) { if s, ok := c.byName[n]; ok { @@ -1483,9 +1490,23 @@ func (c *fakeCluster) WorldVolumeExists(_ context.Context, n string) (bool, erro } func (c *fakeCluster) SetDesiredState(_ context.Context, n string, s v1alpha1.DesiredState) error { + if err := c.wakeErr[n]; err != nil && s == v1alpha1.DesiredRunning { + return err + } c.desired[n] = s return nil } +func (c *fakeCluster) AcquireMaintenance(_ context.Context, n, kind string) error { + if err := c.maintErr[n]; err != nil { + return err + } + c.acquired = append(c.acquired, n+":"+kind) + return nil +} +func (c *fakeCluster) ReleaseMaintenance(_ context.Context, n string) error { + c.released = append(c.released, n) + return nil +} func (c *fakeCluster) CreateServer(_ context.Context, in CreateServerInput) error { if c.createErr != nil { return c.createErr diff --git a/internal/api/cluster.go b/internal/api/cluster.go index d8b6f9e..d722b43 100644 --- a/internal/api/cluster.go +++ b/internal/api/cluster.go @@ -94,8 +94,18 @@ type Cluster interface { // velocity registration pull (spec §7 GET /servers). ListServers(ctx context.Context) ([]ServerInfo, error) // SetDesiredState flips spec.desiredState — the only write the API performs - // against the CRD (spec §9.1). It is idempotent. + // against the CRD (spec §9.1). It is idempotent. Flipping to Running returns a + // *MaintenanceBusyError (errors.Is ErrMaintenanceInProgress) while a restore, + // backup or file write holds the world volume. SetDesiredState(ctx context.Context, name string, state v1alpha1.DesiredState) error + // AcquireMaintenance admits one world-volume operation (internal/maintenance + // kind): ErrNotStopped unless the server is fully stopped, a + // *MaintenanceBusyError while another operation holds the volume. The check + // and the lock are one atomic write against a concurrent wake. + AcquireMaintenance(ctx context.Context, name, kind string) error + // ReleaseMaintenance drops the admission lock once the operation's Job exists + // (or could not be created). It is idempotent. + ReleaseMaintenance(ctx context.Context, name string) error // CreateServer creates a MinecraftServer CRD from the validated form (spec // §15). It returns ErrConflict if a server of that name already exists. CreateServer(ctx context.Context, in CreateServerInput) error diff --git a/internal/api/errors.go b/internal/api/errors.go index cef01e8..5fa34a9 100644 --- a/internal/api/errors.go +++ b/internal/api/errors.go @@ -82,8 +82,26 @@ var ( // sentinels so the handler answers 429 (a transient "too busy, retry" — the cap self-clears // as challenges expire), never a 400 that invites an immediate retry. ErrTooManyDiscoverableChallenges = errors.New("too many discoverable login challenges in flight") + // ErrNotStopped means a world-volume operation was refused because the server is + // not fully stopped: desiredState is not Stopped, or its pod is still shutting + // down (phase Stopping) and holds the volume while it saves. + ErrNotStopped = errors.New("server is not stopped") + // ErrMaintenanceInProgress means another operation holds the server's world + // volume (internal/maintenance). Cluster methods return it wrapped in a + // *MaintenanceBusyError that names the holder. + ErrMaintenanceInProgress = errors.New("world maintenance in progress") ) +// MaintenanceBusyError names what holds a server's world volume. errors.Is +// matches it against ErrMaintenanceInProgress. +type MaintenanceBusyError struct{ Kind string } + +func (e *MaintenanceBusyError) Error() string { + return "world maintenance in progress: " + e.Kind +} + +func (e *MaintenanceBusyError) Is(target error) bool { return target == ErrMaintenanceInProgress } + // apiError is a handler-level error carrying an HTTP status and a stable, // machine-readable code. The error envelope matches the platform convention: // diff --git a/internal/api/handlers_backups.go b/internal/api/handlers_backups.go index 247bd2e..840c001 100644 --- a/internal/api/handlers_backups.go +++ b/internal/api/handlers_backups.go @@ -6,6 +6,7 @@ import ( "strings" "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/maintenance" "felis.lolicon.best/internal/naming" ) @@ -69,7 +70,10 @@ func (a *API) handleListBackups(w http.ResponseWriter, r *http.Request) { // re-claims a released server could otherwise resurrect user A's world (the // backup still carries former_owner=A), a data leak. Admin skips this check. // ⑦ stopped gate: the world PVC must be free, so restore is refused unless the -// server is fully stopped. +// server is fully stopped. The world-volume lock (internal/maintenance) then +// makes that atomic against a wake and refuses a second restore, backup or +// file write on the same world with 409 maintenance_in_progress until the +// restore Job finishes. // ⑧ hand off to the Restorer. Restore is asynchronous (a restore Job, like an // image build Job), so success means "enqueued" and the handler answers 202. // @@ -152,7 +156,9 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) { // Refuse unless the server is fully stopped — Ready means it is up, and any // desiredState other than Stopped means it is up or coming up and still owns the // RWO volume (spec §141 readiness is an RCON probe; DesiredStopped is the - // intent). This yields a specific 409 instead of a restore Job that cannot mount. + // intent). This is the early, readable refusal from a snapshot; acquireWorld + // below is the atomic one (RWO is per node, so on a single node a restore Job + // WOULD mount beside a running server). info, err := a.Cluster.GetServer(r.Context(), name) if err != nil { a.writeLookupError(w, r, err) @@ -186,6 +192,16 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) { return } + // World-volume lock (internal/maintenance). The stopped gate above reads a + // snapshot; this is the atomic check, and it keeps a wake — the owner's, or a + // player's join through velocity — from booting the server on a half-extracted + // world until the restore Job has finished. + release, ok := a.acquireWorld(w, r, name, maintenance.KindRestore, "stop the server before restoring a backup") + if !ok { + return + } + defer release() + if err := a.Restorer.Restore(r.Context(), name, backup.BackupRef); err != nil { // ErrNotFound (server vanished from the execution backend) → 404; else 500. a.writeLookupError(w, r, err) @@ -212,9 +228,10 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) { // ② ServerByName — an unknown server is 404 // ③ owner-or-admin, else 403 (an unowned server passes only for admin, so a // released world can still be snapshotted by an operator before disposal) -// ④ stopped gate: the world PVC is RWO and held by a running server, so a backup -// Job cannot double-mount it — refuse unless the server is fully stopped. This -// also guarantees a quiescent, non-torn archive. +// ④ stopped gate: refuse unless the server is fully stopped, so the archive is +// quiescent and non-torn. The world-volume lock (internal/maintenance) keeps it +// that way until the backup Job finishes: a wake meanwhile, or a second +// restore/backup/file write, gets 409 maintenance_in_progress. // ⑤ hand off to the Backuper. Backup is asynchronous (a backup Job), so success // means "enqueued" and the handler answers 202. // @@ -294,9 +311,10 @@ func (a *API) handleInternalBackup(w http.ResponseWriter, r *http.Request) { // (Principal vs trusted service token) and the audit actor/source — keeping the // security-critical stopped-gate single-sourced so the two faces cannot diverge. func (a *API) enqueueBackup(w http.ResponseWriter, r *http.Request, name string, rec *ServerRecord, actor, source string) { - // Stopped gate: the world PVC is RWO and held by a running server, so a backup - // Job cannot double-mount it (mirrors the restore gate). Ready means it is up; - // any desiredState other than Stopped means it owns the RWO volume. + // Stopped gate (mirrors the restore gate): Ready means it is up; any + // desiredState other than Stopped means it is up or coming up. RWO is per node, + // so on a single node the Job WOULD mount beside a live server and archive a + // torn world; acquireWorld below makes this check atomic. info, err := a.Cluster.GetServer(r.Context(), name) if err != nil { a.writeLookupError(w, r, err) @@ -329,6 +347,14 @@ func (a *API) enqueueBackup(w http.ResponseWriter, r *http.Request, name string, return } + // World-volume lock (internal/maintenance): a server woken mid-backup would + // leave a torn archive that a later restore makes permanent. + release, ok := a.acquireWorld(w, r, name, maintenance.KindBackup, "stop the server before backing up its world") + if !ok { + return + } + defer release() + if err := a.Backuper.Backup(r.Context(), name, rec.OwnerID); err != nil { a.writeLookupError(w, r, err) return diff --git a/internal/api/handlers_files.go b/internal/api/handlers_files.go index d5a1095..9d58754 100644 --- a/internal/api/handlers_files.go +++ b/internal/api/handlers_files.go @@ -7,6 +7,7 @@ import ( "felis.lolicon.best/internal/apis/felis/v1alpha1" "felis.lolicon.best/internal/fileedit" + "felis.lolicon.best/internal/maintenance" "felis.lolicon.best/internal/naming" ) @@ -149,6 +150,15 @@ func (a *API) handleWriteFile(w http.ResponseWriter, r *http.Request) { return } + // A write holds the world volume for its Job's lifetime (internal/maintenance); + // reads and listings do not, since a read-only mount cannot hurt a server + // starting beside it. + release, ok := a.acquireWorld(w, r, name, maintenance.KindFileWrite, "stop the server before editing its files") + if !ok { + return + } + defer release() + if err := a.Files.Write(r.Context(), name, path, *body.Content); err != nil { writeFileEditError(w, r, err) return diff --git a/internal/api/handlers_files_test.go b/internal/api/handlers_files_test.go index 986185d..f2d6c3d 100644 --- a/internal/api/handlers_files_test.go +++ b/internal/api/handlers_files_test.go @@ -63,9 +63,9 @@ func mkFiles() (*API, *fakeRepo, *fakeCluster, *fakeFileEditor) { } // TestFileEditorStoppedGate is the gate this whole subsystem hinges on. The world -// PVC is ReadWriteOnce, so a running server holds it and a file Job physically -// cannot mount it — an ungated request would not fail cleanly, it would hang -// waiting for a Pod that can never be scheduled. Every one of the three routes +// PVC is ReadWriteOnce, but RWO is per node: on a single node a file Job mounts it +// right beside a running server, and a write lands under a live world that the +// server's next save overwrites or tears. Every one of the three routes // must therefore refuse a non-stopped server with 409 not_stopped BEFORE reaching // the executor, which is why each asserts calls == 0 as well as the status. func TestFileEditorStoppedGate(t *testing.T) { diff --git a/internal/api/handlers_internal.go b/internal/api/handlers_internal.go index f37ea5a..1347d12 100644 --- a/internal/api/handlers_internal.go +++ b/internal/api/handlers_internal.go @@ -153,8 +153,10 @@ func (a *API) handleInternalWake(w http.ResponseWriter, r *http.Request) { return } + // A 409 maintenance_in_progress tells velocity nothing is coming up until the + // restore/backup/file write finishes, so it does not enqueue the player. if err := a.Cluster.SetDesiredState(r.Context(), name, v1alpha1.DesiredRunning); err != nil { - writeError(w, r, err) + a.writeLookupError(w, r, err) return } // Consume the shared per-server cooldown only after the wake flips, so a join @@ -365,6 +367,8 @@ func (a *API) writeLookupError(w http.ResponseWriter, r *http.Request, err error writeError(w, r, newError(http.StatusNotFound, "not_found", "not found")) case errors.Is(err, ErrConflict): writeError(w, r, newError(http.StatusConflict, "conflict", "conflict")) + case errors.Is(err, ErrMaintenanceInProgress), errors.Is(err, ErrNotStopped): + writeError(w, r, maintenanceError(err, "stop the server completely first")) default: writeError(w, r, err) } diff --git a/internal/api/handlers_maintenance_test.go b/internal/api/handlers_maintenance_test.go new file mode 100644 index 0000000..31e688d --- /dev/null +++ b/internal/api/handlers_maintenance_test.go @@ -0,0 +1,176 @@ +package api + +import ( + "errors" + "fmt" + "net/http" + "net/http/httptest" + "slices" + "testing" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/maintenance" +) + +// The world-volume lock (internal/maintenance) from the handlers' side: a wake +// refused because a restore/backup/file write holds the world, and those three +// operations refused while another one does. The lock itself (atomicity, stale +// locks, Job-backed holds) is K8sCluster's and is tested in k8scluster_test. + +func TestWakeRefusedDuringMaintenance(t *testing.T) { + busy := &MaintenanceBusyError{Kind: maintenance.KindRestore} + + t.Run("external wake -> 409 maintenance_in_progress, cooldown kept", func(t *testing.T) { + repo := newFakeRepo() + cl := newFakeCluster() + cl.byName["survival"] = &ServerInfo{Name: "survival", AutostartPolicy: "public"} + cl.wakeErr["survival"] = busy + api := newTestAPI(repo, cl) + api.External = staticExternal{p: &Principal{UserID: "u1", Role: "user"}} + + w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/wake", "", nil) + if w.Code != http.StatusConflict || decodeErr(t, w) != "maintenance_in_progress" { + t.Fatalf("code = %d body %s, want 409 maintenance_in_progress", w.Code, w.Body.String()) + } + if _, set := cl.desired["survival"]; set { + t.Fatal("a refused wake must not flip desiredState") + } + + // The refusal did not burn the cooldown: once the restore is done the very + // next wake goes through. + delete(cl.wakeErr, "survival") + if w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/wake", "", nil); w.Code != http.StatusAccepted { + t.Fatalf("wake after maintenance: code = %d body %s", w.Code, w.Body.String()) + } + }) + + t.Run("internal wake (velocity) -> 409 maintenance_in_progress, cooldown kept", func(t *testing.T) { + api, cl := newInternalWakeAPI("public") + cl.wakeErr["survival"] = busy + body := `{"mc_uuid":"` + wakeUUID + `"}` + + w := internalWake(api, body) + if w.Code != http.StatusConflict || decodeErr(t, w) != "maintenance_in_progress" { + t.Fatalf("code = %d body %s, want 409 maintenance_in_progress", w.Code, w.Body.String()) + } + delete(cl.wakeErr, "survival") + if w := internalWake(api, body); w.Code != http.StatusAccepted { + t.Fatalf("wake after maintenance: code = %d body %s", w.Code, w.Body.String()) + } + }) +} + +// maintenanceOp is one world-volume operation as the external face serves it. +type maintenanceOp struct { + name string + kind string + method string + path string + body string + calls func() int +} + +func maintenanceOps() (*API, *fakeCluster, []maintenanceOp) { + repo := newFakeRepo() + repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"} + repo.backups = []fakeBackup{{view: BackupView{ID: "bk1", ServerName: "survival", + FormerOwner: "owner1", Status: "present"}, ref: "world-archive-ref"}} + cl := newFakeCluster() + cl.byName["survival"] = &ServerInfo{Name: "survival", Phase: "Stopped", + DesiredState: string(v1alpha1.DesiredStopped)} + restorer, backuper, files := &fakeRestorer{}, &fakeBackuper{}, &fakeFileEditor{} + api := newTestAPI(repo, cl) + api.Restorer, api.Backuper, api.Files = restorer, backuper, files + api.External = staticExternal{p: &Principal{UserID: "owner1", Email: "owner1@example.net", Role: "user"}} + return api, cl, []maintenanceOp{ + {"restore", maintenance.KindRestore, "POST", "/api/v1/servers/survival/restore-backup", "", + func() int { return restorer.calls }}, + {"backup", maintenance.KindBackup, "POST", "/api/v1/servers/survival/backup", "", + func() int { return backuper.calls }}, + {"file write", maintenance.KindFileWrite, "PUT", "/api/v1/servers/survival/file?path=server.properties", + `{"content":"aGk="}`, func() int { return files.calls }}, + } +} + +func (op maintenanceOp) do(api *API) *httptest.ResponseRecorder { + var hdr map[string]string + if op.body != "" { + hdr = jsonHeader + } + return do(api.ExternalHandler(), op.method, op.path, op.body, hdr) +} + +func TestMaintenanceOpsTakeAndReleaseTheLock(t *testing.T) { + _, _, ops := maintenanceOps() + for i := range ops { + t.Run(ops[i].name, func(t *testing.T) { + api, cl, ops := maintenanceOps() + op := ops[i] + if w := op.do(api); w.Code/100 != 2 { + t.Fatalf("code = %d body %s", w.Code, w.Body.String()) + } + if op.calls() != 1 { + t.Fatalf("executor calls = %d, want 1", op.calls()) + } + if want := []string{"survival:" + op.kind}; !slices.Equal(cl.acquired, want) { + t.Fatalf("acquired = %v, want %v", cl.acquired, want) + } + if want := []string{"survival"}; !slices.Equal(cl.released, want) { + t.Fatalf("released = %v, want %v (the Job is the lock from here on)", cl.released, want) + } + }) + } +} + +func TestMaintenanceOpsRefusedWhileHeld(t *testing.T) { + refusals := []struct { + name string + err error + code string + }{ + {"another holder", &MaintenanceBusyError{Kind: maintenance.KindBackup}, "maintenance_in_progress"}, + // The snapshot said Stopped but the atomic re-check found it waking: the + // wake won the race. + {"server not stopped", fmt.Errorf("wrapped: %w", ErrNotStopped), "not_stopped"}, + } + _, _, ops := maintenanceOps() + for i := range ops { + for _, rf := range refusals { + t.Run(ops[i].name+" / "+rf.name, func(t *testing.T) { + api, cl, ops := maintenanceOps() + op := ops[i] + cl.maintErr["survival"] = rf.err + w := op.do(api) + if w.Code != http.StatusConflict || decodeErr(t, w) != rf.code { + t.Fatalf("code = %d body %s, want 409 %s", w.Code, w.Body.String(), rf.code) + } + if op.calls() != 0 { + t.Fatal("a refused operation must not reach its executor") + } + if len(cl.released) != 0 { + t.Fatalf("released %v a lock that was never taken", cl.released) + } + }) + } + } +} + +func TestMaintenanceErrorMapping(t *testing.T) { + for _, tc := range []struct { + err error + code string + }{ + {&MaintenanceBusyError{Kind: maintenance.KindFileWrite}, "maintenance_in_progress"}, + {fmt.Errorf("x: %w", ErrMaintenanceInProgress), "maintenance_in_progress"}, + {ErrNotStopped, "not_stopped"}, + } { + var ae *apiError + if !errors.As(maintenanceError(tc.err, "stop it"), &ae) || ae.code != tc.code || ae.status != http.StatusConflict { + t.Errorf("%v -> %+v, want 409 %s", tc.err, ae, tc.code) + } + } + other := errors.New("boom") + if got := maintenanceError(other, "stop it"); got != other { + t.Errorf("unrelated error rewritten to %v", got) + } +} diff --git a/internal/api/handlers_user.go b/internal/api/handlers_user.go index dbfdd92..c1c62f5 100644 --- a/internal/api/handlers_user.go +++ b/internal/api/handlers_user.go @@ -56,8 +56,10 @@ func (a *API) handleWake(w http.ResponseWriter, r *http.Request) { return } + // Refused with 409 maintenance_in_progress while a restore, backup or file + // write holds the world volume: starting on a half-written world corrupts it. if err := a.Cluster.SetDesiredState(r.Context(), name, v1alpha1.DesiredRunning); err != nil { - writeError(w, r, err) + a.writeLookupError(w, r, err) return } // The wake actually flipped, so consume the per-server cooldown only now: a 503 diff --git a/internal/api/k8scluster.go b/internal/api/k8scluster.go index d306a3a..4506a6d 100644 --- a/internal/api/k8scluster.go +++ b/internal/api/k8scluster.go @@ -2,13 +2,18 @@ package api import ( "context" + "errors" + "time" "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/maintenance" "felis.lolicon.best/internal/naming" + batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/util/retry" "sigs.k8s.io/controller-runtime/pkg/client" ) @@ -20,6 +25,8 @@ import ( type K8sCluster struct { c client.Client namespace string + // now is injectable for the maintenance-lock tests; nil means time.Now. + now func() time.Time } // NewK8sCluster builds a Cluster over c, scoped to namespace. @@ -142,13 +149,15 @@ func (k *K8sCluster) CreateServer(ctx context.Context, in CreateServerInput) err } // SetDesiredState patches spec.desiredState with a merge patch so concurrent -// status writes by the operator are never clobbered (spec §9.1). +// status writes by the operator are never clobbered (spec §9.1). A stop always +// goes through. A start goes through start, which refuses while a maintenance +// operation holds the world volume. func (k *K8sCluster) SetDesiredState(ctx context.Context, name string, state v1alpha1.DesiredState) error { + if state == v1alpha1.DesiredRunning { + return k.start(ctx, name) + } var ms v1alpha1.MinecraftServer - if err := k.c.Get(ctx, types.NamespacedName{Namespace: k.namespace, Name: name}, &ms); err != nil { - if apierrors.IsNotFound(err) { - return ErrNotFound - } + if err := k.getServer(ctx, name, &ms); err != nil { return err } patch := client.MergeFrom(ms.DeepCopy()) @@ -156,6 +165,140 @@ func (k *K8sCluster) SetDesiredState(ctx context.Context, name string, state v1a return k.c.Patch(ctx, &ms, patch) } +// start flips desiredState to Running unless a restore, backup or file write +// holds the world volume (internal/maintenance), in which case it returns a +// *MaintenanceBusyError. The patch carries the resourceVersion it checked +// against, the same as AcquireMaintenance's: whichever of a racing wake and +// admission writes second gets a conflict, re-reads, and sees the other. +func (k *K8sCluster) start(ctx context.Context, name string) error { + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + var ms v1alpha1.MinecraftServer + if err := k.getServer(ctx, name, &ms); err != nil { + return err + } + kind, held, err := k.maintenanceHolder(ctx, &ms) + if err != nil { + return err + } + if held { + return &MaintenanceBusyError{Kind: kind} + } + patch := client.MergeFromWithOptions(ms.DeepCopy(), client.MergeFromWithOptimisticLock{}) + ms.Spec.DesiredState = v1alpha1.DesiredRunning + // A lock still on the object here no longer holds anything (Holder said + // so): drop it in the same write. + delete(ms.Annotations, maintenance.Annotation) + return k.c.Patch(ctx, &ms, patch) + }) +} + +// AcquireMaintenance admits one world-volume operation of the given kind: the +// server must be fully stopped (desiredState Stopped, phase Stopped, and no game +// pod left, so a pod still saving on its way down is waited out) and nothing else +// may hold the volume. Admission writes the maintenance lock under the resourceVersion it +// checked; the caller creates its Job and then calls ReleaseMaintenance, after +// which the Job itself is the lock. +func (k *K8sCluster) AcquireMaintenance(ctx context.Context, name, kind string) error { + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + var ms v1alpha1.MinecraftServer + if err := k.getServer(ctx, name, &ms); err != nil { + return err + } + desired := ms.Spec.DesiredState + if desired == "" { + desired = v1alpha1.DesiredStopped + } + if desired != v1alpha1.DesiredStopped || ms.Status.Ready || ms.Status.Phase != v1alpha1.PhaseStopped { + return ErrNotStopped + } + if up, err := k.gamePodExists(ctx, name); err != nil { + return err + } else if up { + return ErrNotStopped + } + holder, held, err := k.maintenanceHolder(ctx, &ms) + if err != nil { + return err + } + if held { + return &MaintenanceBusyError{Kind: holder} + } + patch := client.MergeFromWithOptions(ms.DeepCopy(), client.MergeFromWithOptimisticLock{}) + if ms.Annotations == nil { + ms.Annotations = map[string]string{} + } + ms.Annotations[maintenance.Annotation] = maintenance.LockValue(kind, k.clock()) + return k.c.Patch(ctx, &ms, patch) + }) +} + +// ReleaseMaintenance drops the admission lock. It is called once the Job exists +// (or failed to be created); a lock that is never released stops holding after +// maintenance.Grace on its own. +func (k *K8sCluster) ReleaseMaintenance(ctx context.Context, name string) error { + var ms v1alpha1.MinecraftServer + if err := k.getServer(ctx, name, &ms); err != nil { + if errors.Is(err, ErrNotFound) { + return nil + } + return err + } + if _, ok := ms.Annotations[maintenance.Annotation]; !ok { + return nil + } + patch := client.MergeFrom(ms.DeepCopy()) + delete(ms.Annotations, maintenance.Annotation) + return k.c.Patch(ctx, &ms, patch) +} + +// maintenanceHolder reads what holds ms's world volume right now. The Jobs are +// listed through the same direct client as the object, so a Job created before +// the lock was released is always visible here. +func (k *K8sCluster) maintenanceHolder(ctx context.Context, ms *v1alpha1.MinecraftServer) (string, bool, error) { + var jobs batchv1.JobList + if err := k.c.List(ctx, &jobs, client.InNamespace(k.namespace), + client.MatchingLabels{maintenance.LabelServer: ms.Name}); err != nil { + return "", false, err + } + kind, held := maintenance.Holder(ms.Name, ms.Annotations, jobs.Items, k.clock()) + return kind, held, nil +} + +// gamePodComponent is the operator's component label value on a game server's +// pod (internal/operator.ComponentValue; k8scluster_test pins the two). +const gamePodComponent = "server" + +// gamePodExists reports whether the server's game pod still exists, terminating +// or not. Phase Stopped is the operator's reading of the StatefulSet's replica +// counts; the pod object itself is the ground truth for "is anything of the +// server still running its preStop save against the volume". +func (k *K8sCluster) gamePodExists(ctx context.Context, name string) (bool, error) { + var pods corev1.PodList + if err := k.c.List(ctx, &pods, client.InNamespace(k.namespace), client.MatchingLabels{ + v1alpha1.LabelServer: name, v1alpha1.LabelComponent: gamePodComponent, + }); err != nil { + return false, err + } + return len(pods.Items) > 0, nil +} + +func (k *K8sCluster) getServer(ctx context.Context, name string, ms *v1alpha1.MinecraftServer) error { + if err := k.c.Get(ctx, types.NamespacedName{Namespace: k.namespace, Name: name}, ms); err != nil { + if apierrors.IsNotFound(err) { + return ErrNotFound + } + return err + } + return nil +} + +func (k *K8sCluster) clock() time.Time { + if k.now != nil { + return k.now() + } + return time.Now() +} + // PatchServerSpec applies the admin-tier spec mutation (spec §7) with the same // merge-patch discipline as SetDesiredState: read, copy, mutate only the fields // the admin set, patch — so an operator status write racing in parallel survives. diff --git a/internal/api/k8scluster_maintenance_test.go b/internal/api/k8scluster_maintenance_test.go new file mode 100644 index 0000000..d5a2715 --- /dev/null +++ b/internal/api/k8scluster_maintenance_test.go @@ -0,0 +1,237 @@ +package api + +import ( + "context" + "errors" + "testing" + "time" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/maintenance" + "felis.lolicon.best/internal/operator" + "felis.lolicon.best/internal/restore" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + clientgoscheme "k8s.io/client-go/kubernetes/scheme" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +// The world-volume lock against a fake API server. The fake client honours +// resourceVersion on an optimistic-lock patch, so the atomic half is exercised +// for real; what it cannot model is two felis-api replicas racing, which the +// resourceVersion check is precisely the defence against. + +var lockNow = time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC) + +func stoppedServer() *v1alpha1.MinecraftServer { + return &v1alpha1.MinecraftServer{ + ObjectMeta: metav1.ObjectMeta{Name: "survival", Namespace: "minecraft"}, + Spec: v1alpha1.MinecraftServerSpec{DesiredState: v1alpha1.DesiredStopped}, + Status: v1alpha1.MinecraftServerStatus{Phase: v1alpha1.PhaseStopped}, + } +} + +func lockCluster(t *testing.T, objs ...client.Object) (*K8sCluster, client.Client) { + t.Helper() + scheme := runtime.NewScheme() + if err := clientgoscheme.AddToScheme(scheme); err != nil { + t.Fatalf("scheme: %v", err) + } + if err := v1alpha1.AddToScheme(scheme); err != nil { + t.Fatalf("scheme: %v", err) + } + c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(objs...). + WithStatusSubresource(&v1alpha1.MinecraftServer{}).Build() + k := NewK8sCluster(c, "minecraft") + k.now = func() time.Time { return lockNow } + return k, c +} + +func runningRestore(t *testing.T) *batchv1.Job { + t.Helper() + j, err := restore.RestoreJob(restore.JobParams{ + Server: "survival", WorldPVC: "world-survival-0", BackupPVC: "felis-backups", + BackupRef: "/backups/survival/a.tar.gz", ArchiveStore: "tarLocal", + Namespace: "minecraft", ServiceAccount: "felis-restore", Image: "felis:1", + BackupRoot: "/backups", WorldsRoot: "/world", Deadline: time.Minute, + CPULimit: "1", MemLimit: "1Gi", TTLAfterFinished: time.Minute, + }) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + return j +} + +func annotations(t *testing.T, c client.Client) (map[string]string, v1alpha1.DesiredState) { + t.Helper() + var ms v1alpha1.MinecraftServer + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival"}, &ms); err != nil { + t.Fatalf("get: %v", err) + } + return ms.Annotations, ms.Spec.DesiredState +} + +func TestAcquireMaintenance(t *testing.T) { + ctx := context.Background() + + t.Run("stopped and free -> lock written", func(t *testing.T) { + k, c := lockCluster(t, stoppedServer()) + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindRestore); err != nil { + t.Fatalf("AcquireMaintenance: %v", err) + } + ann, _ := annotations(t, c) + if got, want := ann[maintenance.Annotation], maintenance.LockValue(maintenance.KindRestore, lockNow); got != want { + t.Fatalf("lock = %q, want %q", got, want) + } + // A second admission while the first holds (its Job not yet created) is refused. + var busy *MaintenanceBusyError + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindBackup); !errors.As(err, &busy) || busy.Kind != maintenance.KindRestore { + t.Fatalf("second admission: %v, want busy(restore)", err) + } + if err := k.ReleaseMaintenance(ctx, "survival"); err != nil { + t.Fatalf("ReleaseMaintenance: %v", err) + } + if ann, _ := annotations(t, c); ann[maintenance.Annotation] != "" { + t.Fatalf("lock survived release: %q", ann[maintenance.Annotation]) + } + }) + + notStopped := []struct { + name string + mut func(*v1alpha1.MinecraftServer) + }{ + {"desired Running", func(ms *v1alpha1.MinecraftServer) { ms.Spec.DesiredState = v1alpha1.DesiredRunning }}, + {"still Stopping", func(ms *v1alpha1.MinecraftServer) { ms.Status.Phase = v1alpha1.PhaseStopping }}, + {"Ready", func(ms *v1alpha1.MinecraftServer) { ms.Status.Ready = true }}, + {"never reconciled", func(ms *v1alpha1.MinecraftServer) { ms.Status.Phase = "" }}, + } + for _, tc := range notStopped { + t.Run(tc.name+" -> ErrNotStopped", func(t *testing.T) { + ms := stoppedServer() + tc.mut(ms) + k, c := lockCluster(t, ms) + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindRestore); !errors.Is(err, ErrNotStopped) { + t.Fatalf("err = %v, want ErrNotStopped", err) + } + if ann, _ := annotations(t, c); ann[maintenance.Annotation] != "" { + t.Fatal("a refused admission wrote the lock") + } + }) + } + + t.Run("game pod still terminating -> ErrNotStopped", func(t *testing.T) { + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "survival-0", Namespace: "minecraft", + Labels: map[string]string{v1alpha1.LabelServer: "survival", v1alpha1.LabelComponent: gamePodComponent}}} + k, _ := lockCluster(t, stoppedServer(), pod) + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindBackup); !errors.Is(err, ErrNotStopped) { + t.Fatalf("err = %v, want ErrNotStopped", err) + } + }) + + t.Run("a Job's own pod is not the game pod", func(t *testing.T) { + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "files-x", Namespace: "minecraft", + Labels: map[string]string{v1alpha1.LabelServer: "survival"}}} + k, _ := lockCluster(t, stoppedServer(), pod) + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindBackup); err != nil { + t.Fatalf("err = %v", err) + } + }) + + t.Run("running restore Job -> busy", func(t *testing.T) { + k, _ := lockCluster(t, stoppedServer(), runningRestore(t)) + var busy *MaintenanceBusyError + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindFileWrite); !errors.As(err, &busy) || busy.Kind != maintenance.KindRestore { + t.Fatalf("err = %v, want busy(restore)", err) + } + }) + + t.Run("unknown server -> ErrNotFound", func(t *testing.T) { + k, _ := lockCluster(t) + if err := k.AcquireMaintenance(ctx, "survival", maintenance.KindRestore); !errors.Is(err, ErrNotFound) { + t.Fatalf("err = %v, want ErrNotFound", err) + } + if err := k.ReleaseMaintenance(ctx, "survival"); err != nil { + t.Fatalf("release on a deleted server: %v", err) + } + }) +} + +func TestStartRespectsMaintenance(t *testing.T) { + ctx := context.Background() + + t.Run("fresh lock -> busy, desiredState untouched", func(t *testing.T) { + ms := stoppedServer() + ms.Annotations = map[string]string{maintenance.Annotation: maintenance.LockValue(maintenance.KindBackup, lockNow.Add(-10*time.Second))} + k, c := lockCluster(t, ms) + var busy *MaintenanceBusyError + if err := k.SetDesiredState(ctx, "survival", v1alpha1.DesiredRunning); !errors.As(err, &busy) || busy.Kind != maintenance.KindBackup { + t.Fatalf("err = %v, want busy(backup)", err) + } + if !errors.Is(&MaintenanceBusyError{}, ErrMaintenanceInProgress) { + t.Fatal("MaintenanceBusyError must match ErrMaintenanceInProgress") + } + if _, desired := annotations(t, c); desired != v1alpha1.DesiredStopped { + t.Fatalf("desiredState = %q, want Stopped", desired) + } + }) + + t.Run("running restore Job -> busy", func(t *testing.T) { + k, _ := lockCluster(t, stoppedServer(), runningRestore(t)) + if err := k.SetDesiredState(ctx, "survival", v1alpha1.DesiredRunning); !errors.Is(err, ErrMaintenanceInProgress) { + t.Fatalf("err = %v, want maintenance in progress", err) + } + }) + + t.Run("finished restore Job -> starts", func(t *testing.T) { + j := runningRestore(t) + j.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}} + k, c := lockCluster(t, stoppedServer(), j) + if err := k.SetDesiredState(ctx, "survival", v1alpha1.DesiredRunning); err != nil { + t.Fatalf("err = %v", err) + } + if _, desired := annotations(t, c); desired != v1alpha1.DesiredRunning { + t.Fatalf("desiredState = %q, want Running", desired) + } + }) + + t.Run("stale lock -> starts and the lock is dropped", func(t *testing.T) { + ms := stoppedServer() + ms.Annotations = map[string]string{maintenance.Annotation: maintenance.LockValue(maintenance.KindRestore, lockNow.Add(-maintenance.Grace-time.Second))} + k, c := lockCluster(t, ms) + if err := k.SetDesiredState(ctx, "survival", v1alpha1.DesiredRunning); err != nil { + t.Fatalf("err = %v", err) + } + ann, desired := annotations(t, c) + if desired != v1alpha1.DesiredRunning { + t.Fatalf("desiredState = %q, want Running", desired) + } + if _, ok := ann[maintenance.Annotation]; ok { + t.Fatal("the stale lock was left behind") + } + }) + + t.Run("stop ignores the lock", func(t *testing.T) { + ms := stoppedServer() + ms.Spec.DesiredState = v1alpha1.DesiredRunning + ms.Annotations = map[string]string{maintenance.Annotation: maintenance.LockValue(maintenance.KindRestore, lockNow)} + k, c := lockCluster(t, ms) + if err := k.SetDesiredState(ctx, "survival", v1alpha1.DesiredStopped); err != nil { + t.Fatalf("err = %v", err) + } + if _, desired := annotations(t, c); desired != v1alpha1.DesiredStopped { + t.Fatalf("desiredState = %q, want Stopped", desired) + } + }) +} + +// gamePodComponent is a copy of the operator's pod label value; a drift would +// let a restore start beside a terminating server. +func TestGamePodComponentMatchesOperator(t *testing.T) { + if gamePodComponent != operator.ComponentValue { + t.Fatalf("gamePodComponent = %q, operator labels its pods %q", gamePodComponent, operator.ComponentValue) + } +} diff --git a/internal/api/maintenance.go b/internal/api/maintenance.go new file mode 100644 index 0000000..52e32af --- /dev/null +++ b/internal/api/maintenance.go @@ -0,0 +1,64 @@ +package api + +import ( + "context" + "errors" + "log" + "net/http" + "time" + + "felis.lolicon.best/internal/maintenance" +) + +// maintenanceError maps the world-volume lock's refusals onto their 409s: +// maintenance_in_progress while a restore, backup or file write holds the volume, +// not_stopped (with the caller's wording) while the server is not fully down. +// Anything else passes through unchanged. +func maintenanceError(err error, notStopped string) error { + var busy *MaintenanceBusyError + switch { + case errors.As(err, &busy): + return newError(http.StatusConflict, "maintenance_in_progress", + "%s is running on this server's world; retry once it finishes", maintenanceLabel(busy.Kind)) + case errors.Is(err, ErrMaintenanceInProgress): + return newError(http.StatusConflict, "maintenance_in_progress", + "another operation is running on this server's world; retry once it finishes") + case errors.Is(err, ErrNotStopped): + return newError(http.StatusConflict, "not_stopped", "%s", notStopped) + } + return err +} + +func maintenanceLabel(kind string) string { + switch kind { + case maintenance.KindRestore: + return "a restore" + case maintenance.KindBackup: + return "a backup" + case maintenance.KindFileWrite: + return "a file write" + } + return "another operation" +} + +// acquireWorld admits one world-volume operation of the given kind +// (internal/maintenance) and returns the release to run once its Job exists, or +// false after writing the refusal. notStopped is the not_stopped wording the +// calling face uses. +// +// The release runs on a context detached from the request: a client that hangs +// up the moment its 202 is written must not leave the lock behind to refuse the +// owner's next wake for maintenance.Grace. +func (a *API) acquireWorld(w http.ResponseWriter, r *http.Request, name, kind, notStopped string) (func(), bool) { + if err := a.Cluster.AcquireMaintenance(r.Context(), name, kind); err != nil { + a.writeLookupError(w, r, maintenanceError(err, notStopped)) + return nil, false + } + return func() { + ctx, cancel := context.WithTimeout(context.WithoutCancel(r.Context()), 10*time.Second) + defer cancel() + if err := a.Cluster.ReleaseMaintenance(ctx, name); err != nil { + log.Printf("api: release the maintenance lock on %s: %v (it lapses after %s)", name, err, maintenance.Grace) + } + }, true +} diff --git a/internal/fileedit/jobspec.go b/internal/fileedit/jobspec.go index 07aade8..21f10db 100644 --- a/internal/fileedit/jobspec.go +++ b/internal/fileedit/jobspec.go @@ -16,12 +16,15 @@ import ( // way. LabelOpID is the one addition: it is how felis-api finds THIS invocation's // Pod among any others, which matters here in a way it does not for restore — // file operations are interactive and repeated, so several may be in flight or -// lingering inside their TTL at once. +// lingering inside their TTL at once. LabelMode carries the operation so the +// world-volume lock (internal/maintenance) can tell a write, which holds the +// volume, from a read or listing, which does not. const ( LabelManagedBy = "app.kubernetes.io/managed-by" LabelComponent = "app.kubernetes.io/component" LabelServer = "felis.lolicon.best/server" LabelOpID = "felis.lolicon.best/files-op" + LabelMode = "felis.lolicon.best/files-mode" managedByValue = "felis-files" componentValue = "world-files" @@ -80,6 +83,7 @@ func filesLabels(p JobParams) map[string]string { LabelComponent: componentValue, LabelServer: p.Server, LabelOpID: p.OpID, + LabelMode: p.Op, } } diff --git a/internal/maintenance/maintenance.go b/internal/maintenance/maintenance.go new file mode 100644 index 0000000..6183e0d --- /dev/null +++ b/internal/maintenance/maintenance.go @@ -0,0 +1,141 @@ +// Package maintenance is the per-server mutual exclusion between a game server +// and the Jobs that write or snapshot its world volume (restore, backup, file +// write). The world PVC is ReadWriteOnce, and RWO is exclusive per NODE: on a +// single-node cluster the game pod and a restore pod mount it side by side, so +// the access mode alone guards nothing. A server woken mid-restore boots on a +// half-extracted world and the restore then prunes what it wrote; a server woken +// mid-backup produces a torn archive that a later restore makes permanent. +// +// Two signals mark a volume as held: +// +// - an unfinished maintenance Job labelled for the server. Once the Job exists +// it IS the lock, for as long as it runs, however long that is. +// - the Annotation on the MinecraftServer, written by felis-api with an +// optimistic-lock patch before it creates the Job and removed right after. +// It bridges the gap between "admitted" and "the Job is visible", and it is +// what serialises admission against a wake: both write the same object under +// its resourceVersion, so one of two racing writers always loses with a +// conflict and re-checks. +// +// A lock older than Grace with no Job behind it is stale (felis-api died between +// the two writes) and holds nothing. +// +// File reads and listings are not holders. They mount the volume read-only for a +// second or two, and a server starting beside one cannot hurt either side, so +// nobody waits for them. +package maintenance + +import ( + "strings" + "time" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" +) + +const ( + // Annotation is the admission lock on a MinecraftServer. Its value is + // "@" (LockValue). + Annotation = "felis.lolicon.best/maintenance" + // Grace is how long a lock counts as held when no Job backs it. felis-api + // creates the Job within milliseconds of taking the lock and then drops it, so + // a lock this old means the process died in between. + Grace = 2 * time.Minute + + // LabelServer / LabelManagedBy are the labels every maintenance executor puts + // on its Job (internal/restore, internal/backupjob, internal/fileedit keep + // their own copies; maintenance_test pins them against these). + LabelServer = "felis.lolicon.best/server" + LabelManagedBy = "app.kubernetes.io/managed-by" + // LabelFilesMode is the file-editor operation (list, read, write) a files Job + // performs. Only write holds the volume. + LabelFilesMode = "felis.lolicon.best/files-mode" +) + +// Kinds of holder. +const ( + KindRestore = "restore" + KindBackup = "backup" + KindFileWrite = "file-write" +) + +// FilesModeWrite is the LabelFilesMode value of a file write. +const FilesModeWrite = "write" + +// JobKind names the holder a Job represents, or reports false for a Job that +// holds nothing (a file read, a build, anything else in the namespace). A files +// Job without LabelFilesMode predates the label and is counted as a write: it can +// only be an old Job still inside its TTL, and over-counting it for that window +// is the safe side. +func JobKind(j *batchv1.Job) (string, bool) { + switch j.Labels[LabelManagedBy] { + case "felis-restore": + return KindRestore, true + case "felis-backup": + return KindBackup, true + case "felis-files": + mode, ok := j.Labels[LabelFilesMode] + if !ok || mode == FilesModeWrite { + return KindFileWrite, true + } + } + return "", false +} + +// JobFinished reports whether a Job has reached a terminal condition. The +// success/failure-target conditions count as terminal: the Job controller sets +// them the moment the outcome is decided, before it finishes tearing the pods +// down, and waiting for Complete would keep a wake refused for no reason. +func JobFinished(j *batchv1.Job) bool { + for _, c := range j.Status.Conditions { + if c.Status != corev1.ConditionTrue { + continue + } + switch c.Type { + case batchv1.JobComplete, batchv1.JobFailed, batchv1.JobSuccessCriteriaMet, batchv1.JobFailureTarget: + return true + } + } + return false +} + +// LockValue renders the Annotation value for a holder admitted at `at`. +func LockValue(kind string, at time.Time) string { + return kind + "@" + at.UTC().Format(time.RFC3339) +} + +// parseLock splits a lock value. A value that does not parse is stale: a lock +// nobody can date must not be able to hold a server down forever. +func parseLock(v string) (string, time.Time, bool) { + kind, stamp, ok := strings.Cut(v, "@") + if !ok || kind == "" { + return "", time.Time{}, false + } + at, err := time.Parse(time.RFC3339, stamp) + if err != nil { + return "", time.Time{}, false + } + return kind, at, true +} + +// Holder reports what, if anything, holds the server's world volume at `now`: +// the first unfinished maintenance Job among jobs, else a lock in annotations +// younger than Grace. jobs may contain unrelated Jobs; only the server's own +// holders count. +func Holder(server string, annotations map[string]string, jobs []batchv1.Job, now time.Time) (string, bool) { + for i := range jobs { + j := &jobs[i] + if j.Labels[LabelServer] != server || JobFinished(j) { + continue + } + if kind, ok := JobKind(j); ok { + return kind, true + } + } + if v, ok := annotations[Annotation]; ok { + if kind, at, ok := parseLock(v); ok && now.Sub(at) < Grace && at.Sub(now) < Grace { + return kind, true + } + } + return "", false +} diff --git a/internal/maintenance/maintenance_test.go b/internal/maintenance/maintenance_test.go new file mode 100644 index 0000000..a6b72d3 --- /dev/null +++ b/internal/maintenance/maintenance_test.go @@ -0,0 +1,150 @@ +package maintenance + +import ( + "testing" + "time" + + "felis.lolicon.best/internal/backupjob" + "felis.lolicon.best/internal/fileedit" + "felis.lolicon.best/internal/restore" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" +) + +var now = time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC) + +func restoreJob(t *testing.T, server string) batchv1.Job { + t.Helper() + j, err := restore.RestoreJob(restore.JobParams{ + Server: server, WorldPVC: "world-" + server + "-0", BackupPVC: "felis-backups", + BackupRef: "/backups/" + server + "/a.tar.gz", ArchiveStore: "tarLocal", + Namespace: "minecraft", ServiceAccount: "felis-restore", Image: "felis:1", + BackupRoot: "/backups", WorldsRoot: "/world", Deadline: time.Minute, + CPULimit: "1", MemLimit: "1Gi", TTLAfterFinished: time.Minute, + }) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + return *j +} + +func backupJob(t *testing.T, server string) batchv1.Job { + t.Helper() + j, err := backupjob.BackupJob(backupjob.JobParams{ + Server: server, FormerOwner: "u", WorldPVC: "world-" + server + "-0", BackupPVC: "felis-backups", + Namespace: "minecraft", ServiceAccount: "felis-restore", Image: "felis:1", + ConfigSecret: "felis-config", ConfigMount: "/etc/felis", BackupRoot: "/backups", + WorldsRoot: "/world", Deadline: time.Minute, CPULimit: "1", MemLimit: "1Gi", + TTLAfterFinished: time.Minute, + }) + if err != nil { + t.Fatalf("BackupJob: %v", err) + } + return *j +} + +func filesJob(t *testing.T, server, op string) batchv1.Job { + t.Helper() + p := fileedit.JobParams{ + Server: server, OpID: "0011223344556677", Op: op, Path: "server.properties", + WorldPVC: "world-" + server + "-0", Namespace: "minecraft", ServiceAccount: "felis-restore", + Image: "felis:1", WorldsRoot: "/data", Deadline: time.Minute, CPULimit: "500m", + MemLimit: "256Mi", TTLAfterFinished: time.Minute, + } + if op == fileedit.OpWrite { + p.Content = []byte("motd=hi\n") + } + j, err := fileedit.FilesJob(p) + if err != nil { + t.Fatalf("FilesJob: %v", err) + } + return *j +} + +func finished(j batchv1.Job, cond batchv1.JobConditionType) batchv1.Job { + j.Status.Conditions = append(j.Status.Conditions, batchv1.JobCondition{Type: cond, Status: corev1.ConditionTrue}) + return j +} + +// The executors keep their own label copies (they must not import each other or +// this package's callers); this pins what they actually render against JobKind. +func TestJobKindMatchesTheExecutors(t *testing.T) { + for _, tc := range []struct { + name string + job batchv1.Job + kind string + ok bool + }{ + {"restore", restoreJob(t, "survival"), KindRestore, true}, + {"backup", backupJob(t, "survival"), KindBackup, true}, + {"file write", filesJob(t, "survival", fileedit.OpWrite), KindFileWrite, true}, + {"file read", filesJob(t, "survival", fileedit.OpRead), "", false}, + {"file list", filesJob(t, "survival", fileedit.OpList), "", false}, + } { + kind, ok := JobKind(&tc.job) + if kind != tc.kind || ok != tc.ok { + t.Errorf("%s: JobKind = %q, %v; want %q, %v", tc.name, kind, ok, tc.kind, tc.ok) + } + if tc.job.Labels[LabelServer] != "survival" { + t.Errorf("%s: server label = %q", tc.name, tc.job.Labels[LabelServer]) + } + } +} + +func TestUnlabelledFilesJobCountsAsWrite(t *testing.T) { + j := filesJob(t, "survival", fileedit.OpRead) + delete(j.Labels, LabelFilesMode) + if kind, ok := JobKind(&j); !ok || kind != KindFileWrite { + t.Fatalf("a files Job from before the mode label = %q, %v; want file-write", kind, ok) + } +} + +func TestHolderFromJobs(t *testing.T) { + restoreRunning := restoreJob(t, "survival") + if kind, ok := Holder("survival", nil, []batchv1.Job{restoreRunning}, now); !ok || kind != KindRestore { + t.Fatalf("running restore: Holder = %q, %v", kind, ok) + } + for _, cond := range []batchv1.JobConditionType{ + batchv1.JobComplete, batchv1.JobFailed, batchv1.JobSuccessCriteriaMet, batchv1.JobFailureTarget, + } { + if kind, ok := Holder("survival", nil, []batchv1.Job{finished(restoreRunning, cond)}, now); ok { + t.Errorf("restore with %s still holds as %q", cond, kind) + } + } + other := backupJob(t, "creative") + if kind, ok := Holder("survival", nil, []batchv1.Job{other}, now); ok { + t.Fatalf("another server's backup holds survival as %q", kind) + } + read := filesJob(t, "survival", fileedit.OpRead) + if kind, ok := Holder("survival", nil, []batchv1.Job{read}, now); ok { + t.Fatalf("a running file read holds the volume as %q", kind) + } + write := filesJob(t, "survival", fileedit.OpWrite) + if kind, ok := Holder("survival", nil, []batchv1.Job{read, write}, now); !ok || kind != KindFileWrite { + t.Fatalf("running file write: Holder = %q, %v", kind, ok) + } +} + +func TestHolderFromLock(t *testing.T) { + for _, tc := range []struct { + name string + value string + held bool + }{ + {"fresh", LockValue(KindBackup, now.Add(-10*time.Second)), true}, + {"just inside grace", LockValue(KindBackup, now.Add(-Grace+time.Second)), true}, + {"stale", LockValue(KindBackup, now.Add(-Grace)), false}, + {"far future", LockValue(KindBackup, now.Add(time.Hour)), false}, + {"garbage", "yes", false}, + {"no kind", "@" + now.Format(time.RFC3339), false}, + {"bad time", "restore@yesterday", false}, + } { + kind, held := Holder("survival", map[string]string{Annotation: tc.value}, nil, now) + if held != tc.held { + t.Errorf("%s (%q): held = %v, want %v", tc.name, tc.value, held, tc.held) + } + if held && kind != KindBackup { + t.Errorf("%s: kind = %q", tc.name, kind) + } + } +} diff --git a/internal/operator/maintenance_test.go b/internal/operator/maintenance_test.go new file mode 100644 index 0000000..d252266 --- /dev/null +++ b/internal/operator/maintenance_test.go @@ -0,0 +1,108 @@ +package operator_test + +import ( + "context" + "testing" + "time" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/maintenance" + appsv1 "k8s.io/api/apps/v1" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" +) + +func maintenanceJob(managedBy string, finished bool) *batchv1.Job { + j := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Name: "restore-survival", Namespace: "minecraft", + Labels: map[string]string{maintenance.LabelServer: "survival", maintenance.LabelManagedBy: managedBy}}} + if finished { + j.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}} + } + return j +} + +// A desiredState flipped to Running behind felis-api's back (kubectl, a script) +// while a restore Job is unpacking into the world must not get a game pod. +func TestReconcileRunning_HeldDuringMaintenance(t *testing.T) { + r, c := newReconciler(t, fakeProber{}, runningServer(), rconSecret(), maintenanceJob("felis-restore", false)) + r.Jobs = c + + res := reconcile(t, r, "survival") + if res.RequeueAfter <= 0 { + t.Fatal("a held server must requeue to notice the Job finishing") + } + var sts appsv1.StatefulSet + err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival"}, &sts) + if !apierrors.IsNotFound(err) { + t.Fatalf("StatefulSet created while the restore runs (err=%v)", err) + } + s := getServer(t, c, "survival") + cond := meta.FindStatusCondition(s.Status.Conditions, v1alpha1.ConditionReady) + if cond == nil || cond.Reason != "MaintenanceInProgress" { + t.Fatalf("Ready condition = %+v, want reason MaintenanceInProgress", cond) + } + if s.Status.StartRequestedAt != nil { + t.Fatal("the hold must not start the startup-timeout clock") + } + + // The Job finishes: the next pass starts the server. + var j batchv1.Job + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "restore-survival"}, &j); err != nil { + t.Fatalf("get job: %v", err) + } + j.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}} + if err := c.Status().Update(context.Background(), &j); err != nil { + if err := c.Update(context.Background(), &j); err != nil { + t.Fatalf("finish job: %v", err) + } + } + reconcile(t, r, "survival") + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival"}, &sts); err != nil { + t.Fatalf("StatefulSet not created after the restore finished: %v", err) + } +} + +func TestReconcileRunning_FreshLockHolds(t *testing.T) { + ms := runningServer() + ms.Annotations = map[string]string{maintenance.Annotation: maintenance.LockValue(maintenance.KindBackup, fixedNow().Add(-5*time.Second))} + r, c := newReconciler(t, fakeProber{}, ms, rconSecret()) + r.Jobs = c + reconcile(t, r, "survival") + var sts appsv1.StatefulSet + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival"}, &sts); !apierrors.IsNotFound(err) { + t.Fatalf("StatefulSet created under a fresh maintenance lock (err=%v)", err) + } +} + +func TestReconcileRunning_NotHeldByReadsOrFinishedJobs(t *testing.T) { + read := maintenanceJob("felis-files", false) + read.Name = "files-survival-read" + read.Labels[maintenance.LabelFilesMode] = "read" + r, c := newReconciler(t, fakeProber{}, runningServer(), rconSecret(), read, maintenanceJob("felis-backup", true)) + r.Jobs = c + reconcile(t, r, "survival") + var sts appsv1.StatefulSet + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival"}, &sts); err != nil { + t.Fatalf("StatefulSet not created: %v", err) + } +} + +// A server that is already up is never held: its pod is the one on the volume, +// and holding would only stall its readiness bookkeeping. +func TestReconcileRunning_RunningServerNotHeld(t *testing.T) { + r, c := newReconciler(t, fakeProber{}, runningServer(), rconSecret()) + reconcile(t, r, "survival") // creates the StatefulSet at replicas 1 + if err := c.Create(context.Background(), maintenanceJob("felis-restore", false)); err != nil { + t.Fatalf("create job: %v", err) + } + r.Jobs = c + reconcile(t, r, "survival") + cond := meta.FindStatusCondition(getServer(t, c, "survival").Status.Conditions, v1alpha1.ConditionReady) + if cond != nil && cond.Reason == "MaintenanceInProgress" { + t.Fatal("a scaled-up server was held") + } +} diff --git a/internal/operator/reconciler.go b/internal/operator/reconciler.go index e0ee795..ed5bf11 100644 --- a/internal/operator/reconciler.go +++ b/internal/operator/reconciler.go @@ -13,8 +13,10 @@ import ( "time" "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/maintenance" "felis.lolicon.best/internal/metrics" appsv1 "k8s.io/api/apps/v1" + batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" @@ -37,6 +39,7 @@ const ( requeueStopping = 5 * time.Second requeueSecret = 10 * time.Second requeueIdleProbe = 30 * time.Second + requeueMaintenance = 5 * time.Second defaultTimeoutSeconds = 300 defaultReadinessTimeoutSec = 300 ) @@ -69,6 +72,10 @@ type Reconciler struct { // Now is injectable for deterministic timestamps in tests; defaults to // metav1.Now. Now func() metav1.Time + // Jobs reads the minecraft namespace's Jobs for the world-volume lock + // (internal/maintenance). It is the manager's uncached API reader, so the + // operator needs jobs:list and no Job informer. Nil skips the check. + Jobs client.Reader } func (r *Reconciler) now() metav1.Time { @@ -111,6 +118,20 @@ func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Resu } func (r *Reconciler) reconcileRunning(ctx context.Context, server *v1alpha1.MinecraftServer) (ctrl.Result, error) { + if kind, held, err := r.maintenanceHold(ctx, server); err != nil { + return ctrl.Result{}, err + } else if held { + // Leave phase and the start anchor alone: nothing is starting yet, and a + // long restore must not be charged to the startup timeout. + server.Status.ObservedGeneration = server.Generation + r.setCondition(server, v1alpha1.ConditionReady, metav1.ConditionFalse, "MaintenanceInProgress", + "waiting for the "+kind+" on this server's world to finish before starting") + if err := r.patchStatus(ctx, server); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: requeueMaintenance}, nil + } + endpointAddress, err := r.ensureServices(ctx, server) if err != nil { return ctrl.Result{}, err @@ -283,6 +304,35 @@ func (r *Reconciler) reconcileStopped(ctx context.Context, server *v1alpha1.Mine return ctrl.Result{}, r.patchStatus(ctx, server) } +// maintenanceHold reports whether a restore, backup or file write holds the +// server's world volume (internal/maintenance) while its pod is about to be +// created. felis-api already refuses a wake in that state; this is the same rule +// for a desiredState flipped by anything else (kubectl, a script), since a game +// pod scheduled beside a restore Job boots on a half-extracted world. A server +// whose StatefulSet is already scaled up is never held: its pod exists, and +// stopping it here would only lose the players on it. +func (r *Reconciler) maintenanceHold(ctx context.Context, server *v1alpha1.MinecraftServer) (string, bool, error) { + if r.Jobs == nil { + return "", false, nil + } + var sts appsv1.StatefulSet + err := r.Get(ctx, types.NamespacedName{Namespace: server.Namespace, Name: server.Name}, &sts) + switch { + case apierrors.IsNotFound(err): + case err != nil: + return "", false, err + case sts.Spec.Replicas == nil || *sts.Spec.Replicas > 0 || sts.Status.Replicas > 0: + return "", false, nil + } + var jobs batchv1.JobList + if err := r.Jobs.List(ctx, &jobs, client.InNamespace(server.Namespace), + client.MatchingLabels{maintenance.LabelServer: server.Name}); err != nil { + return "", false, err + } + kind, held := maintenance.Holder(server.Name, server.Annotations, jobs.Items, r.now().Time) + return kind, held, nil +} + func (r *Reconciler) ensureServices(ctx context.Context, server *v1alpha1.MinecraftServer) (string, error) { endpointAddress := "" for _, svc := range []*corev1.Service{buildHeadlessService(server), buildClientService(server)} { diff --git a/internal/platform/rbac.go b/internal/platform/rbac.go index 0d4292e..6fd4dde 100644 --- a/internal/platform/rbac.go +++ b/internal/platform/rbac.go @@ -145,8 +145,11 @@ func APIBuildRole(p Params) *rbacv1.Role { // status (Status().Update — `update` only) and patches spec.desiredState to // Stopped for idle auto-stop (spec §8 — the one spec field it may write, using // the same merge patch as the reaper's Stop: without the grant the auto-stop -// call fails closed with a 403), and reads RCON Secrets. It never touches pods, -// PVCs, Events, or finalizers, so none appear here. +// call fails closed with a 403), and reads RCON Secrets. Jobs are list-only, +// through the manager's uncached API reader: before scaling a server up from zero +// the operator checks that no restore/backup/file-write Job holds its world +// (internal/maintenance). It never touches pods, PVCs, Events, or finalizers, so +// none appear here. func OperatorRole(p Params) *rbacv1.Role { p = p.withDefaults() return role(p.MinecraftNamespace, "felis-operator", ComponentOperator, []rbacv1.PolicyRule{ @@ -161,6 +164,9 @@ func OperatorRole(p Params) *rbacv1.Role { // did not already have. No update/delete — the password is written once and // removed by garbage collection through its controller reference. rule([]string{groupCore}, []string{"secrets"}, []string{"get", "list", "watch", "create"}), + // list only: an uncached List (no informer, so no watch) of the world-volume + // maintenance Jobs; the operator never creates or deletes a Job. + rule([]string{groupBatch}, []string{"jobs"}, []string{"list"}), }) } diff --git a/internal/platform/rbac_test.go b/internal/platform/rbac_test.go index a91ca58..784b36b 100644 --- a/internal/platform/rbac_test.go +++ b/internal/platform/rbac_test.go @@ -188,6 +188,15 @@ func TestOperatorRole_ScopeExact(t *testing.T) { t.Errorf("operator must NOT touch core/%s", res) } } + // The world-volume lock check lists Jobs uncached; it never writes one. + if !hasRule(op, groupBatch, "jobs", "list") { + t.Error("operator must have jobs:list (maintenance hold before scale-up)") + } + for _, v := range []string{"get", "watch", "create", "update", "patch", "delete"} { + if hasRule(op, groupBatch, "jobs", v) { + t.Errorf("operator must NOT have jobs:%s (list-only)", v) + } + } } // TestReaperRole_ScopeExact pins the reaper's two destructive powers and confirms diff --git a/panel/src/i18n/resources/en-US/errors.json b/panel/src/i18n/resources/en-US/errors.json index 2af285e..96c4f08 100644 --- a/panel/src/i18n/resources/en-US/errors.json +++ b/panel/src/i18n/resources/en-US/errors.json @@ -14,6 +14,7 @@ "console_unavailable": "Can't reach the server console right now — try again shortly.", "no_backup": "There's no restorable backup for this server yet.", "not_stopped": "Stop the server completely before restoring — a restore overwrites the live world volume.", + "maintenance_in_progress": "This server's world is busy with a restore, backup or file write — try again once it finishes, usually within a minute or two.", "no_world_volume": "This server has no world volume yet — start it once so it is created, then retry.", "restore_unavailable": "Restore isn't available right now — try again later.", "session_expired": "Your session expired — please sign in again.", diff --git a/panel/src/i18n/resources/zh-CN/errors.json b/panel/src/i18n/resources/zh-CN/errors.json index 39d9941..19cc253 100644 --- a/panel/src/i18n/resources/zh-CN/errors.json +++ b/panel/src/i18n/resources/zh-CN/errors.json @@ -14,6 +14,7 @@ "console_unavailable": "暂时无法连接服务器控制台,请稍后重试。", "no_backup": "这台服务器暂时没有可回档的备份。", "not_stopped": "回档会覆盖世界的实时存储卷,请先把服务器完全停止再回档。", + "maintenance_in_progress": "这台服务器的世界正在回档、备份或写入文件——等它完成后再试,通常不超过一两分钟。", "no_world_volume": "这台服务器还没有世界卷——先启动一次让它创建,然后再试。", "restore_unavailable": "回档功能当前不可用,请稍后再试。", "session_expired": "会话已过期——请重新登录。", diff --git a/panel/src/lib/api.ts b/panel/src/lib/api.ts index ccf0e7f..6f522b1 100644 --- a/panel/src/lib/api.ts +++ b/panel/src/lib/api.ts @@ -705,6 +705,11 @@ export function humanizeError(e: unknown): string { return t("no_backup"); case "not_stopped": return t("not_stopped"); + // World-volume lock: a restore, backup or file write is running on this + // server's world, so a wake or a second world operation is refused until the + // Job finishes (internal/maintenance). + case "maintenance_in_progress": + return t("maintenance_in_progress"); case "no_world_volume": return t("no_world_volume"); case "restore_unavailable": diff --git a/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java b/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java index 25c2167..a84ed5f 100644 --- a/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java +++ b/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java @@ -79,7 +79,8 @@ public final class FelisApiClient { * wake pulls the domain-autostart lever for {@code name} on behalf of the * joining player (spec §9.1, §14). The reply (202) carries the current phase * and ready flag so the caller can decide whether to wait. A 403 (policy gate), - * 429 (cooldown), or 503 {@code at_capacity} (running cap) arrives as a + * 409 {@code maintenance_in_progress} (a restore, backup or file write holds the + * world), 429 (cooldown), or 503 {@code at_capacity} (running cap) arrives as a * LinkException the caller branches on. */ public ServerView wake(String name, UUID mcUuid) throws LinkException { diff --git a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java index 0869efd..ab61f25 100644 --- a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java +++ b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java @@ -498,6 +498,20 @@ public final class WaitingRouter { zh ? "你无权启动「" + serverName + "」。" : "You're not allowed to start « " + serverName + " ».", NamedTextColor.RED)); return; + case 409: + if ("maintenance_in_progress".equals(e.errorCode())) { + // A restore, backup or file write owns the world right now and + // the server will not start until it finishes (minutes at most), + // so waiting here would only run into the queue timeout. + player.sendMessage(Component.text( + zh ? "「" + serverName + "」正在维护(回档、备份或改文件),请稍后再试。" + : "« " + serverName + " » is under maintenance (restore, backup or file edit)." + + " Please try again shortly.", + NamedTextColor.YELLOW)); + return; + } + logWakeFailure(player, serverName, zh, e); + return; case 429: break; // a wake is already in flight → join the existing wait case 503: @@ -513,13 +527,10 @@ public final class WaitingRouter { } // A 503 without the at_capacity code is a plain outage, not a // capacity verdict — report it like any other failure. - // fall through + logWakeFailure(player, serverName, zh, e); + return; default: - log.warn("Felis: wake {} failed (status={}): {}", serverName, e.statusCode(), e.getMessage()); - player.sendMessage(Component.text( - zh ? "现在无法启动「" + serverName + "」。请稍后再试。" - : "Couldn't start « " + serverName + " » right now. Try again shortly.", - NamedTextColor.RED)); + logWakeFailure(player, serverName, zh, e); return; } } @@ -531,6 +542,14 @@ public final class WaitingRouter { serverName, System.currentTimeMillis() + WAIT_TIMEOUT_MILLIS, fromMenu)); } + private void logWakeFailure(Player player, String serverName, boolean zh, LinkException e) { + log.warn("Felis: wake {} failed (status={}): {}", serverName, e.statusCode(), e.getMessage()); + player.sendMessage(Component.text( + zh ? "现在无法启动「" + serverName + "」。请稍后再试。" + : "Couldn't start « " + serverName + " » right now. Try again shortly.", + NamedTextColor.RED)); + } + private void transfer(Player player, String serverName, RegisteredServer backend) { player.createConnectionRequest(backend).connect().whenComplete((result, err) -> { if (err != null || (result != null && !result.isSuccessful())) {