fix(backup): 手动备份按服冷却、每服保留上限与独立保留期,容量驱逐不再删除回收世界的唯一副本
This commit is contained in:
30 files changed
+801
-42
No files matched your search
@@ -128,6 +128,14 @@ type API struct {
|
||||
// on the wake lever). Zero disables throttling.
|
||||
WakeCooldown time.Duration
|
||||
|
||||
// BackupCooldown spaces out an owner's on-demand backups of one server, and
|
||||
// BackupStoreCap refuses them once the present backups reach [archive]
|
||||
// max_local_bytes (data-durability-9): each archive lands on the node disk
|
||||
// the worlds and the database share. Admins and the break-glass console are
|
||||
// exempt. Zero disables each lever.
|
||||
BackupCooldown time.Duration
|
||||
BackupStoreCap int64
|
||||
|
||||
// SubmitCreateCooldown / SubmitUploadCooldown throttle the user-modpack
|
||||
// submission lane per user: create bounds how quickly review-queue rows can
|
||||
// appear, upload bounds how often a user may stream a (up to 1 GiB) build
|
||||
|
||||
@@ -48,6 +48,9 @@ type fakeRepo struct {
|
||||
resourceUpdates map[string]ResourceSpec
|
||||
audits []AuditEntry
|
||||
failAudit error // Audit fails with it (a store outage)
|
||||
// backupRequested mirrors the newest backup.create audit row per server,
|
||||
// stamped by Audit with the wall clock (LastBackupRequest).
|
||||
backupRequested map[string]time.Time
|
||||
joins []string
|
||||
// create-server seeding (spec §15)
|
||||
seeded map[string]bool // name -> servers row exists
|
||||
@@ -750,9 +753,32 @@ func (f *fakeRepo) Audit(_ context.Context, e AuditEntry) error {
|
||||
return f.failAudit
|
||||
}
|
||||
f.audits = append(f.audits, e)
|
||||
if e.Action == "backup.create" {
|
||||
if f.backupRequested == nil {
|
||||
f.backupRequested = map[string]time.Time{}
|
||||
}
|
||||
f.backupRequested[e.ServerName] = time.Now()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeRepo) LastBackupRequest(_ context.Context, serverName string, since time.Time) (time.Time, error) {
|
||||
if at, ok := f.backupRequested[serverName]; ok && !at.Before(since) {
|
||||
return at, nil
|
||||
}
|
||||
return time.Time{}, nil
|
||||
}
|
||||
|
||||
func (f *fakeRepo) BackupStoreBytes(context.Context) (int64, error) {
|
||||
var n int64
|
||||
for _, b := range f.backups {
|
||||
if b.view.Status == "present" {
|
||||
n += b.view.SizeBytes
|
||||
}
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// AllBackups / BackupsForUser / LatestBackup mirror the PG queries' contract so
|
||||
// the hermetic tests can't pass against a too-lenient fake: only status='present'
|
||||
// rows are visible, the user scope is the former_owner column, and LatestBackup
|
||||
|
||||
@@ -5,7 +5,9 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
)
|
||||
@@ -191,6 +193,78 @@ func TestBackupNow(t *testing.T) {
|
||||
t.Fatalf("code = %d, want 400", w.Code)
|
||||
}
|
||||
})
|
||||
|
||||
// data-durability-9: an owner's backups are rationed per server; an admin's
|
||||
// are not.
|
||||
t.Run("owner inside the cooldown -> 429 backup_cooldown with Retry-After", func(t *testing.T) {
|
||||
api, _, _, backuper := mk()
|
||||
api.BackupCooldown = 10 * time.Minute
|
||||
api.Now = time.Now // the fake stamps backup.create audits with the wall clock
|
||||
api.External = staticExternal{p: owner}
|
||||
if w := do(api.ExternalHandler(), "POST", path, "", nil); w.Code != http.StatusAccepted {
|
||||
t.Fatalf("first backup: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
w := do(api.ExternalHandler(), "POST", path, "", nil)
|
||||
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "backup_cooldown" {
|
||||
t.Fatalf("second backup: code = %d body %s", w.Code, w.Body.String())
|
||||
}
|
||||
if ra, _ := strconv.Atoi(w.Header().Get("Retry-After")); ra < 590 || ra > 600 {
|
||||
t.Fatalf("Retry-After = %q, want about 600", w.Header().Get("Retry-After"))
|
||||
}
|
||||
if backuper.calls != 1 {
|
||||
t.Fatalf("backuper called %d times, want 1", backuper.calls)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("cooldown elapsed -> 202", func(t *testing.T) {
|
||||
api, repo, _, _ := mk()
|
||||
api.BackupCooldown = 10 * time.Minute
|
||||
api.Now = time.Now
|
||||
repo.backupRequested = map[string]time.Time{"survival": time.Now().Add(-11 * time.Minute)}
|
||||
api.External = staticExternal{p: owner}
|
||||
if w := do(api.ExternalHandler(), "POST", path, "", nil); w.Code != http.StatusAccepted {
|
||||
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("admin bypasses the cooldown and the store cap", func(t *testing.T) {
|
||||
api, repo, _, backuper := mk()
|
||||
api.BackupCooldown = 10 * time.Minute
|
||||
api.BackupStoreCap = 100
|
||||
api.Now = time.Now
|
||||
repo.backupRequested = map[string]time.Time{"survival": time.Now()}
|
||||
repo.backups = []fakeBackup{{view: BackupView{ID: "b1", ServerName: "other", Status: "present", SizeBytes: 500}}}
|
||||
api.External = staticExternal{p: &Principal{UserID: "admin1", Email: "[email protected]",
|
||||
Role: "admin", ViaAdminAccess: true}}
|
||||
if w := do(api.ExternalHandler(), "POST", path, "", nil); w.Code != http.StatusAccepted {
|
||||
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if backuper.calls != 1 {
|
||||
t.Fatal("the admin's backup did not start")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("owner with the store at its cap -> 507 backup_store_full", func(t *testing.T) {
|
||||
api, repo, _, backuper := mk()
|
||||
api.BackupStoreCap = 1000
|
||||
repo.backups = []fakeBackup{
|
||||
{view: BackupView{ID: "b1", ServerName: "other", Status: "present", SizeBytes: 600}},
|
||||
{view: BackupView{ID: "b2", ServerName: "survival", Status: "present", SizeBytes: 400}},
|
||||
{view: BackupView{ID: "b3", ServerName: "survival", Status: "deleted", SizeBytes: 9000}},
|
||||
}
|
||||
api.External = staticExternal{p: owner}
|
||||
w := do(api.ExternalHandler(), "POST", path, "", nil)
|
||||
if w.Code != http.StatusInsufficientStorage || decodeErr(t, w) != "backup_store_full" {
|
||||
t.Fatalf("code = %d body %s", w.Code, w.Body.String())
|
||||
}
|
||||
if backuper.calls != 0 {
|
||||
t.Fatal("a full store still started a backup")
|
||||
}
|
||||
repo.backups[0].view.Status = "deleted"
|
||||
if w := do(api.ExternalHandler(), "POST", path, "", nil); w.Code != http.StatusAccepted {
|
||||
t.Fatalf("below the cap: code = %d (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// TestInternalBackup exercises POST /api/v1/internal/servers/{name}/backup, the
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||
"felis.lolicon.best/internal/maintenance"
|
||||
@@ -255,10 +257,49 @@ func (a *API) handleBackupNow(w http.ResponseWriter, r *http.Request) {
|
||||
writeError(w, r, errForbidden)
|
||||
return
|
||||
}
|
||||
if !p.IsAdmin() {
|
||||
if err := a.backupAllowance(r.Context(), name); err != nil {
|
||||
writeError(w, r, err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
a.enqueueBackup(w, r, name, rec, auditActor(p), "external")
|
||||
}
|
||||
|
||||
// backupAllowance rations an owner's on-demand backups (data-durability-9):
|
||||
// one per BackupCooldown per server, and none while the present backups fill
|
||||
// BackupStoreCap. The owner's older backups are pruned by the Job itself
|
||||
// ([archive] manual_keep), so these two gates bound the rate and the total.
|
||||
func (a *API) backupAllowance(ctx context.Context, name string) error {
|
||||
if a.BackupCooldown > 0 {
|
||||
now := a.now()
|
||||
last, err := a.Repo.LastBackupRequest(ctx, name, now.Add(-a.BackupCooldown))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !last.IsZero() {
|
||||
wait := last.Add(a.BackupCooldown).Sub(now)
|
||||
if wait > 0 {
|
||||
return newError(http.StatusTooManyRequests, "backup_cooldown",
|
||||
"a backup of this server was started %s ago; the next one can start in %s",
|
||||
now.Sub(last).Round(time.Second), wait.Round(time.Second)).retryAfter(wait)
|
||||
}
|
||||
}
|
||||
}
|
||||
if a.BackupStoreCap > 0 {
|
||||
used, err := a.Repo.BackupStoreBytes(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if used >= a.BackupStoreCap {
|
||||
return newError(http.StatusInsufficientStorage, "backup_store_full",
|
||||
"the backup store is full; ask an administrator to free space")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// handleInternalBackup is the internal-face backup trigger. The break-glass console
|
||||
// (root on the node, holding the service token) POSTs here to snapshot a stopped
|
||||
// world while felis-api is alive — it goes through the API rather than direct-to-CRD
|
||||
|
||||
@@ -11,6 +11,9 @@ import (
|
||||
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"
|
||||
)
|
||||
|
||||
type fakeJobStatus struct {
|
||||
@@ -145,3 +148,58 @@ func TestJobToAsyncJob(t *testing.T) {
|
||||
t.Fatal("foreign job must be dropped")
|
||||
}
|
||||
}
|
||||
|
||||
// TestLatestJobsExplainsFailures: a failed Job reports the error its container
|
||||
// exited on (the last line of the terminated message), newest pod first, and
|
||||
// keeps the condition text when no pod explains it.
|
||||
func TestLatestJobsExplainsFailures(t *testing.T) {
|
||||
scheme := runtime.NewScheme()
|
||||
if err := clientgoscheme.AddToScheme(scheme); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
labels := func(job string) map[string]string {
|
||||
return map[string]string{jobServerLabel: "survival", jobManagedByLabel: jobManagedByBackup, "job-name": job}
|
||||
}
|
||||
at := func(min int) metav1.Time { return metav1.NewTime(time.Date(2026, 9, 24, 10, min, 0, 0, time.UTC)) }
|
||||
failedJob := func(name string, min int) *batchv1.Job {
|
||||
start := at(min)
|
||||
return &batchv1.Job{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "minecraft", Labels: labels(name)},
|
||||
Status: batchv1.JobStatus{StartTime: &start, Conditions: []batchv1.JobCondition{{
|
||||
Type: batchv1.JobFailed, Status: corev1.ConditionTrue,
|
||||
Reason: "BackoffLimitExceeded", Message: "Job has reached the specified backoff limit",
|
||||
}}},
|
||||
}
|
||||
}
|
||||
pod := func(name, job string, min int, exit int32, msg string) *corev1.Pod {
|
||||
return &corev1.Pod{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "minecraft", Labels: labels(job), CreationTimestamp: at(min)},
|
||||
Status: corev1.PodStatus{ContainerStatuses: []corev1.ContainerStatus{{
|
||||
Name: "backup", State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{ExitCode: exit, Message: msg}},
|
||||
}}},
|
||||
}
|
||||
}
|
||||
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(
|
||||
failedJob("backup-survival-a", 1),
|
||||
failedJob("backup-survival-b", 2),
|
||||
pod("a-1", "backup-survival-a", 1, 1, "felis backup: first try\n"),
|
||||
pod("a-2", "backup-survival-a", 3, 1,
|
||||
"archiving survival\nfelis backup: backup: not enough free disk for the archive: the world is 2.0 GiB\n"),
|
||||
pod("b-1", "backup-survival-b", 2, 0, "done"),
|
||||
).Build()
|
||||
|
||||
jobs, err := NewK8sJobStatus(c, "minecraft").LatestJobs(context.Background(), "survival")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got := map[string]string{}
|
||||
for _, j := range jobs {
|
||||
got[j.Name] = j.Message
|
||||
}
|
||||
if want := "felis backup: backup: not enough free disk for the archive: the world is 2.0 GiB"; got["backup-survival-a"] != want {
|
||||
t.Errorf("a: message = %q, want %q", got["backup-survival-a"], want)
|
||||
}
|
||||
if want := "Job has reached the specified backoff limit"; got["backup-survival-b"] != want {
|
||||
t.Errorf("b: message = %q, want the condition text", got["backup-survival-b"])
|
||||
}
|
||||
}
|
||||
@@ -3,6 +3,8 @@ package api
|
||||
import (
|
||||
"context"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
@@ -53,9 +55,66 @@ func (k *K8sJobStatus) LatestJobs(ctx context.Context, serverName string) ([]Asy
|
||||
if len(out) > 20 {
|
||||
out = out[:20]
|
||||
}
|
||||
k.explainFailures(ctx, serverName, out)
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// explainFailures replaces a failed Job's condition text ("Job has reached the
|
||||
// specified backoff limit") with the error its container exited on. The
|
||||
// executors set TerminationMessagePolicy FallbackToLogsOnError, so the
|
||||
// terminated state carries the tail of the log, whose last line is the
|
||||
// "felis backup: …" / "felis restore: …" error. The pods carry the Job's
|
||||
// labels, so one list covers every Job of the server; a pod already collected
|
||||
// by the Job TTL, or a list error, leaves the condition text in place.
|
||||
func (k *K8sJobStatus) explainFailures(ctx context.Context, serverName string, jobs []AsyncJob) {
|
||||
failed := map[string]int{}
|
||||
for i, j := range jobs {
|
||||
if j.State == "failed" {
|
||||
failed[j.Name] = i
|
||||
}
|
||||
}
|
||||
if len(failed) == 0 {
|
||||
return
|
||||
}
|
||||
var pods corev1.PodList
|
||||
if err := k.c.List(ctx, &pods, client.InNamespace(k.namespace),
|
||||
client.MatchingLabels{jobServerLabel: serverName}); err != nil {
|
||||
return
|
||||
}
|
||||
newest := map[string]time.Time{}
|
||||
for i := range pods.Items {
|
||||
pod := &pods.Items[i]
|
||||
idx, ok := failed[pod.Labels["job-name"]]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
msg := lastTerminationLine(pod)
|
||||
if msg == "" || !pod.CreationTimestamp.Time.After(newest[pod.Labels["job-name"]]) {
|
||||
continue
|
||||
}
|
||||
newest[pod.Labels["job-name"]] = pod.CreationTimestamp.Time
|
||||
jobs[idx].Message = msg
|
||||
}
|
||||
}
|
||||
|
||||
// lastTerminationLine returns the last non-empty line of the pod's terminated
|
||||
// container message, capped for display.
|
||||
func lastTerminationLine(pod *corev1.Pod) string {
|
||||
for _, cs := range pod.Status.ContainerStatuses {
|
||||
t := cs.State.Terminated
|
||||
if t == nil || t.ExitCode == 0 {
|
||||
continue
|
||||
}
|
||||
lines := strings.Split(strings.TrimSpace(t.Message), "\n")
|
||||
line := strings.TrimSpace(lines[len(lines)-1])
|
||||
if len(line) > 400 {
|
||||
line = line[:400] + "…"
|
||||
}
|
||||
return line
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// jobToAsyncJob projects one Job onto its kind/state/message. Complete condition →
|
||||
// succeeded, Failed → failed with its reason (Job conditions carry the generic
|
||||
// "backoff limit exceeded" text; the pod log holds the underlying error), anything
|
||||
|
||||
@@ -747,6 +747,28 @@ func (p *PGRepo) BackupByID(ctx context.Context, id string) (*BackupRecord, erro
|
||||
return &b, nil
|
||||
}
|
||||
|
||||
// LastBackupRequest reads the newest backup.create audit row for the server
|
||||
// since the given time; the created_at index bounds the scan to that window.
|
||||
func (p *PGRepo) LastBackupRequest(ctx context.Context, serverName string, since time.Time) (time.Time, error) {
|
||||
var at sql.NullTime
|
||||
err := p.db.QueryRowContext(ctx,
|
||||
`SELECT max(created_at) FROM audit_logs
|
||||
WHERE created_at >= $2 AND action = 'backup.create' AND server_name = $1`,
|
||||
serverName, since).Scan(&at)
|
||||
if err != nil {
|
||||
return time.Time{}, err
|
||||
}
|
||||
return at.Time, nil
|
||||
}
|
||||
|
||||
// BackupStoreBytes sums size_bytes over the present world backups.
|
||||
func (p *PGRepo) BackupStoreBytes(ctx context.Context) (int64, error) {
|
||||
var n int64
|
||||
err := p.db.QueryRowContext(ctx,
|
||||
`SELECT COALESCE(sum(size_bytes), 0) FROM world_backups WHERE status = 'present'`).Scan(&n)
|
||||
return n, 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).
|
||||
|
||||
@@ -297,6 +297,14 @@ type Repo interface {
|
||||
// none matches. Like LatestBackup the returned BackupRecord carries the
|
||||
// server-side backup_ref the restore path needs; the client never sees it.
|
||||
BackupByID(ctx context.Context, id string) (*BackupRecord, error)
|
||||
// LastBackupRequest returns when an on-demand backup of the server was last
|
||||
// accepted (its newest backup.create audit row) at or after since, or the zero
|
||||
// time when there was none. The since bound keeps the lookup inside the
|
||||
// cooldown window the caller enforces.
|
||||
LastBackupRequest(ctx context.Context, serverName string, since time.Time) (time.Time, error)
|
||||
// BackupStoreBytes sums the sizes of every present world backup, the figure
|
||||
// [archive] max_local_bytes caps.
|
||||
BackupStoreBytes(ctx context.Context) (int64, error)
|
||||
// SeedServer inserts the business-layer rows for a newly created server (spec
|
||||
// §15): a servers row (owner_id NULL — claimed later, spec §9.3) and its
|
||||
// subdomain alias, both idempotent. The resource cache (cpuMilli, memoryMB,
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
package backup
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"syscall"
|
||||
)
|
||||
|
||||
// MinFreeAfter is the share of the archive filesystem an on-demand backup must
|
||||
// leave free. The archive store sits on the node's disk beside the worlds and
|
||||
// the database; below about a tenth free the kubelet starts evicting pods
|
||||
// (docs/troubleshooting.md §13b), so a backup that would cross it is refused.
|
||||
const MinFreeAfter = 0.10
|
||||
|
||||
// ErrNoRoom is returned by CheckRoom when the archive would leave too little
|
||||
// free space.
|
||||
var ErrNoRoom = errors.New("backup: not enough free disk for the archive")
|
||||
|
||||
// CheckRoom refuses an archive of srcDir into archiveDir that could push the
|
||||
// archive filesystem below minFree free. The world's uncompressed size stands
|
||||
// in for the archive's, which gzip only makes smaller. archiveDir need not
|
||||
// exist yet; its nearest existing parent is measured.
|
||||
func CheckRoom(archiveDir, srcDir string, minFree float64) error {
|
||||
dir := archiveDir
|
||||
for {
|
||||
if _, err := os.Stat(dir); err == nil {
|
||||
break
|
||||
}
|
||||
parent := filepath.Dir(dir)
|
||||
if parent == dir {
|
||||
return fmt.Errorf("backup: no existing directory above %s", archiveDir)
|
||||
}
|
||||
dir = parent
|
||||
}
|
||||
var st syscall.Statfs_t
|
||||
if err := syscall.Statfs(dir, &st); err != nil {
|
||||
return fmt.Errorf("backup: measure %s: %w", dir, err)
|
||||
}
|
||||
bsize := uint64(st.Bsize) // uint32 on darwin
|
||||
total := uint64(st.Blocks) * bsize
|
||||
avail := uint64(st.Bavail) * bsize
|
||||
if total == 0 {
|
||||
return nil
|
||||
}
|
||||
need, err := treeBytes(srcDir)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
floor := uint64(float64(total) * minFree)
|
||||
if avail < need || avail-need < floor {
|
||||
return fmt.Errorf("%w: the world is %s and %s is free of %s, which would leave less than %.0f%% free",
|
||||
ErrNoRoom, byteSize(need), byteSize(avail), byteSize(total), minFree*100)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// treeBytes sums the sizes of the regular files under dir.
|
||||
func treeBytes(dir string) (uint64, error) {
|
||||
var n uint64
|
||||
err := filepath.WalkDir(dir, func(_ string, d fs.DirEntry, err error) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if d.Type().IsRegular() {
|
||||
info, err := d.Info()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
n += uint64(info.Size())
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("backup: measure the world: %w", err)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
func byteSize(b uint64) string {
|
||||
const unit = 1024
|
||||
if b < unit {
|
||||
return fmt.Sprintf("%d B", b)
|
||||
}
|
||||
div, exp := uint64(unit), 0
|
||||
for n := b / unit; n >= unit; n /= unit {
|
||||
div *= unit
|
||||
exp++
|
||||
}
|
||||
return fmt.Sprintf("%.1f %ciB", float64(b)/float64(div), "KMGTPE"[exp])
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
package backup
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestCheckRoom(t *testing.T) {
|
||||
src := t.TempDir()
|
||||
if err := os.WriteFile(filepath.Join(src, "level.dat"), make([]byte, 4096), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
archives := filepath.Join(t.TempDir(), "not", "made", "yet")
|
||||
if err := CheckRoom(archives, src, 0); err != nil {
|
||||
t.Fatalf("no floor: %v", err)
|
||||
}
|
||||
// No disk is ever entirely free, so a 100% floor always refuses.
|
||||
if err := CheckRoom(archives, src, 1); !errors.Is(err, ErrNoRoom) {
|
||||
t.Fatalf("full floor: err = %v, want ErrNoRoom", err)
|
||||
}
|
||||
if err := CheckRoom(archives, filepath.Join(src, "absent"), 0); err == nil {
|
||||
t.Fatal("a missing world measured as empty")
|
||||
}
|
||||
}
|
||||
@@ -159,6 +159,10 @@ func BackupJob(p JobParams) (*batchv1.Job, error) {
|
||||
},
|
||||
}
|
||||
|
||||
// The exit error reaches GET /servers/{name}/jobs through the terminated
|
||||
// state (api.K8sJobStatus), in place of the Job's generic backoff text.
|
||||
container.TerminationMessagePolicy = corev1.TerminationMessageFallbackToLogsOnError
|
||||
|
||||
name := p.JobName
|
||||
if name == "" {
|
||||
name = BackupJobName(p.Server)
|
||||
|
||||
@@ -220,13 +220,20 @@ type RegistryS3Config struct {
|
||||
}
|
||||
|
||||
// ArchiveConfig is the [archive] table plus its [archive.s3] subtable (spec §19).
|
||||
// The manual_* keys bound the owners' on-demand backups, which share the
|
||||
// archive store with the reaper's: how long each is kept (default 30d), how
|
||||
// many per server (default 5, the oldest go first), and how soon an owner may
|
||||
// ask for the next one (default 10m). Empty or zero means the default.
|
||||
type ArchiveConfig struct {
|
||||
Store string `toml:"store"`
|
||||
LocalPath string `toml:"local_path"`
|
||||
Retention string `toml:"retention"`
|
||||
WarnBefore []string `toml:"warn_before"`
|
||||
MaxLocalBytes string `toml:"max_local_bytes"`
|
||||
S3 ArchiveS3Config `toml:"s3"`
|
||||
Store string `toml:"store"`
|
||||
LocalPath string `toml:"local_path"`
|
||||
Retention string `toml:"retention"`
|
||||
WarnBefore []string `toml:"warn_before"`
|
||||
MaxLocalBytes string `toml:"max_local_bytes"`
|
||||
ManualRetention string `toml:"manual_retention"`
|
||||
ManualKeep int `toml:"manual_keep"`
|
||||
ManualCooldown string `toml:"manual_cooldown"`
|
||||
S3 ArchiveS3Config `toml:"s3"`
|
||||
}
|
||||
|
||||
// ArchiveS3Config is the [archive.s3] subtable.
|
||||
|
||||
@@ -5,9 +5,11 @@ package pgint
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/api"
|
||||
"felis.lolicon.best/internal/backup"
|
||||
"felis.lolicon.best/internal/reaper"
|
||||
)
|
||||
@@ -139,3 +141,95 @@ func TestReclaimRestartsReaperClock(t *testing.T) {
|
||||
t.Fatalf("owner = %v, newest backup = %q; want released and %q", owner, newest, ar.archived[0])
|
||||
}
|
||||
}
|
||||
|
||||
// TestManualBackupRationing pins the SQL behind data-durability-9: keep-N
|
||||
// pruning picks a server's oldest on-demand backups, capacity eviction takes
|
||||
// on-demand backups before copied reaper archives and never offers the only
|
||||
// copy of a reaped world, and the API's cooldown and store-size reads see the
|
||||
// rows the rest of the platform writes.
|
||||
func TestManualBackupRationing(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := reaper.NewPGStore(db)
|
||||
name := "ration-" + suffix(t)
|
||||
if _, err := db.ExecContext(ctx,
|
||||
`INSERT INTO servers (name, cached_cpu_milli, cached_memory_mb, cached_storage_mb) VALUES ($1, 100, 128, 1)`,
|
||||
name); err != nil {
|
||||
t.Fatalf("seed server: %v", err)
|
||||
}
|
||||
storeBefore, err := repo.BackupStoreBytes(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("BackupStoreBytes: %v", err)
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
insert := func(id, reason string, created time.Time, offsite bool, size int64) {
|
||||
t.Helper()
|
||||
var offsiteAt any
|
||||
if offsite {
|
||||
offsiteAt = created
|
||||
}
|
||||
if _, err := db.ExecContext(ctx,
|
||||
`INSERT INTO world_backups (id, server_name, backup_ref, size_bytes, reason, status, created_at, expires_at, offsite_at)
|
||||
VALUES ($1, $2, $3, $4, $5, 'present', $6, $7, $8)`,
|
||||
id, name, "/archives/"+id+".tar.gz", size, reason, created, created.Add(90*reaper.Day), offsiteAt); err != nil {
|
||||
t.Fatalf("seed %s: %v", id, err)
|
||||
}
|
||||
}
|
||||
sfx := suffix(t)
|
||||
var manual []string
|
||||
for i := 0; i < 7; i++ {
|
||||
id := "bk-m" + string(rune('0'+i)) + "-" + sfx
|
||||
manual = append(manual, id)
|
||||
insert(id, "manual", now.Add(time.Duration(i-7)*time.Hour), false, 10)
|
||||
}
|
||||
sole := "bk-sole-" + sfx
|
||||
copied := "bk-copied-" + sfx
|
||||
insert(sole, "inactive_15d", now.Add(-100*reaper.Day), false, 1000)
|
||||
insert(copied, "inactive_15d", now.Add(-50*reaper.Day), true, 100)
|
||||
|
||||
excess, err := st.ExcessManualBackups(ctx, name, 5)
|
||||
if err != nil {
|
||||
t.Fatalf("ExcessManualBackups: %v", err)
|
||||
}
|
||||
if len(excess) != 2 || excess[0].ID != manual[0] || excess[1].ID != manual[1] {
|
||||
t.Fatalf("excess = %+v; want the two oldest manual backups, oldest first", excess)
|
||||
}
|
||||
|
||||
all, err := st.EvictableBackups(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("EvictableBackups: %v", err)
|
||||
}
|
||||
var got []string
|
||||
for _, b := range all {
|
||||
if b.ServerName == name {
|
||||
got = append(got, b.ID)
|
||||
}
|
||||
}
|
||||
want := append(append([]string{}, manual...), copied)
|
||||
if strings.Join(got, ",") != strings.Join(want, ",") {
|
||||
t.Fatalf("eviction order = %v\nwant %v (manual oldest first, then the copied archive, never the sole copy)", got, want)
|
||||
}
|
||||
|
||||
storeAfter, err := repo.BackupStoreBytes(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("BackupStoreBytes: %v", err)
|
||||
}
|
||||
if storeAfter-storeBefore != 7*10+1000+100 {
|
||||
t.Fatalf("store grew by %d, want %d", storeAfter-storeBefore, 7*10+1000+100)
|
||||
}
|
||||
|
||||
if at, err := repo.LastBackupRequest(ctx, name, now.Add(-10*time.Minute)); err != nil || !at.IsZero() {
|
||||
t.Fatalf("before any request: LastBackupRequest = (%v, %v)", at, err)
|
||||
}
|
||||
if err := repo.Audit(ctx, api.AuditEntry{Actor: "[email protected]", Source: "external",
|
||||
Action: "backup.create", ServerName: name}); err != nil {
|
||||
t.Fatalf("Audit: %v", err)
|
||||
}
|
||||
at, err := repo.LastBackupRequest(ctx, name, time.Now().Add(-10*time.Minute))
|
||||
if err != nil || at.IsZero() || time.Since(at) > time.Minute {
|
||||
t.Fatalf("after a request: LastBackupRequest = (%v, %v)", at, err)
|
||||
}
|
||||
if at, err := repo.LastBackupRequest(ctx, name, time.Now().Add(time.Minute)); err != nil || !at.IsZero() {
|
||||
t.Fatalf("a request before since still counted: (%v, %v)", at, err)
|
||||
}
|
||||
}
|
||||
@@ -106,14 +106,28 @@ func (s *PGStore) PresentBackupBytes(ctx context.Context) (int64, error) {
|
||||
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`
|
||||
func (s *PGStore) EvictableBackups(ctx context.Context) ([]StoredBackup, error) {
|
||||
const q = `SELECT id, server_name, backup_ref, size_bytes, reason FROM world_backups
|
||||
WHERE status = 'present' AND (reason <> 'inactive_15d' OR offsite_at IS NOT NULL)
|
||||
ORDER BY reason = 'inactive_15d', created_at ASC`
|
||||
return s.queryBackups(ctx, q)
|
||||
}
|
||||
|
||||
// ExcessManualBackups lists server's present on-demand backups beyond the
|
||||
// newest keep, oldest first: what the backup Job removes after adding one.
|
||||
func (s *PGStore) ExcessManualBackups(ctx context.Context, server string, keep int) ([]StoredBackup, error) {
|
||||
const q = `SELECT id, server_name, backup_ref, size_bytes, reason FROM world_backups
|
||||
WHERE server_name = $1 AND status = 'present' AND reason = 'manual'
|
||||
ORDER BY created_at DESC OFFSET $2`
|
||||
out, err := s.queryBackups(ctx, q, server, keep)
|
||||
for i, j := 0, len(out)-1; i < j; i, j = i+1, j-1 {
|
||||
out[i], out[j] = out[j], out[i]
|
||||
}
|
||||
return out, err
|
||||
}
|
||||
|
||||
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
|
||||
const q = `SELECT id, server_name, backup_ref, size_bytes, reason FROM world_backups
|
||||
WHERE status = 'present' AND expires_at < $1 ORDER BY expires_at ASC`
|
||||
return s.queryBackups(ctx, q, now)
|
||||
}
|
||||
@@ -127,7 +141,7 @@ func (s *PGStore) queryBackups(ctx context.Context, q string, args ...any) ([]St
|
||||
var out []StoredBackup
|
||||
for rows.Next() {
|
||||
var b StoredBackup
|
||||
if err := rows.Scan(&b.ID, &b.ServerName, &b.BackupRef, &b.SizeBytes); err != nil {
|
||||
if err := rows.Scan(&b.ID, &b.ServerName, &b.BackupRef, &b.SizeBytes, &b.Reason); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, b)
|
||||
|
||||
@@ -81,6 +81,16 @@ type Config struct {
|
||||
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
|
||||
// ManualRetention is how long an owner's on-demand backup is kept; it is
|
||||
// a restore point for a world that still exists, so it goes sooner than a
|
||||
// reaped world's only archive (default 30d).
|
||||
ManualRetention time.Duration
|
||||
// ManualKeep caps the on-demand backups kept per server; the backup Job
|
||||
// removes the oldest beyond it (default 5).
|
||||
ManualKeep int
|
||||
// ManualCooldown is the shortest gap between two owner-requested backups
|
||||
// of one server (default 10m); operators are not held to it.
|
||||
ManualCooldown time.Duration
|
||||
// RequireOffsite holds each deletion until the world's archive has its
|
||||
// off-site copy ([offsite] configured; internal/offsite records the copy).
|
||||
// The archive is written on the run that finds the world idle, and the
|
||||
@@ -96,6 +106,10 @@ func DefaultConfig() Config {
|
||||
WarnBefore: []time.Duration{3 * Day, 1 * Day},
|
||||
Retention: 90 * Day,
|
||||
MaxLocalBytes: 0,
|
||||
|
||||
ManualRetention: 30 * Day,
|
||||
ManualKeep: 5,
|
||||
ManualCooldown: 10 * time.Minute,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -151,6 +165,7 @@ type StoredBackup struct {
|
||||
ServerName string
|
||||
BackupRef string
|
||||
SizeBytes int64
|
||||
Reason string
|
||||
}
|
||||
|
||||
// AuditRecord is a reaper-sourced audit_logs entry. The PG binding fills
|
||||
@@ -190,9 +205,12 @@ type Store interface {
|
||||
// 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)
|
||||
// EvictableBackups lists the status=present backups that may go before
|
||||
// their expiry when the store is full, in eviction order: on-demand
|
||||
// backups first, then reaper archives that have an off-site copy, oldest
|
||||
// first within each. A reaper archive without an off-site copy is the only
|
||||
// copy of a deleted world and is never listed.
|
||||
EvictableBackups(ctx context.Context) ([]StoredBackup, error)
|
||||
|
||||
// ListExpiredBackups lists status=present backups whose expires_at < now.
|
||||
ListExpiredBackups(ctx context.Context, now time.Time) ([]StoredBackup, error)
|
||||
@@ -450,10 +468,10 @@ func (r *Reaper) ensureCapacity(ctx context.Context, now time.Time, sum *Summary
|
||||
if used < r.Cfg.MaxLocalBytes {
|
||||
return true, nil
|
||||
}
|
||||
r.log().Warn("reaper: backup store at capacity, evicting oldest backups early",
|
||||
r.log().Warn("reaper: backup store at capacity, evicting on-demand and off-site-copied backups early",
|
||||
"used", used, "max", r.Cfg.MaxLocalBytes)
|
||||
|
||||
old, err := r.Store.OldestPresentBackups(ctx)
|
||||
old, err := r.Store.EvictableBackups(ctx)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
@@ -474,6 +492,11 @@ func (r *Reaper) ensureCapacity(ctx context.Context, now time.Time, sum *Summary
|
||||
}
|
||||
used -= b.SizeBytes
|
||||
sum.EvictedEarly++
|
||||
r.log().Warn("reaper: backup evicted early", "id", b.ID, "server", b.ServerName, "reason", b.Reason, "bytes", b.SizeBytes)
|
||||
}
|
||||
if used >= r.Cfg.MaxLocalBytes {
|
||||
r.log().Error("reaper: backup store still full; what remains are the only copies of reaped worlds, kept until they expire",
|
||||
"used", used, "max", r.Cfg.MaxLocalBytes)
|
||||
}
|
||||
return used < r.Cfg.MaxLocalBytes, nil
|
||||
}
|
||||
|
||||
@@ -193,17 +193,22 @@ func (s *fakeStore) PresentBackupBytes(context.Context) (int64, error) {
|
||||
return total, nil
|
||||
}
|
||||
|
||||
func (s *fakeStore) OldestPresentBackups(context.Context) ([]StoredBackup, error) {
|
||||
func (s *fakeStore) EvictableBackups(context.Context) ([]StoredBackup, error) {
|
||||
var ps []*fakeBackup
|
||||
for _, b := range s.backups {
|
||||
if b.status == "present" {
|
||||
if b.status == "present" && (b.reason != ReasonInactive || b.offsite) {
|
||||
ps = append(ps, b)
|
||||
}
|
||||
}
|
||||
sort.Slice(ps, func(i, j int) bool { return ps[i].createdAt.Before(ps[j].createdAt) })
|
||||
sort.SliceStable(ps, func(i, j int) bool {
|
||||
if ri, rj := ps[i].reason == ReasonInactive, ps[j].reason == ReasonInactive; ri != rj {
|
||||
return rj
|
||||
}
|
||||
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})
|
||||
out = append(out, StoredBackup{ID: b.id, ServerName: b.server, BackupRef: b.ref, SizeBytes: b.size, Reason: b.reason})
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@@ -623,8 +628,8 @@ func TestCapacityEvictsOldestThenReaps(t *testing.T) {
|
||||
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)},
|
||||
{id: "old", server: "zzz", ref: "ref-old", reason: "manual", size: 75, status: "present", createdAt: idleBy(40 * Day), expires: testNow.Add(30 * Day)},
|
||||
{id: "new", server: "yyy", ref: "ref-new", reason: "manual", size: 75, status: "present", createdAt: idleBy(5 * Day), expires: testNow.Add(60 * Day)},
|
||||
}
|
||||
|
||||
sum := mustRun(t, r)
|
||||
@@ -669,7 +674,7 @@ func TestCapacityStillFullSkipsReap(t *testing.T) {
|
||||
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)},
|
||||
{id: "stuck", server: "zzz", ref: "ref-stuck", reason: "manual", 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")
|
||||
@@ -689,6 +694,50 @@ func TestCapacityStillFullSkipsReap(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Eviction order: on-demand backups go first, then reaper archives that have an
|
||||
// off-site copy; the only copy of a reaped world is never evicted early, even
|
||||
// when that leaves the store full and the idle world waits.
|
||||
func TestCapacityEvictionOrderSparesSoleCopies(t *testing.T) {
|
||||
cfg := DefaultConfig()
|
||||
cfg.MaxLocalBytes = 100
|
||||
r, st, _, _ := newReaper(cfg,
|
||||
Candidate{Name: "theta", OwnerID: "user-8", LastActiveAt: idleBy(20 * Day)})
|
||||
st.backups = []*fakeBackup{
|
||||
{id: "sole", server: "gone1", ref: "ref-sole", reason: ReasonInactive, size: 60, status: "present", createdAt: idleBy(80 * Day), expires: testNow.Add(10 * Day)},
|
||||
{id: "copied", server: "gone2", ref: "ref-copied", reason: ReasonInactive, offsite: true, size: 30, status: "present", createdAt: idleBy(70 * Day), expires: testNow.Add(20 * Day)},
|
||||
{id: "man", server: "live", ref: "ref-man", reason: "manual", size: 30, status: "present", createdAt: idleBy(2 * Day), expires: testNow.Add(28 * Day)},
|
||||
}
|
||||
sum := mustRun(t, r)
|
||||
status := map[string]string{}
|
||||
for _, b := range st.backups {
|
||||
status[b.id] = b.status
|
||||
}
|
||||
// 120 over a cap of 100: the on-demand backup (30) alone brings it to 90.
|
||||
if status["man"] != "deleted" || status["copied"] != "present" || status["sole"] != "present" {
|
||||
t.Fatalf("evicted %v; want only the on-demand backup", status)
|
||||
}
|
||||
if sum.EvictedEarly != 1 || sum.WorldsReaped != 1 {
|
||||
t.Fatalf("summary = %+v, want 1 evicted, 1 reaped", sum)
|
||||
}
|
||||
|
||||
// Over a cap of 50 with only reaper archives left: the copied one goes,
|
||||
// the sole copy stays, the store is still full, and the idle world waits.
|
||||
cfg.MaxLocalBytes = 50
|
||||
r2, st2, cl2, _ := newReaper(cfg,
|
||||
Candidate{Name: "iota", OwnerID: "user-9", LastActiveAt: idleBy(20 * Day)})
|
||||
st2.backups = []*fakeBackup{
|
||||
{id: "sole", server: "gone1", ref: "ref-sole", reason: ReasonInactive, size: 60, status: "present", createdAt: idleBy(80 * Day), expires: testNow.Add(10 * Day)},
|
||||
{id: "copied", server: "gone2", ref: "ref-copied", reason: ReasonInactive, offsite: true, size: 30, status: "present", createdAt: idleBy(70 * Day), expires: testNow.Add(20 * Day)},
|
||||
}
|
||||
sum2 := mustRun(t, r2)
|
||||
if st2.backups[0].status != "present" || st2.backups[1].status != "deleted" {
|
||||
t.Fatalf("sole=%s copied=%s, want present/deleted", st2.backups[0].status, st2.backups[1].status)
|
||||
}
|
||||
if sum2.StoreFull != 1 || sum2.WorldsReaped != 0 || cl2.deletePVCCalls != 0 {
|
||||
t.Fatalf("summary = %+v, deletes = %d: the world should wait while the store is full", sum2, cl2.deletePVCCalls)
|
||||
}
|
||||
}
|
||||
|
||||
// Retention pass: backups past expires_at are deleted from the backend and
|
||||
// marked deleted; unexpired backups are untouched.
|
||||
func TestExpiredBackupsDeleted(t *testing.T) {
|
||||
|
||||
@@ -138,6 +138,10 @@ func RestoreJob(p JobParams) (*batchv1.Job, error) {
|
||||
},
|
||||
}
|
||||
|
||||
// The exit error reaches GET /servers/{name}/jobs through the terminated
|
||||
// state (api.K8sJobStatus), in place of the Job's generic backoff text.
|
||||
container.TerminationMessagePolicy = corev1.TerminationMessageFallbackToLogsOnError
|
||||
|
||||
job := &batchv1.Job{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: RestoreJobName(p.Server),
|
||||
|
||||
Reference in new issue
Block a user