feat(backup): add backup, restore, and reaper subsystems
Archive-based world backup and restore, plus the reaper that enforces retention and reclaims idle servers.
This commit is contained in:
12 files changed
+2749
No files matched your search
@@ -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-<name>-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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user