fix(api): wake 与回档/备份/改文件按服务器互斥

This commit is contained in:
Lemon-miaow committed 2026-09-24 14:39:02 +08:00
1 parent 7819e5de50
commit abfe60d62d
27 files changed
+1300 -36

No files matched your search

+22 -1
View File
@@ -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
+11 -1
View File
@@ -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
+18
View File
@@ -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:
//
+34 -8
View File
@@ -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
+10
View File
@@ -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
+3 -3
View File
@@ -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) {
+5 -1
View File
@@ -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)
}
+176
View File
@@ -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: "[email protected]", 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)
}
}
+3 -1
View File
@@ -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
+148 -5
View File
@@ -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.
+237
View File
@@ -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)
}
}
+64
View File
@@ -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
}