diff --git a/internal/backup/archiver.go b/internal/backup/archiver.go new file mode 100644 index 0000000..5a3ccb3 --- /dev/null +++ b/internal/backup/archiver.go @@ -0,0 +1,45 @@ +// Package backup implements the WorldArchiver abstraction (spec §19). The +// reaper and the restore endpoint speak only to the interface and never learn +// whether the backend is a tar file, a VolumeSnapshot, or a Longhorn backup — +// ArchiveRef is deliberately opaque. +package backup + +import "context" + +// ArchiveRef is an opaque handle to a stored world archive. Depending on the +// backend it may be a tar path, a VolumeSnapshot name, or a Longhorn backup URL. +type ArchiveRef string + +// WorldArchiver archives, restores, and deletes a server's world. The signature +// is intentionally "archive a world" rather than "write bytes": snapshot +// backends (VolumeSnapshot/Longhorn) cannot produce an io.Reader — they create +// K8s objects referencing the source PVC (spec §19). +type WorldArchiver interface { + // Archive captures the world living on pvc for server and returns an opaque + // ref plus the stored size in bytes. + Archive(ctx context.Context, server, pvc string) (ref ArchiveRef, size int64, err error) + // Restore writes a previously archived world into targetPVC. + Restore(ctx context.Context, ref ArchiveRef, targetPVC string) error + // Delete removes the archive identified by ref. + Delete(ctx context.Context, ref ArchiveRef) error +} + +// PVCResolver maps a PVC name to the local filesystem path where it is mounted. +// In production the reaper Job mounts the source/backup PVCs and supplies a +// resolver over those mount points; tests supply temp dirs. +type PVCResolver func(pvc string) (string, error) + +// StaticResolver resolves PVC names from a fixed map, erroring on unknown names. +func StaticResolver(paths map[string]string) PVCResolver { + return func(pvc string) (string, error) { + if p, ok := paths[pvc]; ok { + return p, nil + } + return "", &UnknownPVCError{PVC: pvc} + } +} + +// UnknownPVCError is returned when a resolver cannot map a PVC name. +type UnknownPVCError struct{ PVC string } + +func (e *UnknownPVCError) Error() string { return "backup: unknown pvc " + e.PVC } diff --git a/internal/backup/tarlocal.go b/internal/backup/tarlocal.go new file mode 100644 index 0000000..928c988 --- /dev/null +++ b/internal/backup/tarlocal.go @@ -0,0 +1,286 @@ +package backup + +import ( + "archive/tar" + "compress/gzip" + "context" + "fmt" + "io" + "io/fs" + "os" + "path" + "path/filepath" + "strings" + "time" +) + +// TarLocal is the zero-storageClass-requirement backend (spec §19): it mounts +// the source PVC, tars+gzips it, and writes the archive into the backup PVC. +// It runs on any StorageClass, including hostPath-style local-path, where the +// snapshot backends cannot. +type TarLocal struct { + // BackupRoot is the directory (backup PVC mount) archives are written into. + BackupRoot string + // Resolve maps a PVC name to its mounted filesystem path. + Resolve PVCResolver + // Now is injectable for deterministic archive names in tests. + Now func() time.Time +} + +func (t *TarLocal) now() time.Time { + if t.Now != nil { + return t.Now() + } + return time.Now() +} + +// Archive tars+gzips the world on pvc into BackupRoot and returns the archive +// path as the opaque ref plus its on-disk size. +func (t *TarLocal) Archive(ctx context.Context, server, pvc string) (ArchiveRef, int64, error) { + srcDir, err := t.Resolve(pvc) + if err != nil { + return "", 0, err + } + if err := os.MkdirAll(t.BackupRoot, 0o750); err != nil { + return "", 0, fmt.Errorf("backup: mkdir backup root: %w", err) + } + name := fmt.Sprintf("%s-%d.tar.gz", server, t.now().UTC().UnixNano()) + dest := filepath.Join(t.BackupRoot, name) + + f, err := os.Create(dest) + if err != nil { + return "", 0, fmt.Errorf("backup: create archive: %w", err) + } + if err := writeTarGz(ctx, f, srcDir); err != nil { + f.Close() + os.Remove(dest) + return "", 0, err + } + if err := f.Close(); err != nil { + os.Remove(dest) + return "", 0, fmt.Errorf("backup: close archive: %w", err) + } + + info, err := os.Stat(dest) + if err != nil { + return "", 0, fmt.Errorf("backup: stat archive: %w", err) + } + return ArchiveRef(dest), info.Size(), nil +} + +// Restore extracts the archive at ref into the world mount for targetPVC, +// replacing the target's contents so the world equals the archive (spec §466 +// "restore PVC": a rollback must not leave stale files the backup lacks — e.g. a +// griefer's chunks). It extracts over the target, then removes any pre-existing +// entry the archive did not contain. +// +// The prune runs only after a fully successful extract: a corrupt or truncated +// archive fails before the prune, leaving the target as a (recoverable) partial +// overlay rather than a destroyed world. The archive is retained on restore, so +// such a failure is recoverable by re-running the Job. +// +// A top-level lost+found is never a prune target. It is a filesystem artifact +// (root-owned, mode 0700) that the non-root restore Pod cannot delete anyway, +// and writeTarGz includes it in the archive, so it is preserved on both axes. +func (t *TarLocal) Restore(ctx context.Context, ref ArchiveRef, targetPVC string) error { + dstDir, err := t.Resolve(targetPVC) + if err != nil { + return err + } + if err := os.MkdirAll(dstDir, 0o750); err != nil { + return fmt.Errorf("backup: mkdir restore target: %w", err) + } + f, err := os.Open(string(ref)) + if err != nil { + return fmt.Errorf("backup: open archive: %w", err) + } + defer f.Close() + + keep, err := readTarGz(ctx, f, dstDir) + if err != nil { + return err + } + return pruneToManifest(dstDir, keep) +} + +// pruneToManifest removes every entry under dstDir whose archive-relative path +// is absent from keep, giving Restore replace semantics. keep holds cleaned, +// forward-slash relative paths (no trailing slash) for every archive entry plus +// all of their ancestor directories, so a kept file's parent dirs are never +// removed. A top-level lost+found is always kept. dstDir (the mount root) is +// never removed. +func pruneToManifest(dstDir string, keep map[string]struct{}) error { + cleanDst := filepath.Clean(dstDir) + return filepath.WalkDir(cleanDst, func(p string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if p == cleanDst { + return nil // never remove the mount root itself + } + relNative, err := filepath.Rel(cleanDst, p) + if err != nil { + return err + } + rel := filepath.ToSlash(relNative) + if rel == "lost+found" { + if d.IsDir() { + return filepath.SkipDir // filesystem artifact: keep and don't descend + } + return nil + } + if _, ok := keep[rel]; ok { + return nil // the archive contained this path: keep it + } + // Stale: present in the target but absent from the archive. + if err := os.RemoveAll(p); err != nil { + return fmt.Errorf("backup: prune stale entry %q: %w", rel, err) + } + if d.IsDir() { + return filepath.SkipDir // already removed; don't descend into it + } + return nil + }) +} + +// Delete removes the tar archive at ref. +func (t *TarLocal) Delete(_ context.Context, ref ArchiveRef) error { + if err := os.Remove(string(ref)); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("backup: delete archive: %w", err) + } + return nil +} + +func writeTarGz(ctx context.Context, w io.Writer, srcDir string) error { + gz := gzip.NewWriter(w) + tw := tar.NewWriter(gz) + + root := filepath.Clean(srcDir) + err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error { + if err != nil { + return err + } + if ctx.Err() != nil { + return ctx.Err() + } + rel, err := filepath.Rel(root, path) + if err != nil { + return err + } + if rel == "." { + return nil // don't archive the root entry itself + } + // Normalize to forward slashes so archives are portable. + name := filepath.ToSlash(rel) + + switch { + case info.IsDir(): + hdr := &tar.Header{Name: name + "/", Mode: 0o750, Typeflag: tar.TypeDir} + return tw.WriteHeader(hdr) + case info.Mode().IsRegular(): + hdr := &tar.Header{Name: name, Mode: 0o640, Size: info.Size(), Typeflag: tar.TypeReg} + if err := tw.WriteHeader(hdr); err != nil { + return err + } + src, err := os.Open(path) + if err != nil { + return err + } + defer src.Close() + _, err = io.Copy(tw, src) + return err + default: + // Skip symlinks/devices/sockets: a world directory should be plain + // files, and refusing the rest avoids surprising archive contents. + return nil + } + }) + if err != nil { + return fmt.Errorf("backup: tar walk: %w", err) + } + if err := tw.Close(); err != nil { + return fmt.Errorf("backup: close tar: %w", err) + } + if err := gz.Close(); err != nil { + return fmt.Errorf("backup: close gzip: %w", err) + } + return nil +} + +// readTarGz extracts the gzip+tar stream into dstDir and returns the keep-set: +// the cleaned, forward-slash relative path of every entry the archive contained +// plus all of their ancestor directories. The caller uses it to prune stale +// target files for replace semantics. On any error the keep-set is incomplete +// and must not be used to prune (a partial manifest would delete live files the +// stream had not yet reached). +func readTarGz(ctx context.Context, r io.Reader, dstDir string) (map[string]struct{}, error) { + gz, err := gzip.NewReader(r) + if err != nil { + return nil, fmt.Errorf("backup: open gzip: %w", err) + } + defer gz.Close() + tr := tar.NewReader(gz) + + keep := make(map[string]struct{}) + cleanDst := filepath.Clean(dstDir) + for { + if ctx.Err() != nil { + return nil, ctx.Err() + } + hdr, err := tr.Next() + if err == io.EOF { + return keep, nil + } + if err != nil { + return nil, fmt.Errorf("backup: read tar: %w", err) + } + + // Guard against path traversal (zip-slip): the resolved target must stay + // within dstDir. + target := filepath.Join(cleanDst, filepath.FromSlash(hdr.Name)) + if target != cleanDst && !strings.HasPrefix(target, cleanDst+string(os.PathSeparator)) { + return nil, fmt.Errorf("backup: archive entry escapes target: %q", hdr.Name) + } + + // Record this entry and its ancestors in the keep-set. Names are stored + // in archive form (forward slash, no trailing slash) to match the + // relative paths pruneToManifest derives from the on-disk walk. + rememberKept(keep, hdr.Name) + + switch hdr.Typeflag { + case tar.TypeDir: + if err := os.MkdirAll(target, 0o750); err != nil { + return nil, err + } + case tar.TypeReg: + if err := os.MkdirAll(filepath.Dir(target), 0o750); err != nil { + return nil, err + } + out, err := os.OpenFile(target, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o640) + if err != nil { + return nil, err + } + if _, err := io.Copy(out, tr); err != nil { + out.Close() + return nil, err + } + if err := out.Close(); err != nil { + return nil, err + } + default: + // Ignore entry types tarLocal never writes. + } + } +} + +// rememberKept adds an archive entry name and every ancestor directory to keep, +// normalized to a cleaned forward-slash path with no trailing slash. Adding +// ancestors guards against archives that list a file without an explicit entry +// for its parent dir: the dir must still survive the prune. +func rememberKept(keep map[string]struct{}, name string) { + rel := path.Clean(strings.TrimSuffix(name, "/")) + for rel != "." && rel != "/" && rel != "" { + keep[rel] = struct{}{} + rel = path.Dir(rel) + } +} diff --git a/internal/backup/tarlocal_test.go b/internal/backup/tarlocal_test.go new file mode 100644 index 0000000..6c9bfc3 --- /dev/null +++ b/internal/backup/tarlocal_test.go @@ -0,0 +1,224 @@ +package backup_test + +import ( + "archive/tar" + "compress/gzip" + "context" + "os" + "path/filepath" + "testing" + + "felis.lolicon.best/internal/backup" +) + +// writeTree creates files (path->content) under root. +func writeTree(t *testing.T, root string, files map[string]string) { + t.Helper() + for rel, content := range files { + full := filepath.Join(root, filepath.FromSlash(rel)) + if err := os.MkdirAll(filepath.Dir(full), 0o750); err != nil { + t.Fatalf("mkdir: %v", err) + } + if err := os.WriteFile(full, []byte(content), 0o640); err != nil { + t.Fatalf("write %s: %v", rel, err) + } + } +} + +func TestTarLocalRoundTrip(t *testing.T) { + src := t.TempDir() + dst := t.TempDir() + backupRoot := t.TempDir() + + want := map[string]string{ + "level.dat": "world-seed-and-spawn", + "region/r.0.0.mca": "chunk-bytes-aaaa", + "data/scoreboard.dat": "{}", + "playerdata/uuid.dat": "player-state", + } + writeTree(t, src, want) + + archiver := &backup.TarLocal{ + BackupRoot: backupRoot, + Resolve: backup.StaticResolver(map[string]string{ + "src-pvc": src, + "dst-pvc": dst, + }), + } + ctx := context.Background() + + ref, size, err := archiver.Archive(ctx, "survival", "src-pvc") + if err != nil { + t.Fatalf("Archive: %v", err) + } + if size <= 0 { + t.Errorf("archive size = %d, want > 0", size) + } + if _, err := os.Stat(string(ref)); err != nil { + t.Fatalf("archive file missing: %v", err) + } + + if err := archiver.Restore(ctx, ref, "dst-pvc"); err != nil { + t.Fatalf("Restore: %v", err) + } + for rel, content := range want { + got, err := os.ReadFile(filepath.Join(dst, filepath.FromSlash(rel))) + if err != nil { + t.Errorf("restored file %s missing: %v", rel, err) + continue + } + if string(got) != content { + t.Errorf("restored %s = %q, want %q", rel, got, content) + } + } + + if err := archiver.Delete(ctx, ref); err != nil { + t.Fatalf("Delete: %v", err) + } + if _, err := os.Stat(string(ref)); !os.IsNotExist(err) { + t.Errorf("archive still present after Delete: %v", err) + } + // Delete of an already-gone archive is a no-op. + if err := archiver.Delete(ctx, ref); err != nil { + t.Errorf("second Delete should be a no-op, got %v", err) + } +} + +// TestTarLocalRestoreReplacesTarget pins replace semantics (spec §466): after a +// restore the world must equal the archive, not be merged onto whatever the +// target already held. It restores over a populated target and asserts that +// +// (a) archive files are present with the archive's content (overwriting stale +// copies), +// (b) files the archive did not contain are gone — including a stale chunk +// inside a directory the archive *does* keep, which proves per-file prune +// within a surviving dir and exercises the rel-path normalization, and +// (c) a pre-existing lost+found/ with a file inside survives untouched — the +// never-delete invariant for the filesystem artifact a non-root restore +// Pod cannot remove. +// +// Honesty: this runs as the test user (which *can* delete anything), so it +// proves the prune logic and the lost+found skip but does NOT exercise the +// non-root / FSGroup runtime path. "Restore works as a non-root Pod on a real +// ext4 PVC" remains code-complete-but-unverified (same bucket as the K8s E2E). +func TestTarLocalRestoreReplacesTarget(t *testing.T) { + src := t.TempDir() + dst := t.TempDir() + backupRoot := t.TempDir() + + archived := map[string]string{ + "level.dat": "new-seed", + "region/r.0.0.mca": "good-chunk-00", + "region/nested/a.mca": "good-chunk-nested", + "playerdata/uuid.dat": "player-state", + } + writeTree(t, src, archived) + + // The target already holds an older, divergent world: a stale copy of a file + // the archive also has, a griefer chunk inside a kept dir, and a whole stale + // directory the archive never mentions. + writeTree(t, dst, map[string]string{ + "level.dat": "OLD-seed-overwrite-me", + "region/r.9.9.mca": "griefer-chunk-must-vanish", + "oldworld/junk.dat": "whole-stale-dir-must-vanish", + }) + // A pre-existing lost+found with content the prune must never touch. + if err := os.MkdirAll(filepath.Join(dst, "lost+found"), 0o700); err != nil { + t.Fatalf("mkdir lost+found: %v", err) + } + if err := os.WriteFile(filepath.Join(dst, "lost+found", "0001"), []byte("fsck-recovered"), 0o600); err != nil { + t.Fatalf("seed lost+found: %v", err) + } + + archiver := &backup.TarLocal{ + BackupRoot: backupRoot, + Resolve: backup.StaticResolver(map[string]string{ + "src-pvc": src, + "dst-pvc": dst, + }), + } + ctx := context.Background() + + ref, _, err := archiver.Archive(ctx, "survival", "src-pvc") + if err != nil { + t.Fatalf("Archive: %v", err) + } + if err := archiver.Restore(ctx, ref, "dst-pvc"); err != nil { + t.Fatalf("Restore: %v", err) + } + + // (a) Every archive file present with the archive's content. + for rel, want := range archived { + got, err := os.ReadFile(filepath.Join(dst, filepath.FromSlash(rel))) + if err != nil { + t.Errorf("archive file %s missing after restore: %v", rel, err) + continue + } + if string(got) != want { + t.Errorf("restored %s = %q, want %q", rel, got, want) + } + } + + // (b) Files absent from the archive are gone — both the stale chunk inside the + // kept region/ dir and the whole stale directory. + for _, gone := range []string{"region/r.9.9.mca", "oldworld/junk.dat", "oldworld"} { + if _, err := os.Stat(filepath.Join(dst, filepath.FromSlash(gone))); !os.IsNotExist(err) { + t.Errorf("stale entry %s survived the restore (err=%v); replace semantics broken", gone, err) + } + } + + // (c) lost+found and its contents survive untouched. + lf, err := os.ReadFile(filepath.Join(dst, "lost+found", "0001")) + if err != nil { + t.Errorf("lost+found content was removed: %v", err) + } else if string(lf) != "fsck-recovered" { + t.Errorf("lost+found content = %q, want %q", lf, "fsck-recovered") + } +} + +func TestTarLocalUnknownPVC(t *testing.T) { + archiver := &backup.TarLocal{ + BackupRoot: t.TempDir(), + Resolve: backup.StaticResolver(map[string]string{}), + } + if _, _, err := archiver.Archive(context.Background(), "x", "missing"); err == nil { + t.Fatal("expected error for unknown pvc") + } +} + +// TestTarLocalRejectsZipSlip crafts a malicious archive whose entry escapes the +// target directory and asserts Restore refuses it. +func TestTarLocalRejectsZipSlip(t *testing.T) { + backupRoot := t.TempDir() + dst := t.TempDir() + evil := filepath.Join(backupRoot, "evil.tar.gz") + + f, err := os.Create(evil) + if err != nil { + t.Fatalf("create evil archive: %v", err) + } + gz := gzip.NewWriter(f) + tw := tar.NewWriter(gz) + body := []byte("pwned") + if err := tw.WriteHeader(&tar.Header{Name: "../escape.txt", Mode: 0o640, Size: int64(len(body)), Typeflag: tar.TypeReg}); err != nil { + t.Fatalf("write header: %v", err) + } + if _, err := tw.Write(body); err != nil { + t.Fatalf("write body: %v", err) + } + tw.Close() + gz.Close() + f.Close() + + archiver := &backup.TarLocal{ + BackupRoot: backupRoot, + Resolve: backup.StaticResolver(map[string]string{"dst-pvc": dst}), + } + if err := archiver.Restore(context.Background(), backup.ArchiveRef(evil), "dst-pvc"); err == nil { + t.Fatal("Restore must reject a path-traversal archive") + } + // Ensure nothing was written outside the target. + if _, err := os.Stat(filepath.Join(filepath.Dir(dst), "escape.txt")); !os.IsNotExist(err) { + t.Errorf("zip-slip wrote outside target: %v", err) + } +} diff --git a/internal/reaper/k8scluster.go b/internal/reaper/k8scluster.go new file mode 100644 index 0000000..1fe7cda --- /dev/null +++ b/internal/reaper/k8scluster.go @@ -0,0 +1,74 @@ +package reaper + +import ( + "context" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/naming" + 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" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// WorldPVCName returns the world PVC name for a server. The convention +// ("world--0") is owned by internal/naming because the operator, the +// reaper, and restore all depend on it; this is a thin alias kept so existing +// reaper call sites read naturally. +func WorldPVCName(server string) string { + return naming.WorldPVCName(server) +} + +// K8sCluster is the production Cluster backed by a controller-runtime client +// (spec §4, §18). It reads spec.reaperExempt, deletes the world PVC, and flips +// spec.desiredState to Stopped — nothing else. It is integration-tested against +// a live cluster, not the hermetic reaper_test.go suite. +type K8sCluster struct { + c client.Client + namespace string +} + +// NewK8sCluster builds a Cluster over c, scoped to namespace. +func NewK8sCluster(c client.Client, namespace string) *K8sCluster { + return &K8sCluster{c: c, namespace: namespace} +} + +func (k *K8sCluster) Inspect(ctx context.Context, name string) (ServerCRD, error) { + 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 ServerCRD{}, ErrNotFound + } + return ServerCRD{}, err + } + return ServerCRD{Exempt: ms.Spec.ReaperExempt, PVC: WorldPVCName(name)}, nil +} + +// DeletePVC deletes the world PersistentVolumeClaim. A missing PVC is not an +// error: the reap is idempotent and a re-run after a partial failure must still +// converge. +func (k *K8sCluster) DeletePVC(ctx context.Context, pvc string) error { + obj := &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{Namespace: k.namespace, Name: pvc}, + } + if err := k.c.Delete(ctx, obj); err != nil && !apierrors.IsNotFound(err) { + return err + } + return nil +} + +// Stop sets spec.desiredState=Stopped with a merge patch so a concurrent status +// write by the operator is never clobbered (spec §9.1). +func (k *K8sCluster) Stop(ctx context.Context, name string) error { + 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 + } + return err + } + patch := client.MergeFrom(ms.DeepCopy()) + ms.Spec.DesiredState = v1alpha1.DesiredStopped + return k.c.Patch(ctx, &ms, patch) +} diff --git a/internal/reaper/pgstore.go b/internal/reaper/pgstore.go new file mode 100644 index 0000000..8694592 --- /dev/null +++ b/internal/reaper/pgstore.go @@ -0,0 +1,155 @@ +package reaper + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "time" +) + +// PGStore is the production Store backed by Postgres (spec §6, §18). The SQL +// here is exercised by integration tests against a live database, not the +// hermetic reaper_test.go suite. Every statement is the narrow operation §18 +// requires; there is no generic UPDATE escape hatch. +type PGStore struct { + db *sql.DB +} + +// NewPGStore wraps an existing pool (from store.PostgresDriver.DB()). +func NewPGStore(db *sql.DB) *PGStore { return &PGStore{db: db} } + +func (s *PGStore) ListActiveServers(ctx context.Context) ([]Candidate, error) { + const q = `SELECT name, owner_id, last_active_at, warned_3d_at, warned_1d_at + FROM servers WHERE deleted_at IS NULL ORDER BY name` + rows, err := s.db.QueryContext(ctx, q) + if err != nil { + return nil, err + } + defer rows.Close() + var out []Candidate + for rows.Next() { + var ( + c Candidate + owner sql.NullString + w3, w1 sql.NullTime + ) + if err := rows.Scan(&c.Name, &owner, &c.LastActiveAt, &w3, &w1); err != nil { + return nil, err + } + c.OwnerID = owner.String + if w3.Valid { + c.Warned3dAt = w3.Time + } + if w1.Valid { + c.Warned1dAt = w1.Time + } + out = append(out, c) + } + return out, rows.Err() +} + +func (s *PGStore) FreshBackup(ctx context.Context, server string, since time.Time) (string, bool, error) { + const q = `SELECT backup_ref FROM world_backups + WHERE server_name = $1 AND status = 'present' AND created_at >= $2 + ORDER BY created_at DESC LIMIT 1` + var ref string + switch err := s.db.QueryRowContext(ctx, q, server, since).Scan(&ref); { + case err == sql.ErrNoRows: + return "", false, nil + case err != nil: + return "", false, err + } + return ref, true, nil +} + +func (s *PGStore) InsertBackup(ctx context.Context, rec BackupRecord) error { + const q = `INSERT INTO world_backups + (id, server_name, former_owner, backup_ref, size_bytes, reason, status, created_at, expires_at) + VALUES ($1, $2, NULLIF($3, ''), $4, $5, $6, 'present', now(), $7)` + _, err := s.db.ExecContext(ctx, q, + rec.ID, rec.ServerName, rec.FormerOwner, rec.BackupRef, rec.SizeBytes, rec.Reason, rec.ExpiresAt) + return err +} + +// ReleaseWorld releases ownership and resets the activity clock and warnings — +// without deleting the row (red line ②). +func (s *PGStore) ReleaseWorld(ctx context.Context, name string, at time.Time) error { + const q = `UPDATE servers + SET owner_id = NULL, last_active_at = $2, warned_3d_at = NULL, warned_1d_at = NULL + WHERE name = $1 AND deleted_at IS NULL` + _, err := s.db.ExecContext(ctx, q, name, at) + return err +} + +func (s *PGStore) MarkWarned(ctx context.Context, name string, tier Tier, at time.Time) error { + // The column is one of two fixed identifiers, never user input. + col := "warned_3d_at" + if tier == Tier1d { + col = "warned_1d_at" + } + q := fmt.Sprintf(`UPDATE servers SET %s = $2 WHERE name = $1 AND deleted_at IS NULL`, col) + _, err := s.db.ExecContext(ctx, q, name, at) + return err +} + +func (s *PGStore) PresentBackupBytes(ctx context.Context) (int64, error) { + var n int64 + err := s.db.QueryRowContext(ctx, + `SELECT COALESCE(SUM(size_bytes), 0) FROM world_backups WHERE status = 'present'`).Scan(&n) + return n, err +} + +func (s *PGStore) OldestPresentBackups(ctx context.Context) ([]StoredBackup, error) { + const q = `SELECT id, server_name, backup_ref, size_bytes FROM world_backups + WHERE status = 'present' ORDER BY created_at ASC` + return s.queryBackups(ctx, q) +} + +func (s *PGStore) ListExpiredBackups(ctx context.Context, now time.Time) ([]StoredBackup, error) { + const q = `SELECT id, server_name, backup_ref, size_bytes FROM world_backups + WHERE status = 'present' AND expires_at < $1 ORDER BY expires_at ASC` + return s.queryBackups(ctx, q, now) +} + +func (s *PGStore) queryBackups(ctx context.Context, q string, args ...any) ([]StoredBackup, error) { + rows, err := s.db.QueryContext(ctx, q, args...) + if err != nil { + return nil, err + } + defer rows.Close() + var out []StoredBackup + for rows.Next() { + var b StoredBackup + if err := rows.Scan(&b.ID, &b.ServerName, &b.BackupRef, &b.SizeBytes); err != nil { + return nil, err + } + out = append(out, b) + } + return out, rows.Err() +} + +func (s *PGStore) MarkBackupDeleted(ctx context.Context, id string, at time.Time) error { + _, err := s.db.ExecContext(ctx, + `UPDATE world_backups SET status = 'deleted', deleted_at = $2 WHERE id = $1`, id, at) + return err +} + +// Audit writes a reaper-sourced row. The actor/source are the system identity +// "reaper" (no human Access email applies, spec §14); former_owner has no +// dedicated column so it goes into the payload jsonb. +func (s *PGStore) Audit(ctx context.Context, rec AuditRecord) error { + payload := []byte("{}") + if rec.FormerOwner != "" { + p, err := json.Marshal(map[string]string{"former_owner": rec.FormerOwner}) + if err != nil { + return err + } + payload = p + } + _, err := s.db.ExecContext(ctx, + `INSERT INTO audit_logs (actor, source, action, server_name, payload) + VALUES ('reaper', 'reaper', $1, NULLIF($2, ''), $3)`, + rec.Action, rec.ServerName, string(payload)) + return err +} diff --git a/internal/reaper/reaper.go b/internal/reaper/reaper.go new file mode 100644 index 0000000..665c28b --- /dev/null +++ b/internal/reaper/reaper.go @@ -0,0 +1,491 @@ +// Package reaper implements the world reaper / three-clock retention batch +// (spec §18). It is the only component that deletes a player's world data, so +// every step is ordered and gated to honor §18's six red lines: +// +// ① system-server exemption — spec.reaperExempt servers are never touched +// (losing the lobby would be total ingress loss). +// ② delete PVC ≠ delete the servers row — the row stays so the subdomain +// remains reserved and a re-claim yields the same-named empty world. +// ③ world_backups is NOT FK'd to servers — a backup must outlive the world +// it came from (3-month retention from deletion). +// ④ back up BEFORE deleting — the PVC is only deleted after the archive is +// both written and recorded; an archive failure preserves the world. +// ⑤ warnings are best-effort — a delivery failure never blocks a reap, and a +// server with no linked owner is still reaped on time. +// ⑥ a real join resets the clock — RecordJoin (spec §7) refreshes +// last_active_at and clears warned_*, so renewal restarts the countdown. +// +// The logic here is pure and hermetically testable: all I/O is behind the +// Store, Cluster, and Warner interfaces plus the backup.WorldArchiver, with an +// injectable clock and id generator. The Postgres and Kubernetes bindings live +// in pgstore.go / k8scluster.go and are integration-tested, not unit-tested. +package reaper + +import ( + "context" + "crypto/rand" + "encoding/hex" + "errors" + "fmt" + "log/slog" + "sort" + "strconv" + "time" + + "felis.lolicon.best/internal/backup" +) + +// Day is a calendar day; the retention windows in §18 are expressed in days. +const Day = 24 * time.Hour + +// ReasonInactive is the canonical world_backups.reason for an idle-reaped world +// (spec §18). It is a stable label, not a literal restatement of the deadline. +const ReasonInactive = "inactive_15d" + +// Audit actions emitted by the reaper. The actor/source are a system identity +// ("reaper") because no human Access email is in play here (spec §14). +const ( + ActionReapWorld = "reap_world" + ActionEvictBackup = "evict_backup_early" +) + +// ErrNotFound is returned by Cluster.Inspect when the MinecraftServer CRD for a +// servers row no longer exists. The reaper treats it as "skip" — it will not +// delete world data it cannot first inspect for the exemption flag. +var ErrNotFound = errors.New("reaper: server not found") + +// errStoreFull is an internal sentinel: the backup store is at capacity and +// could not be freed, so the world is preserved rather than deleted without a +// backup (red line ④). It is never returned to callers. +var errStoreFull = errors.New("reaper: backup store full, world preserved") + +// Tier identifies which warning column a notice corresponds to. The tiers are +// positional: Tier3d is the first configured WarnBefore offset (the earlier, +// larger one) and maps to servers.warned_3d_at; Tier1d is the second and maps +// to warned_1d_at. The names follow the schema columns, not fixed durations — +// the actual thresholds are derived from Config.IdleBeforeReap. +type Tier int + +const ( + Tier3d Tier = iota + Tier1d +) + +// Config holds the retention windows. The warning thresholds are derived from +// the deadline (threshold = IdleBeforeReap - offset) rather than hardcoded, so +// changing IdleBeforeReap moves the warnings with it and the toml's +// warn_before=["3d","1d"] maps 1:1 onto WarnBefore. +type Config struct { + IdleBeforeReap time.Duration // §18: reap after this much per-server idle (default 15d) + WarnBefore []time.Duration // §24: warn these long before the deadline (default 3d, 1d) + Retention time.Duration // §18: keep a backup this long after deletion (default 3mo≈90d) + MaxLocalBytes int64 // §26: backup store soft cap; 0 = unlimited +} + +// DefaultConfig is the spec's §24 default window set. +func DefaultConfig() Config { + return Config{ + IdleBeforeReap: 15 * Day, + WarnBefore: []time.Duration{3 * Day, 1 * Day}, + Retention: 90 * Day, + MaxLocalBytes: 0, + } +} + +// Candidate is a servers row in the reaper's view (only deleted_at IS NULL rows +// are listed). Activity is per-server: LastActiveAt = max(last human join, +// created_at); stop/idle/wake do NOT move it (spec §18). +type Candidate struct { + Name string + OwnerID string // "" when unowned — still reaped, but never warned (red line ⑤) + LastActiveAt time.Time + Warned3dAt time.Time // zero = not yet sent + Warned1dAt time.Time // zero = not yet sent +} + +func (c Candidate) warnedAt(t Tier) time.Time { + if t == Tier1d { + return c.Warned1dAt + } + return c.Warned3dAt +} + +// ServerCRD is the slice of the MinecraftServer CRD the reaper needs: the +// exemption flag (red line ①) and the world PVC to archive then delete. +type ServerCRD struct { + Exempt bool + PVC string +} + +// BackupRecord is a world_backups insert. FormerOwner is captured so the +// backup, which outlives the server row, still records who it belonged to +// (red line ③). It is NOT a foreign key. +type BackupRecord struct { + ID string + ServerName string + FormerOwner string + BackupRef string + SizeBytes int64 + Reason string + ExpiresAt time.Time +} + +// StoredBackup is an existing world_backups row, used by both the expiry pass +// and the capacity-eviction path. +type StoredBackup struct { + ID string + ServerName string + BackupRef string + SizeBytes int64 +} + +// AuditRecord is a reaper-sourced audit_logs entry. The PG binding fills +// actor="reaper", source="reaper", and puts FormerOwner into the payload jsonb +// (audit_logs has no former_owner column). +type AuditRecord struct { + Action string + ServerName string + FormerOwner string +} + +// Store is the business-layer (Postgres) face the reaper needs. It deliberately +// exposes only the narrow operations §18 performs, never a generic UPDATE. +type Store interface { + // ListActiveServers returns every servers row with deleted_at IS NULL. + ListActiveServers(ctx context.Context) ([]Candidate, error) + + // FreshBackup reports an existing present backup for server whose world is + // still current — created at or after since (the world's last_active_at). + // It makes a reap idempotent across a DeletePVC failure: the retry reuses + // the archive instead of writing a duplicate. + FreshBackup(ctx context.Context, server string, since time.Time) (ref string, ok bool, err error) + + // InsertBackup records a world_backups row (status=present). + InsertBackup(ctx context.Context, rec BackupRecord) error + + // ReleaseWorld is the post-delete business mutation: owner_id→NULL, + // last_active_at→at (clock reset), warned_*→NULL. It does NOT delete the + // row (red line ②). + ReleaseWorld(ctx context.Context, name string, at time.Time) error + + // MarkWarned stamps the warned_3d_at / warned_1d_at column for tier. + MarkWarned(ctx context.Context, name string, tier Tier, at time.Time) error + + // PresentBackupBytes is the total size of status=present backups (§26 cap). + PresentBackupBytes(ctx context.Context) (int64, error) + + // OldestPresentBackups lists status=present backups oldest-first, for + // early eviction when the store is full. + OldestPresentBackups(ctx context.Context) ([]StoredBackup, error) + + // ListExpiredBackups lists status=present backups whose expires_at < now. + ListExpiredBackups(ctx context.Context, now time.Time) ([]StoredBackup, error) + + // MarkBackupDeleted flips a backup to status=deleted, deleted_at=at. + MarkBackupDeleted(ctx context.Context, id string, at time.Time) error + + // Audit appends a reaper-sourced audit_logs row. + Audit(ctx context.Context, rec AuditRecord) error +} + +// Cluster is the lifecycle (Kubernetes) face: read the CRD, delete the world +// PVC, and flip desiredState to Stopped. These are the only cluster operations +// §18 performs. +type Cluster interface { + // Inspect returns the exemption flag and world PVC name for a server, or + // ErrNotFound if the CRD is gone. + Inspect(ctx context.Context, name string) (ServerCRD, error) + // DeletePVC deletes the world PersistentVolumeClaim. + DeletePVC(ctx context.Context, pvc string) error + // Stop sets spec.desiredState=Stopped. + Stop(ctx context.Context, name string) error +} + +// Warner delivers an impending-reap notice. It is optional and best-effort: a +// nil Warner or a delivery error never blocks a reap (red line ⑤). +type Warner interface { + Warn(ctx context.Context, ownerID, server, remaining string) error +} + +// Reaper runs the §18 batch. Now and IDGen are injectable for hermetic tests; +// Log defaults to slog.Default(); Warner may be nil. +type Reaper struct { + Cfg Config + Store Store + Cluster Cluster + Archiver backup.WorldArchiver + Warner Warner + Log *slog.Logger + Now func() time.Time + IDGen func() string +} + +// Summary is the per-run tally (feeds §23 metrics). +type Summary struct { + Evaluated int + WorldsReaped int + Warned int + Skipped int // exempt, CRD gone, or could not back up + EvictedEarly int + BackupsExpired int +} + +func (r *Reaper) now() time.Time { + if r.Now != nil { + return r.Now() + } + return time.Now() +} + +func (r *Reaper) log() *slog.Logger { + if r.Log != nil { + return r.Log + } + return slog.Default() +} + +func (r *Reaper) id() string { + if r.IDGen != nil { + return r.IDGen() + } + var b [16]byte + if _, err := rand.Read(b[:]); err != nil { + return "bk-" + strconv.FormatInt(r.now().UnixNano(), 16) + } + return "bk-" + hex.EncodeToString(b[:]) +} + +// RunOnce executes one full batch: a world pass over active servers, then a +// retention pass over expired backups. It is idempotent and restart-safe, so a +// Kubernetes CronJob can drive the daily cadence (spec §18). Per-server +// failures are logged and counted as Skipped without aborting the batch; only +// an inability to list servers is a hard error. +func (r *Reaper) RunOnce(ctx context.Context) (Summary, error) { + var sum Summary + + cands, err := r.Store.ListActiveServers(ctx) + if err != nil { + return sum, fmt.Errorf("reaper: list active servers: %w", err) + } + + // Warning thresholds derive from the deadline, so the offsets must be + // largest-first (earliest warning first) to honor §18's elif precedence. + offs := append([]time.Duration(nil), r.Cfg.WarnBefore...) + sort.Slice(offs, func(i, j int) bool { return offs[i] > offs[j] }) + + now := r.now() + for _, c := range cands { + sum.Evaluated++ + if err := r.evaluate(ctx, now, offs, c, &sum); err != nil { + r.log().Error("reaper: skipping server", "server", c.Name, "err", err) + sum.Skipped++ + } + } + + r.expireBackups(ctx, now, &sum) + return sum, nil +} + +// evaluate handles one server: exemption, reap, or warning. A returned error +// means the server was skipped (counted by the caller); nil covers the normal +// outcomes including "exempt" and "warned". +func (r *Reaper) evaluate(ctx context.Context, now time.Time, offs []time.Duration, c Candidate, sum *Summary) error { + crd, err := r.Cluster.Inspect(ctx, c.Name) + if err != nil { + if errors.Is(err, ErrNotFound) { + // CRD gone but the row lingers — nothing safe to do; not a failure. + r.log().Warn("reaper: CRD missing, skipping", "server", c.Name) + return nil + } + return fmt.Errorf("inspect: %w", err) + } + if crd.Exempt { + // Red line ①: system servers (lobby/proxy) are never reaped. + return nil + } + + idle := now.Sub(c.LastActiveAt) + if idle > r.Cfg.IdleBeforeReap { + return r.reap(ctx, now, c, crd, sum) + } + r.maybeWarn(ctx, now, idle, offs, c, sum) + return nil +} + +// reap archives the world, records the backup, and only then deletes the PVC, +// releases ownership, and stops the server — the strict ordering of red line ④. +func (r *Reaper) reap(ctx context.Context, now time.Time, c Candidate, crd ServerCRD, sum *Summary) error { + // §26 soft cap: free space before adding a backup. If the store cannot be + // brought under cap, preserve the world rather than delete it unbacked. + if r.Cfg.MaxLocalBytes > 0 { + ok, err := r.ensureCapacity(ctx, now, sum) + if err != nil { + return fmt.Errorf("ensure capacity: %w", err) + } + if !ok { + r.log().Error("reaper: backup store full, world preserved", "server", c.Name) + return errStoreFull + } + } + + // Idempotent archive: if a prior run already archived this (unchanged) + // world but failed before deleting the PVC, reuse that backup rather than + // writing a duplicate. The world has not changed since last_active_at, so + // any present backup created after it still describes the current world. + ref, ok, err := r.Store.FreshBackup(ctx, c.Name, c.LastActiveAt) + if err != nil { + return fmt.Errorf("lookup fresh backup: %w", err) + } + if !ok { + aref, size, err := r.Archiver.Archive(ctx, c.Name, crd.PVC) + if err != nil { + // Red line ④: archive failed → the PVC is untouched, the world + // survives, and this server is retried next run. + return fmt.Errorf("archive: %w", err) + } + rec := BackupRecord{ + ID: r.id(), + ServerName: c.Name, + FormerOwner: c.OwnerID, + BackupRef: string(aref), + SizeBytes: size, + Reason: ReasonInactive, + ExpiresAt: now.Add(r.Cfg.Retention), + } + if err := r.Store.InsertBackup(ctx, rec); err != nil { + // The archive exists but is untracked. Delete the orphan so it does + // not leak, then fail without touching the PVC. + if derr := r.Archiver.Delete(ctx, aref); derr != nil { + r.log().Error("reaper: orphan archive cleanup failed", "server", c.Name, "ref", aref, "err", derr) + } + return fmt.Errorf("insert backup: %w", err) + } + ref = string(aref) + } + + // World is safely archived and recorded — now (and only now) delete it. + if err := r.Cluster.DeletePVC(ctx, crd.PVC); err != nil { + // The backup row persists; next run's FreshBackup reuses it and retries + // the delete, so no duplicate archive is created. + return fmt.Errorf("delete pvc: %w", err) + } + if err := r.Store.ReleaseWorld(ctx, c.Name, now); err != nil { + return fmt.Errorf("release world: %w", err) + } + if err := r.Cluster.Stop(ctx, c.Name); err != nil { + // The world is already deleted and ownership released; the desiredState + // flip is cosmetic by comparison. Log, but the reap stands. + r.log().Error("reaper: set desiredState=Stopped failed", "server", c.Name, "err", err) + } + if err := r.Store.Audit(ctx, AuditRecord{Action: ActionReapWorld, ServerName: c.Name, FormerOwner: c.OwnerID}); err != nil { + r.log().Error("reaper: audit reap_world failed", "server", c.Name, "err", err) + } + + sum.WorldsReaped++ + r.log().Info("reaper: world reaped", "server", c.Name, "former_owner", c.OwnerID, "backup_ref", ref) + return nil +} + +// ensureCapacity frees the backup store down under MaxLocalBytes by evicting the +// oldest present backups early. Early eviction is destructive (it removes +// not-yet-expired backups), so each eviction is alerted and audited. It returns +// whether the store is now under cap. +func (r *Reaper) ensureCapacity(ctx context.Context, now time.Time, sum *Summary) (bool, error) { + used, err := r.Store.PresentBackupBytes(ctx) + if err != nil { + return false, err + } + if used < r.Cfg.MaxLocalBytes { + return true, nil + } + r.log().Warn("reaper: backup store at capacity, evicting oldest backups early", + "used", used, "max", r.Cfg.MaxLocalBytes) + + old, err := r.Store.OldestPresentBackups(ctx) + if err != nil { + return false, err + } + for _, b := range old { + if used < r.Cfg.MaxLocalBytes { + break + } + if err := r.Archiver.Delete(ctx, backup.ArchiveRef(b.BackupRef)); err != nil { + r.log().Error("reaper: early-evict delete failed", "id", b.ID, "err", err) + continue + } + if err := r.Store.MarkBackupDeleted(ctx, b.ID, now); err != nil { + r.log().Error("reaper: early-evict mark failed", "id", b.ID, "err", err) + continue + } + if err := r.Store.Audit(ctx, AuditRecord{Action: ActionEvictBackup, ServerName: b.ServerName}); err != nil { + r.log().Error("reaper: audit evict failed", "id", b.ID, "err", err) + } + used -= b.SizeBytes + sum.EvictedEarly++ + } + return used < r.Cfg.MaxLocalBytes, nil +} + +// maybeWarn sends at most one impending-reap notice per run, honoring §18's +// elif precedence (earliest unsent warning first). Unowned servers are never +// warned but are still reaped at the deadline (red line ⑤). A warner delivery +// failure is logged but the warned_* stamp still advances so the notice is not +// retried forever; a real join (RecordJoin) is what clears the stamps. +func (r *Reaper) maybeWarn(ctx context.Context, now time.Time, idle time.Duration, offs []time.Duration, c Candidate, sum *Summary) { + if c.OwnerID == "" { + return + } + for i := 0; i < len(offs) && i < 2; i++ { + threshold := r.Cfg.IdleBeforeReap - offs[i] + if idle <= threshold { + continue + } + tier := Tier(i) + if !c.warnedAt(tier).IsZero() { + continue // already sent this tier + } + if r.Warner != nil { + if err := r.Warner.Warn(ctx, c.OwnerID, c.Name, formatRemaining(offs[i])); err != nil { + r.log().Warn("reaper: warn delivery failed (best-effort)", "server", c.Name, "err", err) + } + } + if err := r.Store.MarkWarned(ctx, c.Name, tier, now); err != nil { + r.log().Error("reaper: mark warned failed", "server", c.Name, "err", err) + return + } + sum.Warned++ + return // one warning per run + } +} + +// expireBackups is the retention pass: delete archives whose expires_at has +// passed and mark them deleted. Per-backup failures are logged, not fatal. +func (r *Reaper) expireBackups(ctx context.Context, now time.Time, sum *Summary) { + exp, err := r.Store.ListExpiredBackups(ctx, now) + if err != nil { + r.log().Error("reaper: list expired backups", "err", err) + return + } + for _, b := range exp { + if err := r.Archiver.Delete(ctx, backup.ArchiveRef(b.BackupRef)); err != nil { + r.log().Error("reaper: delete expired archive", "id", b.ID, "err", err) + continue + } + if err := r.Store.MarkBackupDeleted(ctx, b.ID, now); err != nil { + r.log().Error("reaper: mark expired deleted", "id", b.ID, "err", err) + continue + } + sum.BackupsExpired++ + } +} + +// formatRemaining renders an offset as the human-facing time left before reap. +func formatRemaining(d time.Duration) string { + if d%Day == 0 { + return strconv.FormatInt(int64(d/Day), 10) + "d" + } + if d%time.Hour == 0 { + return strconv.FormatInt(int64(d/time.Hour), 10) + "h" + } + return d.String() +} diff --git a/internal/reaper/reaper_test.go b/internal/reaper/reaper_test.go new file mode 100644 index 0000000..026bf2a --- /dev/null +++ b/internal/reaper/reaper_test.go @@ -0,0 +1,615 @@ +package reaper + +import ( + "context" + "errors" + "fmt" + "io" + "log/slog" + "reflect" + "sort" + "testing" + "time" + + "felis.lolicon.best/internal/backup" +) + +// testNow is the frozen clock for every hermetic case. Idle is expressed as an +// offset back from here. +var testNow = time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + +func idleBy(d time.Duration) time.Time { return testNow.Add(-d) } + +// recorder captures the cross-fake call order so a test can assert the strict +// archive→insert→deletePVC→release→stop→audit sequence of red line ④. +type recorder struct{ events []string } + +func (r *recorder) add(e string) { r.events = append(r.events, e) } + +// ---- fake backup.WorldArchiver ------------------------------------------- + +type fakeArchiver struct { + rec *recorder + archiveErr error + deleteErr error + archives int + deletes []backup.ArchiveRef + seq int +} + +func (f *fakeArchiver) Archive(_ context.Context, server, _ string) (backup.ArchiveRef, int64, error) { + if f.archiveErr != nil { + return "", 0, f.archiveErr + } + f.archives++ + f.seq++ + f.rec.add("archive") + return backup.ArchiveRef(fmt.Sprintf("ref-%s-%d", server, f.seq)), 10, nil +} + +func (f *fakeArchiver) Restore(context.Context, backup.ArchiveRef, string) error { return nil } + +func (f *fakeArchiver) Delete(_ context.Context, ref backup.ArchiveRef) error { + if f.deleteErr != nil { + return f.deleteErr + } + f.deletes = append(f.deletes, ref) + f.rec.add("delete") + return nil +} + +// ---- fake Cluster --------------------------------------------------------- + +type fakeCluster struct { + rec *recorder + crds map[string]ServerCRD + inspectErr map[string]error + deletePVCErr error + deletePVCCalls int + deletedPVCs []string + stopped []string +} + +func (c *fakeCluster) Inspect(_ context.Context, name string) (ServerCRD, error) { + if e := c.inspectErr[name]; e != nil { + return ServerCRD{}, e + } + crd, ok := c.crds[name] + if !ok { + return ServerCRD{}, ErrNotFound + } + return crd, nil +} + +func (c *fakeCluster) DeletePVC(_ context.Context, pvc string) error { + c.deletePVCCalls++ + if c.deletePVCErr != nil { + return c.deletePVCErr + } + c.deletedPVCs = append(c.deletedPVCs, pvc) + c.rec.add("deletePVC") + return nil +} + +func (c *fakeCluster) Stop(_ context.Context, name string) error { + c.stopped = append(c.stopped, name) + c.rec.add("stop") + return nil +} + +// ---- fake Store ----------------------------------------------------------- + +type fakeBackup struct { + id, server, ref string + size int64 + status string // present | deleted + createdAt, expires time.Time +} + +type fakeStore struct { + rec *recorder + clock time.Time + order []string + byName map[string]*Candidate + backups []*fakeBackup + audits []AuditRecord + released []string + listErr error + insertErr error + idn int +} + +func (s *fakeStore) ListActiveServers(context.Context) ([]Candidate, error) { + if s.listErr != nil { + return nil, s.listErr + } + out := make([]Candidate, 0, len(s.order)) + for _, n := range s.order { + out = append(out, *s.byName[n]) + } + return out, nil +} + +func (s *fakeStore) FreshBackup(_ context.Context, server string, since time.Time) (string, bool, error) { + for _, b := range s.backups { + if b.server == server && b.status == "present" && !b.createdAt.Before(since) { + return b.ref, true, nil + } + } + return "", false, nil +} + +func (s *fakeStore) InsertBackup(_ context.Context, rec BackupRecord) error { + if s.insertErr != nil { + return s.insertErr + } + s.backups = append(s.backups, &fakeBackup{ + id: rec.ID, server: rec.ServerName, ref: rec.BackupRef, size: rec.SizeBytes, + status: "present", createdAt: s.clock, expires: rec.ExpiresAt, + }) + s.rec.add("insert") + return nil +} + +func (s *fakeStore) ReleaseWorld(_ context.Context, name string, at time.Time) error { + c := s.byName[name] + c.OwnerID = "" + c.LastActiveAt = at + c.Warned3dAt = time.Time{} + c.Warned1dAt = time.Time{} + s.released = append(s.released, name) + s.rec.add("release") + return nil +} + +func (s *fakeStore) MarkWarned(_ context.Context, name string, tier Tier, at time.Time) error { + c := s.byName[name] + if tier == Tier1d { + c.Warned1dAt = at + } else { + c.Warned3dAt = at + } + return nil +} + +func (s *fakeStore) PresentBackupBytes(context.Context) (int64, error) { + var total int64 + for _, b := range s.backups { + if b.status == "present" { + total += b.size + } + } + return total, nil +} + +func (s *fakeStore) OldestPresentBackups(context.Context) ([]StoredBackup, error) { + var ps []*fakeBackup + for _, b := range s.backups { + if b.status == "present" { + ps = append(ps, b) + } + } + sort.Slice(ps, func(i, j int) bool { return ps[i].createdAt.Before(ps[j].createdAt) }) + out := make([]StoredBackup, 0, len(ps)) + for _, b := range ps { + out = append(out, StoredBackup{ID: b.id, ServerName: b.server, BackupRef: b.ref, SizeBytes: b.size}) + } + return out, nil +} + +func (s *fakeStore) ListExpiredBackups(_ context.Context, now time.Time) ([]StoredBackup, error) { + var out []StoredBackup + for _, b := range s.backups { + if b.status == "present" && b.expires.Before(now) { + out = append(out, StoredBackup{ID: b.id, ServerName: b.server, BackupRef: b.ref, SizeBytes: b.size}) + } + } + return out, nil +} + +func (s *fakeStore) MarkBackupDeleted(_ context.Context, id string, at time.Time) error { + for _, b := range s.backups { + if b.id == id { + b.status = "deleted" + return nil + } + } + return fmt.Errorf("no backup %s", id) +} + +func (s *fakeStore) Audit(_ context.Context, rec AuditRecord) error { + s.audits = append(s.audits, rec) + s.rec.add("audit:" + rec.Action) + return nil +} + +// ---- fixture -------------------------------------------------------------- + +func newReaper(cfg Config, cands ...Candidate) (*Reaper, *fakeStore, *fakeCluster, *fakeArchiver) { + rec := &recorder{} + st := &fakeStore{rec: rec, clock: testNow, byName: map[string]*Candidate{}} + cl := &fakeCluster{rec: rec, crds: map[string]ServerCRD{}, inspectErr: map[string]error{}} + for i := range cands { + cc := cands[i] + st.byName[cc.Name] = &cc + st.order = append(st.order, cc.Name) + cl.crds[cc.Name] = ServerCRD{PVC: "world-" + cc.Name + "-0"} + } + ar := &fakeArchiver{rec: rec} + r := &Reaper{ + Cfg: cfg, Store: st, Cluster: cl, Archiver: ar, + Log: slog.New(slog.NewTextHandler(io.Discard, nil)), + Now: func() time.Time { return testNow }, + IDGen: func() string { st.idn++; return fmt.Sprintf("bk-%d", st.idn) }, + } + return r, st, cl, ar +} + +func mustRun(t *testing.T, r *Reaper) Summary { + t.Helper() + sum, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + return sum +} + +// ---- tests ---------------------------------------------------------------- + +// Red line ①: a reaperExempt server is never archived, deleted, or warned no +// matter how idle it is — losing the lobby would be total ingress loss. +func TestReapExemptServerNeverTouched(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "lobby", OwnerID: "", LastActiveAt: idleBy(100 * Day)}) + cl.crds["lobby"] = ServerCRD{Exempt: true, PVC: "world-lobby-0"} + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 || sum.Skipped != 0 || ar.archives != 0 || cl.deletePVCCalls != 0 { + t.Fatalf("exempt server was touched: %+v archives=%d deletePVC=%d", sum, ar.archives, cl.deletePVCCalls) + } + if st.byName["lobby"].OwnerID != "" || len(st.released) != 0 { + t.Fatalf("exempt server ownership mutated") + } +} + +// The happy path, asserting the strict ordering of red line ④ and that the +// reap audit carries former_owner (spec §18 audit(reap_world, s, former_owner)). +func TestReapIdleWorldFullSequence(t *testing.T) { + r, st, cl, _ := newReaper(DefaultConfig(), + Candidate{Name: "alpha", OwnerID: "user-7", LastActiveAt: idleBy(20 * Day)}) + + sum := mustRun(t, r) + + if sum.WorldsReaped != 1 { + t.Fatalf("WorldsReaped = %d, want 1", sum.WorldsReaped) + } + want := []string{"archive", "insert", "deletePVC", "release", "stop", "audit:" + ActionReapWorld} + if !reflect.DeepEqual(st.rec.events, want) { + t.Fatalf("call order = %v, want %v", st.rec.events, want) + } + if got := cl.deletedPVCs; len(got) != 1 || got[0] != "world-alpha-0" { + t.Fatalf("deleted PVCs = %v, want [world-alpha-0]", got) + } + // Backup carries the former owner and a retention deadline 3mo out. + if len(st.backups) != 1 || st.backups[0].server != "alpha" { + t.Fatalf("backup not recorded: %+v", st.backups) + } + if want := testNow.Add(DefaultConfig().Retention); !st.backups[0].expires.Equal(want) { + t.Fatalf("backup expires = %v, want %v", st.backups[0].expires, want) + } + if len(st.audits) != 1 || st.audits[0].FormerOwner != "user-7" || st.audits[0].Action != ActionReapWorld { + t.Fatalf("audit = %+v, want reap_world former_owner=user-7", st.audits) + } + // Red line ②: the row survives (ownership released, not deleted). + c := st.byName["alpha"] + if c.OwnerID != "" || !c.LastActiveAt.Equal(testNow) || !c.Warned3dAt.IsZero() { + t.Fatalf("post-reap state wrong: %+v", c) + } +} + +// CENTERPIECE — red line ④: when the archive fails, the PVC is never deleted, +// ownership is untouched, no backup row is written, and the same server is +// retried (state unchanged) on the next run. +func TestReapArchiveFailurePreservesWorld(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "beta", OwnerID: "user-1", LastActiveAt: idleBy(20 * Day)}) + ar.archiveErr = errors.New("backend offline") + + sum := mustRun(t, r) + + if sum.WorldsReaped != 0 || sum.Skipped != 1 { + t.Fatalf("summary = %+v, want 0 reaped / 1 skipped", sum) + } + if cl.deletePVCCalls != 0 { + t.Fatalf("DeletePVC was called %d times despite archive failure", cl.deletePVCCalls) + } + if len(st.backups) != 0 { + t.Fatalf("a backup row was written despite archive failure: %+v", st.backups) + } + if len(st.released) != 0 { + t.Fatalf("ReleaseWorld ran despite archive failure") + } + c := st.byName["beta"] + if c.OwnerID != "user-1" || !c.LastActiveAt.Equal(idleBy(20*Day)) { + t.Fatalf("server state changed despite archive failure: %+v", c) + } + + // Recovery: backend returns, next run reaps cleanly. + ar.archiveErr = nil + sum2 := mustRun(t, r) + if sum2.WorldsReaped != 1 || cl.deletePVCCalls != 1 { + t.Fatalf("recovery run: summary=%+v deletePVC=%d, want 1 reaped / 1 delete", sum2, cl.deletePVCCalls) + } +} + +// Red line ④ (recording arm): if the archive succeeds but recording it fails, +// the orphan archive is cleaned up and the PVC is still never deleted. +func TestReapInsertBackupFailurePreservesWorld(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "gamma", OwnerID: "user-2", LastActiveAt: idleBy(20 * Day)}) + st.insertErr = errors.New("db down") + + sum := mustRun(t, r) + + if sum.WorldsReaped != 0 || sum.Skipped != 1 { + t.Fatalf("summary = %+v, want 0 reaped / 1 skipped", sum) + } + if cl.deletePVCCalls != 0 { + t.Fatalf("DeletePVC called despite insert failure") + } + if ar.archives != 1 || len(ar.deletes) != 1 { + t.Fatalf("orphan archive not cleaned up: archives=%d deletes=%d", ar.archives, len(ar.deletes)) + } + if len(st.backups) != 0 { + t.Fatalf("backup row present despite insert failure") + } +} + +// Idempotency (the deliberate disk-growth choice): a DeletePVC failure leaves a +// recorded backup; the retry reuses it via FreshBackup instead of writing a +// duplicate archive. +func TestReapDeletePVCFailureIsIdempotent(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "delta", OwnerID: "user-3", LastActiveAt: idleBy(20 * Day)}) + cl.deletePVCErr = errors.New("apiserver timeout") + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 || sum.Skipped != 1 { + t.Fatalf("run1 summary = %+v, want 0 reaped / 1 skipped", sum) + } + if ar.archives != 1 || len(st.backups) != 1 { + t.Fatalf("run1: archives=%d backups=%d, want 1/1", ar.archives, len(st.backups)) + } + + // apiserver recovers; the retry must NOT re-archive. + cl.deletePVCErr = nil + sum2 := mustRun(t, r) + if sum2.WorldsReaped != 1 { + t.Fatalf("run2 WorldsReaped = %d, want 1", sum2.WorldsReaped) + } + if ar.archives != 1 { + t.Fatalf("retry re-archived: archives=%d, want 1 (reuse via FreshBackup)", ar.archives) + } + if len(st.backups) != 1 { + t.Fatalf("retry duplicated the backup row: %d rows, want 1", len(st.backups)) + } +} + +// Red line ⑤ (reap arm): an unowned server is still reaped on time; the audit +// records an empty former_owner. +func TestReapUnownedServerStillReaped(t *testing.T) { + r, st, cl, _ := newReaper(DefaultConfig(), + Candidate{Name: "orphan", OwnerID: "", LastActiveAt: idleBy(20 * Day)}) + + sum := mustRun(t, r) + if sum.WorldsReaped != 1 || cl.deletePVCCalls != 1 { + t.Fatalf("unowned server not reaped: %+v deletePVC=%d", sum, cl.deletePVCCalls) + } + if len(st.audits) != 1 || st.audits[0].FormerOwner != "" { + t.Fatalf("audit former_owner = %q, want empty", st.audits[0].FormerOwner) + } +} + +func TestNoReapBeforeDeadline(t *testing.T) { + r, _, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "fresh", OwnerID: "user-4", LastActiveAt: idleBy(10 * Day)}) + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 || sum.Warned != 0 || ar.archives != 0 || cl.deletePVCCalls != 0 { + t.Fatalf("acted on a 10d-idle server with a 15d deadline: %+v", sum) + } +} + +// Warning thresholds are DERIVED from the deadline, not hardcoded. Using a +// non-default 10d deadline, warnings must fire at 10d-3d=7d and 10d-1d=9d, in +// elif precedence, deduped per tier, and never for an unowned server. +func TestWarningsDerivedFromNonDefaultDeadline(t *testing.T) { + cfg := Config{ + IdleBeforeReap: 10 * Day, + WarnBefore: []time.Duration{1 * Day, 3 * Day}, // intentionally unsorted + Retention: 90 * Day, + } + r, st, cl, _ := newReaper(cfg, + // past 7d, under 9d, nothing sent -> 3d (Tier3d) warning + Candidate{Name: "a", OwnerID: "u-a", LastActiveAt: idleBy(8 * Day)}, + // past 9d, 3d already sent -> 1d (Tier1d) warning + Candidate{Name: "b", OwnerID: "u-b", LastActiveAt: idleBy(95 * Day / 10), Warned3dAt: idleBy(2 * Day)}, + // past 7d, 3d already sent, under 9d -> no new warning (dedup) + Candidate{Name: "c", OwnerID: "u-c", LastActiveAt: idleBy(8 * Day), Warned3dAt: idleBy(2 * Day)}, + // under 7d -> no warning yet + Candidate{Name: "d", OwnerID: "u-d", LastActiveAt: idleBy(6 * Day)}, + // past 7d but unowned -> never warned (red line ⑤) + Candidate{Name: "e", OwnerID: "", LastActiveAt: idleBy(8 * Day)}, + ) + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 { + t.Fatalf("nothing should be reaped under a 10d deadline at <=9.5d idle: %+v", sum) + } + if sum.Warned != 2 { + t.Fatalf("Warned = %d, want 2 (a:3d, b:1d)", sum.Warned) + } + if cl.deletePVCCalls != 0 { + t.Fatalf("a warning path deleted a PVC") + } + if st.byName["a"].Warned3dAt.IsZero() || !st.byName["a"].Warned1dAt.IsZero() { + t.Fatalf("server a: expected only a 3d warning, got %+v", st.byName["a"]) + } + if st.byName["b"].Warned1dAt.IsZero() { + t.Fatalf("server b: expected a 1d warning, got %+v", st.byName["b"]) + } + if !st.byName["e"].Warned3dAt.IsZero() { + t.Fatalf("unowned server e was warned") + } +} + +// Red line ⑤ (best-effort): a Warner delivery error does not abort the run, and +// the warned_* stamp still advances (a real join, not a failed warn, is what +// resets the clock). +func TestWarningBestEffortOnDeliveryFailure(t *testing.T) { + r, st, _, _ := newReaper(DefaultConfig(), + Candidate{Name: "h", OwnerID: "u-h", LastActiveAt: idleBy(13 * Day)}) + r.Warner = failWarner{} + + sum := mustRun(t, r) + if sum.Warned != 1 { + t.Fatalf("Warned = %d, want 1 despite delivery failure", sum.Warned) + } + if st.byName["h"].Warned3dAt.IsZero() { + t.Fatalf("warned_3d_at not stamped after best-effort warn") + } +} + +type failWarner struct{} + +func (failWarner) Warn(context.Context, string, string, string) error { + return errors.New("smtp unavailable") +} + +// §26 capacity: when the store is over its cap, the oldest backup is evicted +// early (destructive — audited) to make room, then the reap proceeds. +func TestCapacityEvictsOldestThenReaps(t *testing.T) { + cfg := DefaultConfig() + cfg.MaxLocalBytes = 100 + r, st, _, ar := newReaper(cfg, + Candidate{Name: "epsilon", OwnerID: "user-5", LastActiveAt: idleBy(20 * Day)}) + // Two present backups of 75 each = 150 > 100. Oldest must be evicted first. + st.backups = []*fakeBackup{ + {id: "old", server: "zzz", ref: "ref-old", size: 75, status: "present", createdAt: idleBy(40 * Day), expires: testNow.Add(30 * Day)}, + {id: "new", server: "yyy", ref: "ref-new", size: 75, status: "present", createdAt: idleBy(5 * Day), expires: testNow.Add(60 * Day)}, + } + + sum := mustRun(t, r) + if sum.EvictedEarly != 1 { + t.Fatalf("EvictedEarly = %d, want 1", sum.EvictedEarly) + } + if sum.WorldsReaped != 1 { + t.Fatalf("WorldsReaped = %d, want 1 after eviction freed room", sum.WorldsReaped) + } + var oldStatus, newStatus string + for _, b := range st.backups { + switch b.id { + case "old": + oldStatus = b.status + case "new": + newStatus = b.status + } + } + if oldStatus != "deleted" || newStatus != "present" { + t.Fatalf("eviction hit wrong backup: old=%s new=%s, want deleted/present", oldStatus, newStatus) + } + // The destructive eviction must be audited. + var evicted bool + for _, a := range st.audits { + if a.Action == ActionEvictBackup { + evicted = true + } + } + if !evicted { + t.Fatalf("early eviction was not audited") + } + if ar.archives != 1 { + t.Fatalf("reap did not archive after eviction: archives=%d", ar.archives) + } +} + +// §26 capacity + red line ④: if eviction cannot free enough space, the reap is +// skipped and the world is preserved rather than deleted unbacked. +func TestCapacityStillFullSkipsReap(t *testing.T) { + cfg := DefaultConfig() + cfg.MaxLocalBytes = 100 + r, st, cl, ar := newReaper(cfg, + Candidate{Name: "zeta", OwnerID: "user-6", LastActiveAt: idleBy(20 * Day)}) + st.backups = []*fakeBackup{ + {id: "stuck", server: "zzz", ref: "ref-stuck", size: 150, status: "present", createdAt: idleBy(40 * Day), expires: testNow.Add(30 * Day)}, + } + // The archive backend can't delete, so eviction cannot free space. + r.Archiver.(*fakeArchiver).deleteErr = errors.New("evict unavailable") + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 || sum.Skipped != 1 { + t.Fatalf("summary = %+v, want 0 reaped / 1 skipped (store full)", sum) + } + if cl.deletePVCCalls != 0 { + t.Fatalf("world was deleted while the store was full") + } + if ar.archives != 0 { + t.Fatalf("archived into a full store") + } + if st.byName["zeta"].OwnerID != "user-6" { + t.Fatalf("ownership changed while store full") + } +} + +// Retention pass: backups past expires_at are deleted from the backend and +// marked deleted; unexpired backups are untouched. +func TestExpiredBackupsDeleted(t *testing.T) { + r, st, _, ar := newReaper(DefaultConfig()) + st.backups = []*fakeBackup{ + {id: "gone", server: "s1", ref: "ref-gone", size: 5, status: "present", createdAt: idleBy(120 * Day), expires: idleBy(1 * Day)}, + {id: "keep", server: "s2", ref: "ref-keep", size: 5, status: "present", createdAt: idleBy(10 * Day), expires: testNow.Add(80 * Day)}, + } + + sum := mustRun(t, r) + if sum.BackupsExpired != 1 { + t.Fatalf("BackupsExpired = %d, want 1", sum.BackupsExpired) + } + if len(ar.deletes) != 1 || ar.deletes[0] != "ref-gone" { + t.Fatalf("deleted archives = %v, want [ref-gone]", ar.deletes) + } + byID := map[string]string{} + for _, b := range st.backups { + byID[b.id] = b.status + } + if byID["gone"] != "deleted" || byID["keep"] != "present" { + t.Fatalf("expiry hit wrong rows: %v", byID) + } +} + +// A servers row whose CRD has been deleted is skipped (not a failure) — the +// reaper never deletes world data it cannot first inspect for the exemption. +func TestMissingCRDSkipped(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "ghost", OwnerID: "u", LastActiveAt: idleBy(20 * Day)}) + delete(cl.crds, "ghost") // CRD gone, servers row lingers + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 || sum.Skipped != 0 { + t.Fatalf("summary = %+v, want 0/0 (skipped without error)", sum) + } + if ar.archives != 0 || cl.deletePVCCalls != 0 { + t.Fatalf("acted on a server with no CRD") + } + if st.byName["ghost"].OwnerID != "u" { + t.Fatalf("mutated a server with no CRD") + } +} + +// A failure to list servers is the one hard error that aborts the batch. +func TestListErrorAbortsBatch(t *testing.T) { + r, st, _, _ := newReaper(DefaultConfig()) + st.listErr = errors.New("db unreachable") + if _, err := r.RunOnce(context.Background()); err == nil { + t.Fatal("expected a hard error when listing servers fails") + } +} diff --git a/internal/restore/jobspec.go b/internal/restore/jobspec.go new file mode 100644 index 0000000..31930aa --- /dev/null +++ b/internal/restore/jobspec.go @@ -0,0 +1,225 @@ +package restore + +import ( + "fmt" + "time" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// Label keys applied to restore objects, mirroring internal/build so the two +// executors are observable the same way. +const ( + LabelManagedBy = "app.kubernetes.io/managed-by" + LabelComponent = "app.kubernetes.io/component" + LabelServer = "felis.lolicon.best/server" + + managedByValue = "felis-restore" + componentValue = "world-restore" + + worldVolume = "world" + backupVolume = "backup" +) + +// JobParams are the rendered inputs to a restore Job, derived from a server + +// archive ref + Config by the Restorer. jobspec is a pure function of them so +// the security-critical Job shape is unit-tested without a cluster. +type JobParams struct { + Server string + WorldPVC string + BackupPVC string + BackupRef string + ArchiveStore string + Namespace string + ServiceAccount string + Image string + BackupRoot string + WorldsRoot string + Deadline time.Duration + CPULimit string + MemLimit string + RunAsUser int64 + RunAsGroup int64 + FSGroup int64 + + TTLAfterFinished time.Duration +} + +// RestoreJobName is the deterministic Job name for a server's restore. It is a +// pure function of the server name, which is how CreateRestoreJob detects a +// restore already in flight (AlreadyExists) and how the orchestrator's +// idempotency holds. +func RestoreJobName(server string) string { return "restore-" + server } + +func restoreLabels(p JobParams) map[string]string { + return map[string]string{ + LabelManagedBy: managedByValue, + LabelComponent: componentValue, + LabelServer: p.Server, + } +} + +// RestoreJob renders the world-restore Job (spec §7, §16, §22). Every isolation +// guarantee lives here and is asserted by jobspec_test.go, because no cluster +// runs in this environment: +// +// - runs under the weak felis-restore SA (never the felis-api SA) with its +// token auto-mount disabled, so it cannot reach the K8s API (spec §16, §21); +// - mounts EXACTLY two volumes — the world PVC read-write and the backup PVC +// read-only — and NO Secret/ConfigMap, so a poisoned archive cannot reach +// the felis database or any credential (the four-power red line, spec §22); +// - runs as a non-root, fixed uid/gid with an fsGroup so the files it writes +// are owned by the same identity the minecraft server later runs as; +// - no privilege, no privilege escalation, read-only root filesystem, drop ALL +// capabilities — all writes go to the mounted world PVC, nothing else; +// - activeDeadlineSeconds + backoffLimit=0 so a wedged or malicious archive +// cannot loop or run forever; ttlSecondsAfterFinished GCs the finished Job. +// +// The container runs `felis restore` (cmd/felis), which extracts the archive at +// BackupRef from the backup mount into the world mount. BackupRef is an absolute +// path, so the backup PVC MUST be mounted at BackupRoot — the same path the +// reaper wrote it under — for the ref to resolve. +func RestoreJob(p JobParams) (*batchv1.Job, error) { + if p.Image == "" { + return nil, fmt.Errorf("restore: image is empty") + } + if p.WorldPVC == "" || p.BackupPVC == "" { + return nil, fmt.Errorf("restore: world and backup PVC names are required") + } + limits, err := resourceLimits(p.CPULimit, p.MemLimit) + if err != nil { + return nil, err + } + deadline := int64(p.Deadline / time.Second) + if deadline <= 0 { + deadline = int64(defaultDeadline / time.Second) + } + ttl := int32(p.TTLAfterFinished / time.Second) + if ttl <= 0 { + ttl = int32(defaultTTL / time.Second) + } + + container := corev1.Container{ + Name: "restore", + Image: p.Image, + Command: []string{"felis", "restore"}, + Args: []string{ + "--server", p.Server, + "--ref", p.BackupRef, + "--archive-store", p.ArchiveStore, + "--backup-root", p.BackupRoot, + "--worlds-root", p.WorldsRoot, + }, + VolumeMounts: []corev1.VolumeMount{ + {Name: worldVolume, MountPath: p.WorldsRoot}, + // The archive is only ever read; mounting it read-only means a + // compromised restore process cannot mutate other servers' backups. + {Name: backupVolume, MountPath: p.BackupRoot, ReadOnly: true}, + }, + Resources: corev1.ResourceRequirements{Limits: limits, Requests: limits}, + SecurityContext: &corev1.SecurityContext{ + Privileged: boolPtr(false), + AllowPrivilegeEscalation: boolPtr(false), + ReadOnlyRootFilesystem: boolPtr(true), + Capabilities: &corev1.Capabilities{Drop: []corev1.Capability{"ALL"}}, + }, + } + + job := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: RestoreJobName(p.Server), + Namespace: p.Namespace, + Labels: restoreLabels(p), + }, + Spec: batchv1.JobSpec{ + // One shot: a bad archive must not loop. The TTL GCs the finished Job + // so a later restore of the same server is not blocked forever by a + // stale completed Job. + BackoffLimit: int32Ptr(0), + ActiveDeadlineSeconds: int64Ptr(deadline), + TTLSecondsAfterFinished: int32Ptr(ttl), + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{Labels: restoreLabels(p)}, + Spec: corev1.PodSpec{ + RestartPolicy: corev1.RestartPolicyNever, + ServiceAccountName: p.ServiceAccount, + AutomountServiceAccountToken: boolPtr(false), + SecurityContext: &corev1.PodSecurityContext{ + RunAsNonRoot: boolPtr(true), + RunAsUser: int64Ptr(p.RunAsUser), + RunAsGroup: int64Ptr(p.RunAsGroup), + FSGroup: int64Ptr(p.FSGroup), + }, + Containers: []corev1.Container{container}, + Volumes: []corev1.Volume{ + { + Name: worldVolume, + VolumeSource: corev1.VolumeSource{ + PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ + ClaimName: p.WorldPVC, + }, + }, + }, + { + Name: backupVolume, + VolumeSource: corev1.VolumeSource{ + PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ + ClaimName: p.BackupPVC, + ReadOnly: true, + }, + }, + }, + }, + }, + }, + }, + } + return job, nil +} + +// RestoreServiceAccount renders the weak restore SA (spec §16, §21). Like the +// build SA it is created bare: no secrets, token auto-mounting disabled, and — +// by having no Role or RoleBinding anywhere — zero K8s API permissions. Its only +// capability is filesystem access to the two PVCs the Job mounts. +func RestoreServiceAccount(namespace, name string) *corev1.ServiceAccount { + return &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + Labels: map[string]string{ + LabelManagedBy: managedByValue, + LabelComponent: componentValue, + }, + }, + AutomountServiceAccountToken: boolPtr(false), + } +} + +// resourceLimits parses the CPU/memory limits into a ResourceList. +func resourceLimits(cpu, mem string) (corev1.ResourceList, error) { + if cpu == "" { + cpu = defaultCPULimit + } + if mem == "" { + mem = defaultMemLimit + } + cpuQty, err := resource.ParseQuantity(cpu) + if err != nil { + return nil, fmt.Errorf("restore: invalid cpu limit %q: %w", cpu, err) + } + memQty, err := resource.ParseQuantity(mem) + if err != nil { + return nil, fmt.Errorf("restore: invalid memory limit %q: %w", mem, err) + } + return corev1.ResourceList{ + corev1.ResourceCPU: cpuQty, + corev1.ResourceMemory: memQty, + }, nil +} + +func boolPtr(b bool) *bool { return &b } +func int32Ptr(i int32) *int32 { return &i } +func int64Ptr(i int64) *int64 { return &i } diff --git a/internal/restore/jobspec_test.go b/internal/restore/jobspec_test.go new file mode 100644 index 0000000..649721f --- /dev/null +++ b/internal/restore/jobspec_test.go @@ -0,0 +1,282 @@ +package restore + +import ( + "testing" + "time" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" +) + +func sampleJobParams() JobParams { + return JobParams{ + Server: "survival", + WorldPVC: "world-survival-0", + BackupPVC: "felis-backups", + BackupRef: "/backups/survival/2026-06-25.tar.gz", + ArchiveStore: "tarLocal", + Namespace: defaultNamespace, + ServiceAccount: defaultServiceAccount, + Image: "registry.felis.svc:5000/felis:1.0", + BackupRoot: "/backups", + WorldsRoot: "/world", + Deadline: 30 * time.Minute, + CPULimit: "1", + MemLimit: "1Gi", + RunAsUser: 1000, + RunAsGroup: 1000, + FSGroup: 1000, + TTLAfterFinished: 10 * time.Minute, + } +} + +// The restore Pod must run under the weak felis-restore SA — never the +// felis-api identity — with its token un-mounted, so it cannot reach the K8s +// API. This is the §16/§22 red line asserted on the rendered spec because no +// cluster runs here. +func TestRestoreJobRunsUnderWeakSA(t *testing.T) { + job, err := RestoreJob(sampleJobParams()) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + sa := job.Spec.Template.Spec.ServiceAccountName + if sa != defaultServiceAccount { + t.Errorf("service account = %q, want %q", sa, defaultServiceAccount) + } + if sa == "felis-api" { + t.Fatal("restore Pod must NOT run as the felis-api SA") + } + if amt := job.Spec.Template.Spec.AutomountServiceAccountToken; amt == nil || *amt { + t.Error("AutomountServiceAccountToken must be explicitly false") + } +} + +// The four-power red line: a restore Pod handles a (potentially poisoned) +// archive, so it must mount EXACTLY the two PVCs — world read-write, backup +// read-only — and NO Secret or ConfigMap, so it can never reach the felis +// database or any credential. +func TestRestoreJobMountsOnlyTheTwoPVCsAndNoSecrets(t *testing.T) { + job, err := RestoreJob(sampleJobParams()) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + vols := job.Spec.Template.Spec.Volumes + if len(vols) != 2 { + t.Fatalf("expected exactly 2 volumes (world + backup), got %d: %+v", len(vols), vols) + } + var world, backup *corev1.Volume + for i := range vols { + v := &vols[i] + // The forbidden volume kinds: anything that could carry DB creds or + // reach the API. + if v.Secret != nil { + t.Errorf("volume %q is a Secret — a restore Pod must never mount a Secret", v.Name) + } + if v.ConfigMap != nil { + t.Errorf("volume %q is a ConfigMap — no config/credential injection allowed", v.Name) + } + if v.Projected != nil || v.DownwardAPI != nil { + t.Errorf("volume %q is a projected/downward volume — could surface the SA token", v.Name) + } + if v.HostPath != nil { + t.Errorf("volume %q is a hostPath — no node filesystem access allowed", v.Name) + } + if v.PersistentVolumeClaim == nil { + t.Errorf("volume %q is not a PVC; only the world and backup PVCs are permitted", v.Name) + continue + } + switch v.PersistentVolumeClaim.ClaimName { + case "world-survival-0": + world = v + case "felis-backups": + backup = v + default: + t.Errorf("unexpected PVC %q mounted", v.PersistentVolumeClaim.ClaimName) + } + } + if world == nil { + t.Fatal("world PVC not mounted") + } + if backup == nil { + t.Fatal("backup PVC not mounted") + } + // The backup PVC must be read-only at the volume source: a restore must not + // be able to mutate the archive store. + if !backup.PersistentVolumeClaim.ReadOnly { + t.Error("backup PVC volume source must be ReadOnly") + } + + // And the container's mounts must agree: backup read-only, world writable. + c := singleContainer(t, job) + var backupMount, worldMount *corev1.VolumeMount + for i := range c.VolumeMounts { + m := &c.VolumeMounts[i] + switch m.Name { + case backupVolume: + backupMount = m + case worldVolume: + worldMount = m + } + } + if backupMount == nil || !backupMount.ReadOnly { + t.Error("backup mount must be ReadOnly") + } + if worldMount == nil || worldMount.ReadOnly { + t.Error("world mount must be writable (the archive extracts into it)") + } + if backupMount != nil && backupMount.MountPath != "/backups" { + t.Errorf("backup mount path = %q, want /backups (absolute refs resolve here)", backupMount.MountPath) + } + + // Volumes are only half the red line: a single Env var (e.g. a DATABASE_URL) + // or an EnvFrom pulling a whole Secret/ConfigMap into the environment would + // hand the restore Pod a credential without ever mounting one. The container + // gets ALL of its input from the command flags, so both must be empty. + if len(c.Env) != 0 { + t.Errorf("restore container must carry no env vars, got %+v", c.Env) + } + if len(c.EnvFrom) != 0 { + t.Errorf("restore container must carry no envFrom sources (no Secret/ConfigMap injection), got %+v", c.EnvFrom) + } +} + +// A poisoned archive must terminate and not loop or run unbounded; the finished +// Job must self-GC. +func TestRestoreJobIsBoundedOneShotAndSelfCleaning(t *testing.T) { + job, err := RestoreJob(sampleJobParams()) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + if job.Spec.BackoffLimit == nil || *job.Spec.BackoffLimit != 0 { + t.Error("BackoffLimit must be 0 — a bad archive must not retry") + } + if job.Spec.ActiveDeadlineSeconds == nil || *job.Spec.ActiveDeadlineSeconds != 1800 { + t.Errorf("ActiveDeadlineSeconds must be 1800, got %v", job.Spec.ActiveDeadlineSeconds) + } + if job.Spec.TTLSecondsAfterFinished == nil || *job.Spec.TTLSecondsAfterFinished != 600 { + t.Errorf("TTLSecondsAfterFinished must be 600, got %v", job.Spec.TTLSecondsAfterFinished) + } + if job.Spec.Template.Spec.RestartPolicy != corev1.RestartPolicyNever { + t.Error("RestartPolicy must be Never") + } +} + +// The container must be non-root, non-privileged, escalation-proof, read-only +// root, drop ALL caps, and carry resource limits. +func TestRestoreJobContainerIsHardened(t *testing.T) { + job, err := RestoreJob(sampleJobParams()) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + pod := job.Spec.Template.Spec + if pod.SecurityContext == nil || pod.SecurityContext.RunAsNonRoot == nil || !*pod.SecurityContext.RunAsNonRoot { + t.Error("pod must set runAsNonRoot=true") + } + if pod.SecurityContext == nil || pod.SecurityContext.FSGroup == nil || *pod.SecurityContext.FSGroup != 1000 { + t.Error("pod must set an fsGroup so restored files are group-owned by the server identity") + } + c := singleContainer(t, job) + sc := c.SecurityContext + if sc == nil { + t.Fatal("container has no security context") + } + if sc.Privileged == nil || *sc.Privileged { + t.Error("container must not be privileged") + } + if sc.AllowPrivilegeEscalation == nil || *sc.AllowPrivilegeEscalation { + t.Error("container must set allowPrivilegeEscalation=false") + } + if sc.ReadOnlyRootFilesystem == nil || !*sc.ReadOnlyRootFilesystem { + t.Error("container must set readOnlyRootFilesystem=true (writes go only to the world PVC)") + } + if sc.Capabilities == nil || len(sc.Capabilities.Drop) == 0 || string(sc.Capabilities.Drop[0]) != "ALL" { + t.Errorf("container must drop ALL capabilities, got %v", sc.Capabilities) + } + if c.Resources.Limits.Cpu().IsZero() || c.Resources.Limits.Memory().IsZero() { + t.Error("container must carry CPU+memory limits") + } +} + +// The container must invoke `felis restore` with the archive parameters as +// plain flags — and crucially the world PVC name the operator/reaper agree on. +func TestRestoreJobInvokesFelisRestoreWithParams(t *testing.T) { + p := sampleJobParams() + job, err := RestoreJob(p) + if err != nil { + t.Fatalf("RestoreJob: %v", err) + } + c := singleContainer(t, job) + if len(c.Command) < 2 || c.Command[0] != "felis" || c.Command[1] != "restore" { + t.Errorf("command = %v, want [felis restore ...]", c.Command) + } + if !argPairPresent(c.Args, "--server", p.Server) { + t.Errorf("args must carry --server %q, got %v", p.Server, c.Args) + } + if !argPairPresent(c.Args, "--ref", p.BackupRef) { + t.Errorf("args must carry --ref %q, got %v", p.BackupRef, c.Args) + } + if !argPairPresent(c.Args, "--archive-store", p.ArchiveStore) { + t.Errorf("args must carry --archive-store %q, got %v", p.ArchiveStore, c.Args) + } + if !argPairPresent(c.Args, "--backup-root", p.BackupRoot) { + t.Errorf("args must carry --backup-root %q, got %v", p.BackupRoot, c.Args) + } + if !argPairPresent(c.Args, "--worlds-root", p.WorldsRoot) { + t.Errorf("args must carry --worlds-root %q, got %v", p.WorldsRoot, c.Args) + } + if job.Name != "restore-survival" { + t.Errorf("job name = %q, want restore-survival (deterministic for idempotency)", job.Name) + } +} + +// An empty image must be rejected rather than render an unrunnable Job; this is +// what lets cmd/felis fall back to a 503 instead of enqueuing junk. +func TestRestoreJobRequiresImage(t *testing.T) { + p := sampleJobParams() + p.Image = "" + if _, err := RestoreJob(p); err == nil { + t.Error("expected error for an empty image") + } +} + +func TestRestoreJobRejectsBadResourceLimit(t *testing.T) { + p := sampleJobParams() + p.MemLimit = "not-a-quantity" + if _, err := RestoreJob(p); err == nil { + t.Error("expected error for an unparseable memory limit") + } +} + +// The restore SA must be bare: no secrets, token automount disabled. +func TestRestoreServiceAccountIsBare(t *testing.T) { + sa := RestoreServiceAccount(defaultNamespace, defaultServiceAccount) + if sa.AutomountServiceAccountToken == nil || *sa.AutomountServiceAccountToken { + t.Error("SA must disable token automounting") + } + if len(sa.Secrets) != 0 { + t.Errorf("SA must carry no secrets, got %d", len(sa.Secrets)) + } + if len(sa.ImagePullSecrets) != 0 { + t.Errorf("SA must carry no image-pull secrets, got %d", len(sa.ImagePullSecrets)) + } +} + +// ---- helpers ---- + +func singleContainer(t *testing.T, job *batchv1.Job) corev1.Container { + t.Helper() + cs := job.Spec.Template.Spec.Containers + if len(cs) != 1 { + t.Fatalf("expected exactly one restore container, got %d", len(cs)) + } + return cs[0] +} + +func argPairPresent(args []string, flag, val string) bool { + for i := 0; i < len(args)-1; i++ { + if args[i] == flag && args[i+1] == val { + return true + } + } + return false +} diff --git a/internal/restore/k8sjobs.go b/internal/restore/k8sjobs.go new file mode 100644 index 0000000..1ed6d5a --- /dev/null +++ b/internal/restore/k8sjobs.go @@ -0,0 +1,46 @@ +package restore + +import ( + "context" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// K8sJobs is the production Jobs backed by a controller-runtime client (spec +// §7, §16). It creates the world-restore Job — nothing more: the restore Job is +// one-shot and self-cleaning (ttlSecondsAfterFinished), so there is no phase or +// cancel seam, and thus no config to hold (unlike build.K8sJobs, which needs the +// namespace to read and delete its Job). Every restore parameter arrives in the +// JobParams the Restorer builds from its own (defaulted) Config. The +// cluster-bootstrap objects (the weak felis-restore SA) are installed once by +// the deployment manifests (spec §21), not per restore, so this binding never +// creates them. It is integration-tested against a live cluster, not the +// hermetic restore_test.go suite. +type K8sJobs struct { + c client.Client +} + +// NewK8sJobs builds a Jobs over c. The restore Job's parameters all travel in +// JobParams, so there is no Config to retain here. +func NewK8sJobs(c client.Client) *K8sJobs { + return &K8sJobs{c: c} +} + +// CreateRestoreJob renders and applies the restore Job. Its name is a +// deterministic function of the server (RestoreJobName), so a concurrent restore +// of the same server collides on Create; that collision is mapped to +// ErrAlreadyExists, which the Restorer treats as success (idempotent enqueue). +func (k *K8sJobs) CreateRestoreJob(ctx context.Context, p JobParams) error { + job, err := RestoreJob(p) + if err != nil { + return err + } + if err := k.c.Create(ctx, job); err != nil { + if apierrors.IsAlreadyExists(err) { + return ErrAlreadyExists + } + return err + } + return nil +} diff --git a/internal/restore/restore.go b/internal/restore/restore.go new file mode 100644 index 0000000..e2e4062 --- /dev/null +++ b/internal/restore/restore.go @@ -0,0 +1,222 @@ +// Package restore implements the world-restore executor (spec §7 +// POST /servers/{name}/restore-backup, spec §466: a former owner who re-claims a +// released server within the retention window restores their archived world). +// +// felis-api cannot restore a world in-process: the world PVC is RWO and owned by +// the operator's StatefulSet, so the API has nothing to mount at request time +// (see internal/api.Restorer). This package is the production executor it hands +// off to — a one-shot Kubernetes Job in the minecraft namespace that mounts the +// target world PVC and the backup store, then runs `felis restore` (cmd/felis) +// to extract the archive into the world volume. +// +// Trust model, mirroring internal/build's weak-SA isolation (spec §16, §21, §22): +// the restore Pod runs under a deliberately weak service account with its token +// auto-mount disabled, so it cannot reach the K8s API; it is handed ONLY the two +// PVCs and the archive parameters as plain flags, never a database URL or any +// Secret — the four-power red line that a build/restore Pod must not touch the +// felis database or the K8s API. felis-api owns the database and the +// authorization decision (handlers_backups.go); this Pod only moves bytes from +// the backup PVC onto the world PVC. Every isolation guarantee lives in the pure +// jobspec (jobspec.go) and is asserted by jobspec_test.go, because no cluster +// runs in this environment. +// +// The Restorer depends on the Jobs interface, so the orchestration (idempotent +// enqueue, error mapping) is unit-tested against an in-memory fake; the +// controller-runtime implementation (k8sjobs.go) compiles here but is exercised +// only by integration tests against a live cluster. +package restore + +import ( + "context" + "errors" + "time" + + "felis.lolicon.best/internal/naming" +) + +// ErrAlreadyExists is returned by a Jobs implementation when a restore Job for a +// server already exists (a restore is already in flight). The Restorer treats it +// as success — see Restore. +var ErrAlreadyExists = errors.New("restore: job already exists") + +// Jobs is the cluster-side restore lifecycle the Restorer depends on. It is an +// interface so the orchestration is tested against a fake; the controller-runtime +// implementation (K8sJobs) is integration-tested only — it requires a live +// cluster. The restore Job is one-shot and self-cleaning (TTL), so unlike the +// build subsystem there is no phase-polling or cancel seam: kicking it off is the +// whole contract, exactly matching the asynchronous 202 the handler answers. +type Jobs interface { + // CreateRestoreJob renders and applies the restore Job for p. It returns + // ErrAlreadyExists if a Job of the same (deterministic) name already exists. + CreateRestoreJob(ctx context.Context, p JobParams) error +} + +// Config parameterises the restore executor. Deployment-specific values that +// have no safe default — the felis Image to run and the BackupPVC to mount — are +// supplied by the caller (cmd/felis sources them from the environment); when +// either is empty the caller leaves the API's Restorer nil so the endpoint +// reports 503 rather than enqueuing a Job that cannot run. +type Config struct { + // Namespace is where the world PVCs live and the restore Job runs (the + // minecraft namespace). The Job is intentionally co-located with the world it + // restores; it never runs in the felis control-plane namespace. + Namespace string + // ServiceAccount is the weak SA the restore Pod runs as. Like felis-build it + // MUST NOT be the felis-api SA and has no Role/RoleBinding anywhere. + ServiceAccount string + // Image is the felis binary image; the Job runs `felis restore` from it. + Image string + // ArchiveStore selects the backup backend. Only "tarLocal" is implemented in + // this build, mirroring the reaper (cmd/felis buildArchiver). + ArchiveStore string + // BackupPVC is the name of the backup PVC the archives live on. It must be + // RWX so the reaper and concurrent restores can mount it (a helm-slice + // contract); restore mounts it read-only. + BackupPVC string + // BackupRoot is the in-Pod mount path of BackupPVC. It MUST equal the path + // the reaper wrote archives under (cfg.Archive.LocalPath), because tarLocal + // archive refs are absolute paths — mounting the PVC anywhere else would make + // the stored ref unresolvable inside the Pod. + BackupRoot string + // WorldsRoot is the in-Pod mount path of the world PVC the archive extracts + // into. + WorldsRoot string + // Deadline caps the restore Pod's wall-clock (activeDeadlineSeconds). + Deadline time.Duration + // CPULimit / MemLimit cap the restore container. + CPULimit string + MemLimit string + // RunAsUser / RunAsGroup / FSGroup are the Pod's runtime identity. FSGroup in + // particular MUST match the operator StatefulSet's runtime group so the files + // the restore Pod writes are readable by the minecraft server that later + // mounts the same world PVC. The default matches the conventional minecraft + // container uid; a deployment that runs minecraft as another id overrides it. + RunAsUser int64 + RunAsGroup int64 + FSGroup int64 + // TTLAfterFinished is how long a finished restore Job lingers before the Job + // controller garbage-collects it. There is no cancel path, so the TTL is the + // only cleanup; it also bounds the window in which a re-restore sees a stale + // completed Job as ErrAlreadyExists. + TTLAfterFinished time.Duration +} + +// defaults applied when a Config field is left zero. Image and BackupPVC have no +// default on purpose — see Config. +const ( + defaultNamespace = "minecraft" + defaultServiceAccount = "felis-restore" + defaultArchiveStore = "tarLocal" + defaultBackupRoot = "/backups" + defaultWorldsRoot = "/world" + defaultDeadline = 30 * time.Minute + defaultCPULimit = "1" + defaultMemLimit = "1Gi" + defaultRunAsID = int64(1000) + defaultTTL = 10 * time.Minute +) + +// withDefaults returns a copy of c with zero fields filled, so a partially +// configured Config (or the zero value, in tests) is always usable. +func (c Config) withDefaults() Config { + if c.Namespace == "" { + c.Namespace = defaultNamespace + } + if c.ServiceAccount == "" { + c.ServiceAccount = defaultServiceAccount + } + if c.ArchiveStore == "" { + c.ArchiveStore = defaultArchiveStore + } + if c.BackupRoot == "" { + c.BackupRoot = defaultBackupRoot + } + if c.WorldsRoot == "" { + c.WorldsRoot = defaultWorldsRoot + } + if c.Deadline <= 0 { + c.Deadline = defaultDeadline + } + if c.CPULimit == "" { + c.CPULimit = defaultCPULimit + } + if c.MemLimit == "" { + c.MemLimit = defaultMemLimit + } + if c.RunAsUser == 0 { + c.RunAsUser = defaultRunAsID + } + if c.RunAsGroup == 0 { + c.RunAsGroup = defaultRunAsID + } + if c.FSGroup == 0 { + c.FSGroup = defaultRunAsID + } + if c.TTLAfterFinished <= 0 { + c.TTLAfterFinished = defaultTTL + } + return c +} + +// Restorer is the production internal/api.Restorer (the compile-time proof of +// that is in internal/api's test, which imports this package; this package never +// imports api). It holds no mutable state. +type Restorer struct { + Jobs Jobs + Config Config +} + +// Restore enqueues a restore Job that extracts the archive at backupRef into +// serverName's world PVC. It returns once the Job is created — the extraction +// runs in the Pod — so the handler's 202 ("restoring") is honest. +// +// It is idempotent: if a restore Job for this server already exists (a restore +// is already in flight, or a just-finished one has not yet hit its TTL), the +// duplicate enqueue is treated as success rather than surfaced as an error. +// +// The coalescing key is the Job name (RestoreJobName), which depends only on the +// server, NOT on backupRef — so a second request that arrives while one is in +// flight is absorbed regardless of the ref it carries, and if the two refs +// differ the second is silently dropped (the in-flight restore wins). That is +// acceptable here: restore runs only for a Stopped server (handler gate ⑥) and +// the handler always passes the latest backup, which for a stopped server does +// not change, so concurrent requests carry the same ref in practice. A caller +// that genuinely needs a different archive can re-request after the Job clears +// its TTL. This keeps the handler's 202 honest without it having to map "already +// in progress" onto a 500. +func (r *Restorer) Restore(ctx context.Context, serverName, backupRef string) error { + if err := r.Jobs.CreateRestoreJob(ctx, r.jobParams(serverName, backupRef)); err != nil { + if errors.Is(err, ErrAlreadyExists) { + return nil // already enqueued — idempotent + } + return err + } + return nil +} + +// jobParams projects the server, archive ref, and config onto the inputs +// jobspec.go renders. The world PVC name is derived from the single shared +// naming convention (naming.WorldPVCName), the same one the operator created it +// under and the reaper deletes it by. +func (r *Restorer) jobParams(serverName, backupRef string) JobParams { + cfg := r.Config.withDefaults() + return JobParams{ + Server: serverName, + WorldPVC: naming.WorldPVCName(serverName), + BackupPVC: cfg.BackupPVC, + BackupRef: backupRef, + ArchiveStore: cfg.ArchiveStore, + Namespace: cfg.Namespace, + ServiceAccount: cfg.ServiceAccount, + Image: cfg.Image, + BackupRoot: cfg.BackupRoot, + WorldsRoot: cfg.WorldsRoot, + Deadline: cfg.Deadline, + CPULimit: cfg.CPULimit, + MemLimit: cfg.MemLimit, + RunAsUser: cfg.RunAsUser, + RunAsGroup: cfg.RunAsGroup, + FSGroup: cfg.FSGroup, + TTLAfterFinished: cfg.TTLAfterFinished, + } +} diff --git a/internal/restore/restore_test.go b/internal/restore/restore_test.go new file mode 100644 index 0000000..56887d3 --- /dev/null +++ b/internal/restore/restore_test.go @@ -0,0 +1,84 @@ +package restore_test + +import ( + "context" + "errors" + "testing" + + "felis.lolicon.best/internal/restore" +) + +// fakeJobs is an in-memory Jobs that records the params it was handed and +// returns a programmable error, so the orchestration is tested without a +// cluster. +type fakeJobs struct { + calls []restore.JobParams + err error +} + +func (f *fakeJobs) CreateRestoreJob(_ context.Context, p restore.JobParams) error { + f.calls = append(f.calls, p) + return f.err +} + +func TestRestoreEnqueuesJobWithDerivedParams(t *testing.T) { + jobs := &fakeJobs{} + r := &restore.Restorer{ + Jobs: jobs, + Config: restore.Config{ + Image: "registry.internal/felis:test", + BackupPVC: "felis-backups", + }, + } + + if err := r.Restore(context.Background(), "survival", "/backups/survival/2026.tar.gz"); err != nil { + t.Fatalf("Restore: %v", err) + } + if len(jobs.calls) != 1 { + t.Fatalf("CreateRestoreJob called %d times, want 1", len(jobs.calls)) + } + got := jobs.calls[0] + if got.Server != "survival" { + t.Errorf("Server = %q, want survival", got.Server) + } + // The world PVC must come from the shared naming convention, not be invented + // here: it is the same name the operator created and the reaper deletes. + if got.WorldPVC != "world-survival-0" { + t.Errorf("WorldPVC = %q, want world-survival-0", got.WorldPVC) + } + if got.BackupRef != "/backups/survival/2026.tar.gz" { + t.Errorf("BackupRef = %q, want the passed ref", got.BackupRef) + } + if got.BackupPVC != "felis-backups" { + t.Errorf("BackupPVC = %q, want felis-backups", got.BackupPVC) + } + if got.Image != "registry.internal/felis:test" { + t.Errorf("Image = %q, want the configured image", got.Image) + } + // withDefaults must have filled the unset fields. + if got.Namespace == "" || got.ServiceAccount == "" || got.WorldsRoot == "" || got.BackupRoot == "" { + t.Errorf("defaults not applied: %+v", got) + } +} + +func TestRestoreIsIdempotentOnAlreadyExists(t *testing.T) { + jobs := &fakeJobs{err: restore.ErrAlreadyExists} + r := &restore.Restorer{Jobs: jobs, Config: restore.Config{Image: "img", BackupPVC: "pvc"}} + + // A restore already in flight is success, not an error: the handler must be + // able to answer 202 for a coalesced duplicate request. + if err := r.Restore(context.Background(), "survival", "ref"); err != nil { + t.Fatalf("Restore on AlreadyExists = %v, want nil (idempotent)", err) + } +} + +func TestRestorePropagatesGenericError(t *testing.T) { + sentinel := errors.New("apiserver exploded") + jobs := &fakeJobs{err: sentinel} + r := &restore.Restorer{Jobs: jobs, Config: restore.Config{Image: "img", BackupPVC: "pvc"}} + + err := r.Restore(context.Background(), "survival", "ref") + if !errors.Is(err, sentinel) { + t.Fatalf("Restore error = %v, want the underlying error propagated", err) + } +}