feat(backup): 玩过的世界每天在停服后自动打一个定时恢复点,按服按 owner 保留 7 个、90 天过期

This commit is contained in:
Lemon-miaow committed 2026-09-27 00:02:49 +08:00
1 parent d6e1a6464c
commit 8ef7112dab
29 files changed
+1028 -59

No files matched your search

+3
View File
@@ -26,6 +26,9 @@ type AsyncJob struct {
ThenRestore string `json:"then_restore,omitempty"`
ThenRestoreReason string `json:"then_restore_reason,omitempty"`
RestoreBackupID string `json:"restore_backup_id,omitempty"`
// Scheduled marks a backup felis-api took on its own (BackupScheduler), so
// the owner can tell it from one somebody asked for.
Scheduled bool `json:"scheduled,omitempty"`
}
// JobStatusReader reads the newest backup/restore Jobs for a server, newest
+23
View File
@@ -25,6 +25,11 @@ const (
jobManagedByBackup = "felis-backup"
jobManagedByRestore = "felis-restore"
// jobBackupReasonLabel marks the backup Job of a scheduled restore point
// (backupjob.LabelReason).
jobBackupReasonLabel = "felis.lolicon.best/backup-reason"
jobBackupReasonScheduled = "scheduled"
)
// K8sJobStatus reads the async Jobs the executors created, by the server label
@@ -137,9 +142,27 @@ func jobToAsyncJob(j *batchv1.Job) (AsyncJob, bool) {
}
}
}
aj.Scheduled = aj.Kind == "backup" && j.Labels[jobBackupReasonLabel] == jobBackupReasonScheduled
return aj, true
}
// RunningWorldJobs counts the backup and restore Jobs of every server that
// have yet to finish (BackupScheduler waits for them).
func (k *K8sJobStatus) RunningWorldJobs(ctx context.Context) (int, error) {
var list batchv1.JobList
if err := k.c.List(ctx, &list, client.InNamespace(k.namespace), client.HasLabels{jobManagedByLabel}); err != nil {
return 0, err
}
n := 0
for i := range list.Items {
j := &list.Items[i]
if _, ok := jobOutcome(j); ok && !maintenance.JobFinished(j) {
n++
}
}
return n, nil
}
// PendingRestoreChains lists the safety snapshots felis-api has yet to settle,
// across every server (the label selector keeps it to them).
func (k *K8sJobStatus) PendingRestoreChains(ctx context.Context) ([]RestoreChain, error) {
+33
View File
@@ -782,6 +782,39 @@ func (p *PGRepo) BackupStoreBytes(ctx context.Context) (int64, error) {
return n, err
}
// ScheduledBackupCandidates is the ScheduleStore behind BackupScheduler. A
// world's point is its current owner's newest intact scheduled backup taken
// since they claimed it: a previous owner's backups say nothing about the world
// the new owner has built, and a corrupt one restores nothing. last_active_at
// moves on every join, so a world nobody joined since its point is skipped.
func (p *PGRepo) ScheduledBackupCandidates(ctx context.Context, before time.Time) ([]ScheduledCandidate, error) {
rows, err := p.db.QueryContext(ctx,
`SELECT s.name, s.owner_id FROM servers s
LEFT JOIN LATERAL (
SELECT max(b.created_at) AS at FROM world_backups b
WHERE b.server_name = s.name AND b.reason = 'scheduled' AND b.status = 'present'
AND b.corrupt_at IS NULL AND b.former_owner = s.owner_id
AND b.created_at >= COALESCE(s.claimed_at, '-infinity')
) pt ON true
WHERE s.deleted_at IS NULL AND s.owner_id IS NOT NULL
AND s.last_active_at > COALESCE(pt.at, '-infinity')
AND COALESCE(pt.at, '-infinity') < $1
ORDER BY pt.at NULLS FIRST, s.name`, before)
if err != nil {
return nil, err
}
defer rows.Close()
var out []ScheduledCandidate
for rows.Next() {
var c ScheduledCandidate
if err := rows.Scan(&c.Name, &c.OwnerID); err != nil {
return nil, err
}
out = append(out, c)
}
return out, rows.Err()
}
func (p *PGRepo) Audit(ctx context.Context, e AuditEntry) error {
// A nil Payload must land as SQL NULL, not the text "null"; a non-nil Payload is
// passed as a JSON text the jsonb column parses (same idiom as reaper.PGStore).
+170
View File
@@ -0,0 +1,170 @@
package api
import (
"context"
"errors"
"fmt"
"log"
"time"
"felis.lolicon.best/internal/maintenance"
)
// Scheduled backups. A world played every day never idles long enough for the
// reaper to archive it, and an owner's own backups are only the ones they
// remembered to take, so felis-api takes a daily restore point of every world
// played since its last one. The backup Job mounts the world volume, which a
// running server holds, so the point is taken once the server is stopped: idle
// auto-stop brings a played world down minutes after its last player leaves,
// and the point lands the same day. A server kept running around the clock gets
// one the next time it stops.
//
// The point is an ordinary world_backups row of reason scheduled: listed and
// restorable by its owner, copied offsite like any other, pruned to [archive]
// scheduled_keep per owner and expired after scheduled_retention, so it never
// takes the place of a backup the owner asked for.
// ScheduledBackuper enqueues the backup Job of a scheduled restore point
// (backupjob.Backuper).
type ScheduledBackuper interface {
BackupScheduled(ctx context.Context, serverName, formerOwner string) error
}
// ScheduledCandidate is an owned world due a scheduled backup.
type ScheduledCandidate struct {
Name string
OwnerID string
}
// ScheduleStore lists the worlds due a scheduled backup: owned, joined since
// their owner's newest intact scheduled backup, and without one taken after
// before. Worlds never given one come first, then the longest waiting.
type ScheduleStore interface {
ScheduledBackupCandidates(ctx context.Context, before time.Time) ([]ScheduledCandidate, error)
}
// WorldJobCounter counts the backup and restore Jobs still running.
type WorldJobCounter interface {
RunningWorldJobs(ctx context.Context) (int, error)
}
// Audit identity of a scheduled backup. The action is its own so the owner's
// manual cooldown (LastBackupRequest, backup.create) never counts it.
const (
scheduledBackupActor = "scheduler"
scheduledBackupAction = "backup.scheduled"
)
// BackupScheduler takes the scheduled backups. Tick launches at most one
// backup Job and only while no backup or restore Job is running, so the points
// queue behind each other and behind the owners' own operations instead of
// loading the node with archives all at once.
type BackupScheduler struct {
API *API
Store ScheduleStore
Jobs WorldJobCounter
// Every is how often a played world gets a point ([archive]
// scheduled_every). Zero or less turns the scheduler off.
Every time.Duration
// tried is when each world last had a Job launched or was found without a
// volume, so one whose backup keeps failing is retried every Every/4
// instead of every tick.
tried map[string]time.Time
// full remembers that the store was at max_local_bytes, so the pause and
// the resume are logged once each.
full bool
}
// Tick launches the backup of the world waiting longest, if any may run now.
func (s *BackupScheduler) Tick(ctx context.Context) error {
b, ok := s.API.Backuper.(ScheduledBackuper)
if !ok || s.Every <= 0 {
return nil
}
now := s.API.now()
running, err := s.Jobs.RunningWorldJobs(ctx)
if err != nil {
return fmt.Errorf("count running backup and restore jobs: %w", err)
}
if running > 0 {
return nil
}
// A full store is the reaper's to evict; a scheduled point added now would
// only push out an older backup of somebody else's.
if limit := s.API.BackupStoreCap; limit > 0 {
used, err := s.API.Repo.BackupStoreBytes(ctx)
if err != nil {
return fmt.Errorf("read the backup store size: %w", err)
}
full := used >= limit
if full != s.full {
if full {
log.Printf("api: scheduled backups paused: the backup store holds %d of its %d bytes ([archive] max_local_bytes)", used, limit)
} else {
log.Printf("api: scheduled backups resumed: the backup store is below [archive] max_local_bytes")
}
s.full = full
}
if full {
return nil
}
}
due, err := s.Store.ScheduledBackupCandidates(ctx, now.Add(-s.Every))
if err != nil {
return fmt.Errorf("list the worlds due a scheduled backup: %w", err)
}
if s.tried == nil {
s.tried = map[string]time.Time{}
}
for name, at := range s.tried {
if now.Sub(at) >= s.Every {
delete(s.tried, name)
}
}
for _, c := range due {
if at, ok := s.tried[c.Name]; ok && now.Sub(at) < s.Every/4 {
continue
}
// A server never started has no world yet, and one whose volume is gone
// has nothing left to save; the Job would sit Pending on the claim.
exists, err := s.API.Cluster.WorldVolumeExists(ctx, c.Name)
if err != nil {
return fmt.Errorf("look up the world volume of %s: %w", c.Name, err)
}
if !exists {
s.tried[c.Name] = now
continue
}
var busy *MaintenanceBusyError
switch err := s.API.Cluster.AcquireMaintenance(ctx, c.Name, maintenance.KindBackup); {
case errors.Is(err, ErrNotStopped), errors.As(err, &busy), errors.Is(err, ErrMaintenanceInProgress):
continue // running, or somebody else has the world: next tick
case errors.Is(err, ErrNotFound):
s.tried[c.Name] = now
continue
case err != nil:
return fmt.Errorf("lock the world of %s: %w", c.Name, err)
}
err = b.BackupScheduled(ctx, c.Name, c.OwnerID)
// Once the Job exists it holds the world; the annotation only covered
// the gap.
if rerr := s.API.Cluster.ReleaseMaintenance(context.WithoutCancel(ctx), c.Name); rerr != nil {
log.Printf("api: release the maintenance lock on %s: %v (it lapses after %s)", c.Name, rerr, maintenance.Grace)
}
s.tried[c.Name] = now
if err != nil {
return fmt.Errorf("start the scheduled backup of %s: %w", c.Name, err)
}
log.Printf("api: started the scheduled backup of %s", c.Name)
s.API.writeAudit(ctx, AuditEntry{
Actor: scheduledBackupActor, Source: scheduledBackupActor,
Action: scheduledBackupAction, ServerName: c.Name,
})
return nil
}
return nil
}
+263
View File
@@ -0,0 +1,263 @@
package api
import (
"context"
"errors"
"testing"
"time"
"felis.lolicon.best/internal/backupjob"
"felis.lolicon.best/internal/maintenance"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
var _ ScheduledBackuper = (*backupjob.Backuper)(nil)
// fakeScheduledBackuper is a Backuper that can also take a scheduled backup.
type fakeScheduledBackuper struct {
fakeBackuper
err error
scheduled []ScheduledCandidate
}
func (f *fakeScheduledBackuper) BackupScheduled(_ context.Context, name, formerOwner string) error {
f.scheduled = append(f.scheduled, ScheduledCandidate{Name: name, OwnerID: formerOwner})
return f.err
}
type fakeScheduleStore struct {
due []ScheduledCandidate
gotBefore time.Time
calls int
}
func (f *fakeScheduleStore) ScheduledBackupCandidates(_ context.Context, before time.Time) ([]ScheduledCandidate, error) {
f.calls++
f.gotBefore = before
return f.due, nil
}
type fakeWorldJobs struct {
running int
err error
}
func (f *fakeWorldJobs) RunningWorldJobs(context.Context) (int, error) { return f.running, f.err }
func TestBackupScheduler(t *testing.T) {
const every = 24 * time.Hour
type rig struct {
a *API
repo *fakeRepo
cl *fakeCluster
b *fakeScheduledBackuper
store *fakeScheduleStore
jobs *fakeWorldJobs
s *BackupScheduler
}
mk := func() rig {
repo := newFakeRepo()
cl := newFakeCluster()
for _, n := range []string{"alpha", "bravo"} {
cl.byName[n] = &ServerInfo{Name: n, Phase: "Stopped"}
}
a := newTestAPI(repo, cl)
b := &fakeScheduledBackuper{}
a.Backuper = b
store := &fakeScheduleStore{due: []ScheduledCandidate{{"alpha", "usr-a"}, {"bravo", "usr-b"}}}
jobs := &fakeWorldJobs{}
return rig{a, repo, cl, b, store, jobs, &BackupScheduler{API: a, Store: store, Jobs: jobs, Every: every}}
}
tick := func(t *testing.T, r rig) {
t.Helper()
if err := r.s.Tick(context.Background()); err != nil {
t.Fatalf("Tick: %v", err)
}
}
launched := func(r rig) []ScheduledCandidate { return r.b.scheduled }
t.Run("backs up the first due world as its owner, one per tick", func(t *testing.T) {
r := mk()
tick(t, r)
if got := launched(r); len(got) != 1 || got[0] != (ScheduledCandidate{"alpha", "usr-a"}) {
t.Fatalf("launched = %+v; want alpha as usr-a only", got)
}
if want := r.a.now().Add(-every); !r.store.gotBefore.Equal(want) {
t.Fatalf("asked for points before %v; want %v", r.store.gotBefore, want)
}
if len(r.cl.acquired) != 1 || r.cl.acquired[0] != "alpha:"+maintenance.KindBackup ||
len(r.cl.released) != 1 || r.cl.released[0] != "alpha" {
t.Fatalf("lock: acquired %v released %v; want alpha held as a backup and let go", r.cl.acquired, r.cl.released)
}
if len(r.repo.audits) != 1 || r.repo.audits[0].Action != "backup.scheduled" || r.repo.audits[0].ServerName != "alpha" {
t.Fatalf("audits = %+v; want one backup.scheduled of alpha", r.repo.audits)
}
if at, _ := r.repo.LastBackupRequest(context.Background(), "alpha", time.Time{}); !at.IsZero() {
t.Fatal("the scheduled backup started the owner's manual cooldown")
}
if r.b.calls != 0 {
t.Fatal("the scheduled backup went through the manual Backup")
}
})
t.Run("a running or busy world waits for the next tick", func(t *testing.T) {
for _, err := range []error{ErrNotStopped, &MaintenanceBusyError{Kind: maintenance.KindRestore}, ErrMaintenanceInProgress} {
r := mk()
r.cl.maintErr["alpha"] = err
tick(t, r)
if got := launched(r); len(got) != 1 || got[0].Name != "bravo" {
t.Fatalf("%v: launched = %+v; want bravo", err, got)
}
delete(r.cl.maintErr, "alpha")
r.b.scheduled = nil
tick(t, r)
if got := launched(r); len(got) != 1 || got[0].Name != "alpha" {
t.Fatalf("%v: once alpha is free, launched = %+v; want alpha", err, got)
}
}
})
t.Run("waits while a backup or restore job runs", func(t *testing.T) {
r := mk()
r.jobs.running = 1
tick(t, r)
if len(launched(r)) != 0 || len(r.cl.acquired) != 0 {
t.Fatalf("launched %+v, acquired %v beside a running job", launched(r), r.cl.acquired)
}
r.jobs.running = 0
tick(t, r)
if len(launched(r)) != 1 {
t.Fatalf("launched = %+v once the job finished", launched(r))
}
})
t.Run("pauses while the store is full", func(t *testing.T) {
r := mk()
r.a.BackupStoreCap = 100
r.repo.backups = []fakeBackup{{view: BackupView{ID: "bk1", ServerName: "bravo", Status: "present", SizeBytes: 100}}}
tick(t, r)
if len(launched(r)) != 0 {
t.Fatalf("launched = %+v into a full store", launched(r))
}
r.repo.backups[0].view.SizeBytes = 99
tick(t, r)
if len(launched(r)) != 1 {
t.Fatalf("launched = %+v below the cap", launched(r))
}
})
t.Run("a world whose backup failed is retried after a quarter period", func(t *testing.T) {
r := mk()
r.store.due = r.store.due[:1]
r.b.err = errors.New("apiserver down")
if err := r.s.Tick(context.Background()); err == nil {
t.Fatal("Tick hid the failed launch")
}
if len(r.cl.released) != 1 {
t.Fatalf("released = %v; the lock must go when the launch fails", r.cl.released)
}
r.b.err = nil
start := r.a.now()
r.a.Now = func() time.Time { return start.Add(every/4 - time.Minute) }
tick(t, r)
if len(launched(r)) != 1 {
t.Fatalf("launched = %+v; retried before a quarter period", launched(r))
}
r.a.Now = func() time.Time { return start.Add(every / 4) }
tick(t, r)
if len(launched(r)) != 2 {
t.Fatalf("launched = %+v; not retried after a quarter period", launched(r))
}
})
t.Run("a world without a volume is passed over", func(t *testing.T) {
r := mk()
r.cl.noWorld["alpha"] = true
tick(t, r)
if got := launched(r); len(got) != 1 || got[0].Name != "bravo" || len(r.cl.acquired) != 1 {
t.Fatalf("launched = %+v, acquired %v; want bravo alone", got, r.cl.acquired)
}
})
t.Run("a world deleted since the listing is passed over", func(t *testing.T) {
r := mk()
r.cl.maintErr["alpha"] = ErrNotFound
tick(t, r)
if got := launched(r); len(got) != 1 || got[0].Name != "bravo" {
t.Fatalf("launched = %+v; want bravo", got)
}
})
t.Run("a lock failure stops the tick", func(t *testing.T) {
r := mk()
r.cl.maintErr["alpha"] = errors.New("conflict storm")
if err := r.s.Tick(context.Background()); err == nil || len(launched(r)) != 0 {
t.Fatalf("Tick = %v, launched %+v; want the error and nothing started", err, launched(r))
}
})
t.Run("off without a scheduling backuper or a period", func(t *testing.T) {
r := mk()
r.a.Backuper = &fakeBackuper{}
tick(t, r)
r2 := mk()
r2.s.Every = 0
tick(t, r2)
if r.store.calls != 0 || r2.store.calls != 0 || len(r2.b.scheduled) != 0 {
t.Fatal("the scheduler ran while off")
}
})
}
// The jobs route marks the executor's scheduled backups, and the scheduler
// counts every unfinished backup and restore Job and nothing else.
func TestK8sScheduledBackupJobs(t *testing.T) {
scheme := runtime.NewScheme()
if err := clientgoscheme.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
params := func(name string, scheduled bool) backupjob.JobParams {
return backupjob.JobParams{
Server: "survival", JobName: name, WorldPVC: "world-survival-0", BackupPVC: "felis-backups",
Namespace: "minecraft", Image: "felis:1", ConfigSecret: "felis-config", ConfigMount: "/etc/felis",
Scheduled: scheduled,
}
}
scheduled, err := backupjob.BackupJob(params("backup-survival-aa", true))
if err != nil {
t.Fatal(err)
}
plain, err := backupjob.BackupJob(params("backup-survival-bb", false))
if err != nil {
t.Fatal(err)
}
plain.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}}
restoring := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Namespace: "minecraft", Name: "restore-other-cc",
Labels: map[string]string{jobServerLabel: "other", jobManagedByLabel: jobManagedByRestore}}}
foreign := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Namespace: "minecraft", Name: "files-survival-dd",
Labels: map[string]string{jobServerLabel: "survival", jobManagedByLabel: "felis-files"}}}
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(scheduled, plain, restoring, foreign).
WithStatusSubresource(&batchv1.Job{}).Build()
k := NewK8sJobStatus(c, "minecraft")
ctx := context.Background()
if n, err := k.RunningWorldJobs(ctx); err != nil || n != 2 {
t.Fatalf("RunningWorldJobs = %d, %v; want the scheduled backup and the restore", n, err)
}
jobs, err := k.LatestJobs(ctx, "survival")
if err != nil {
t.Fatal(err)
}
marks := map[string]bool{}
for _, j := range jobs {
marks[j.Name] = j.Scheduled
}
if len(marks) != 2 || !marks["backup-survival-aa"] || marks["backup-survival-bb"] {
t.Fatalf("scheduled marks = %v; want only backup-survival-aa", marks)
}
}