feat(restore): 恢复前默认为当前世界做安全快照并串接恢复 Job,快照失败则不恢复;并发恢复另一备份返回 409

This commit is contained in:
Lemon-miaow committed 2026-09-25 02:52:18 +08:00
1 parent d44243c0ac
commit 489eff4494
33 files changed
+1260 -105

No files matched your search

+6
View File
@@ -75,6 +75,12 @@ type API struct {
// endpoints only answer 202). Optional: nil → that route reports 503.
JobStatus JobStatusReader
// RestoreChains finds and settles the safety snapshots that run in front of a
// restore (restorechain.go). A restore takes a snapshot first only when this
// is set, the Backuper can chain a restore and the Restorer is wired, since
// something has to start the restore once the snapshot is done.
RestoreChains RestoreChains
// Files is the server file editor (list / read / write a file in a stopped
// server's world volume — the "one wrong line in server.properties" repair).
// Like Restorer and Backuper it is optional: when nil the file routes report
+10
View File
@@ -23,3 +23,13 @@ import "context"
type Backuper interface {
Backup(ctx context.Context, serverName, formerOwner string) error
}
// RestoreSnapshotter is the Backuper's safety-snapshot lever: it enqueues a backup
// of the world as it is now, labelled with the restore to run once that backup
// has succeeded (backupID, backupRef). RestoreChains settles it: it starts the
// restore after a successful snapshot and gives the restore up after a failed
// one, so a restore never overwrites a world that has no copy. Until it is
// settled the snapshot holds the world volume as a restore (internal/maintenance).
type RestoreSnapshotter interface {
BackupThenRestore(ctx context.Context, serverName, formerOwner, backupID, backupRef string) error
}
+32 -7
View File
@@ -76,8 +76,14 @@ func (a *API) handleListBackups(w http.ResponseWriter, r *http.Request) {
// makes that atomic against a wake and refuses a second restore, backup or
// file write on the same world with 409 maintenance_in_progress until the
// restore Job finishes.
// ⑧ hand off to the Restorer. Restore is asynchronous (a restore Job, like an
// image build Job), so success means "enqueued" and the handler answers 202.
// ⑧ hand off. By default the restore starts with a safety snapshot of the world
// as it is (restorechain.go): a pre_restore backup Job that the restore Job
// follows once it succeeds, so a wrong pick can be walked back from the
// backup list. "safety_snapshot": false in the body restores straight away.
// Either way the work is asynchronous (Jobs, like an image build), so
// success means "enqueued" and the handler answers 202, saying in
// safety_snapshot which of the two it did. A restore of another backup still
// running on the world is 409 restore_in_progress.
//
// The opaque backup_ref is resolved server-side from the backup and handed to the
// Restorer directly; the client never names a backup by handle (spec §286
@@ -104,8 +110,10 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) {
}
// Optional backup_id in the JSON body; absent → LatestBackup (backward compat).
// safety_snapshot defaults to true.
var body struct {
BackupID string `json:"backup_id"`
BackupID string `json:"backup_id"`
SafetySnapshot *bool `json:"safety_snapshot"`
}
if strings.HasPrefix(r.Header.Get("Content-Type"), "application/json") {
if err := decodeJSON(w, r, &body); err != nil {
@@ -204,7 +212,23 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) {
}
defer release()
if err := a.Restorer.Restore(r.Context(), name, backup.BackupRef); err != nil {
snapshotter, snapshot := a.snapshotFirst()
if body.SafetySnapshot != nil && !*body.SafetySnapshot {
snapshot = false
}
if snapshot {
// The snapshot records the current owner, like an on-demand backup, so
// the way back is theirs to take.
err = snapshotter.BackupThenRestore(r.Context(), name, rec.OwnerID, backup.ID, backup.BackupRef)
} else {
err = a.Restorer.Restore(r.Context(), name, backup.BackupRef)
}
if err != nil {
if isRestoreInProgress(err) {
writeError(w, r, newError(http.StatusConflict, "restore_in_progress",
"a restore of another backup is still running on this server's world; retry once it finishes"))
return
}
// ErrNotFound (server vanished from the execution backend) → 404; else 500.
a.writeLookupError(w, r, err)
return
@@ -212,9 +236,10 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) {
a.audit(r, "backup.restore", name)
writeJSON(w, http.StatusAccepted, map[string]any{
"name": name,
"status": "restoring",
"backup_id": backup.ID,
"name": name,
"status": "restoring",
"backup_id": backup.ID,
"safety_snapshot": snapshot,
})
}
+7
View File
@@ -19,6 +19,13 @@ type AsyncJob struct {
Message string `json:"message,omitempty"`
StartedAt time.Time `json:"started_at,omitzero"`
FinishedAt time.Time `json:"finished_at,omitzero"`
// ThenRestore is set on a restore's safety snapshot (a backup Job): "pending"
// until the restore behind it starts ("started") or is given up
// ("abandoned", with a ChainAbandon* code in ThenRestoreReason and its
// English in Message). RestoreBackupID is the backup that restore extracts.
ThenRestore string `json:"then_restore,omitempty"`
ThenRestoreReason string `json:"then_restore_reason,omitempty"`
RestoreBackupID string `json:"restore_backup_id,omitempty"`
}
// JobStatusReader reads the newest backup/restore Jobs for a server, newest
+78 -2
View File
@@ -2,12 +2,16 @@ package api
import (
"context"
"encoding/json"
"sort"
"strings"
"time"
"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/types"
"sigs.k8s.io/controller-runtime/pkg/client"
)
@@ -115,11 +119,83 @@ func lastTerminationLine(pod *corev1.Pod) string {
return ""
}
// jobToAsyncJob projects one Job onto its kind/state/message. Complete condition →
// jobToAsyncJob projects one Job onto the AsyncJob the jobs route answers with.
func jobToAsyncJob(j *batchv1.Job) (AsyncJob, bool) {
aj, ok := jobOutcome(j)
if !ok {
return aj, false
}
// A safety snapshot says what became of the restore behind it; a chain given
// up after a snapshot that succeeded explains itself in the message.
if state := j.Labels[maintenance.LabelThenRestore]; state != "" && aj.Kind == "backup" {
aj.ThenRestore = state
aj.RestoreBackupID = j.Annotations[maintenance.AnnotationRestoreBackupID]
if reason := j.Annotations[maintenance.AnnotationThenRestoreReason]; reason != "" {
aj.ThenRestoreReason = reason
if aj.Message == "" {
aj.Message = chainAbandonMessage(reason)
}
}
}
return aj, true
}
// 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) {
var list batchv1.JobList
if err := k.c.List(ctx, &list, client.InNamespace(k.namespace), client.MatchingLabels{
jobManagedByLabel: jobManagedByBackup,
maintenance.LabelThenRestore: maintenance.ThenRestorePending,
}); err != nil {
return nil, err
}
out := make([]RestoreChain, 0, len(list.Items))
for i := range list.Items {
j := &list.Items[i]
snapshot := ChainSnapshotRunning
for _, c := range j.Status.Conditions {
if c.Status != corev1.ConditionTrue {
continue
}
switch c.Type {
case batchv1.JobComplete, batchv1.JobSuccessCriteriaMet:
snapshot = ChainSnapshotSucceeded
case batchv1.JobFailed, batchv1.JobFailureTarget:
snapshot = ChainSnapshotFailed
}
}
out = append(out, RestoreChain{
Job: j.Name,
Server: j.Labels[jobServerLabel],
BackupID: j.Annotations[maintenance.AnnotationRestoreBackupID],
BackupRef: j.Annotations[maintenance.AnnotationRestoreRef],
Snapshot: snapshot,
})
}
return out, nil
}
// SettleRestoreChain records how a chain was settled on its snapshot Job, which
// releases the world volume the chain was holding.
func (k *K8sJobStatus) SettleRestoreChain(ctx context.Context, job, state, reason string) error {
meta := map[string]any{"labels": map[string]string{maintenance.LabelThenRestore: state}}
if reason != "" {
meta["annotations"] = map[string]string{maintenance.AnnotationThenRestoreReason: reason}
}
patch, err := json.Marshal(map[string]any{"metadata": meta})
if err != nil {
return err
}
obj := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Namespace: k.namespace, Name: job}}
return k.c.Patch(ctx, obj, client.RawPatch(types.MergePatchType, patch))
}
// jobOutcome 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
// else is still running.
func jobToAsyncJob(j *batchv1.Job) (AsyncJob, bool) {
func jobOutcome(j *batchv1.Job) (AsyncJob, bool) {
kind := ""
switch j.Labels[jobManagedByLabel] {
case jobManagedByBackup:
+158
View File
@@ -0,0 +1,158 @@
package api
import (
"context"
"errors"
"fmt"
"log"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/maintenance"
)
// A restore overwrites the world, so by default it starts with a safety snapshot:
// POST /servers/{name}/restore-backup enqueues a backup Job carrying the restore
// (RestoreSnapshotter), and SettleRestoreChains starts that restore once the
// snapshot has succeeded. The snapshot is an ordinary pre_restore row in the
// backup list, which is the way back from a wrong pick.
//
// The chain lives on the backup Job (labels and annotations), so a felis-api
// restart picks it up where it was. While it is pending the Job holds the world
// volume as a restore, which keeps the server down in the gap between the two
// Jobs; if felis-api never settles it, the lock goes with the Job at its TTL and
// the world is left as it was.
// Snapshot Job states a RestoreChain reports.
const (
ChainSnapshotRunning = "running"
ChainSnapshotSucceeded = "succeeded"
ChainSnapshotFailed = "failed"
)
// RestoreChain is one safety snapshot whose restore is still pending.
type RestoreChain struct {
Job string // the snapshot's backup Job
Server string
BackupID string // the backup the restore extracts
BackupRef string
Snapshot string // ChainSnapshotRunning / Succeeded / Failed
}
// Why a chain was abandoned: the code recorded on the snapshot Job and reported
// by the jobs route as then_restore_reason, for the panel to word in its own
// language. chainAbandonText is the English the log and the job's message carry.
const (
ChainAbandonSnapshotFailed = "snapshot_failed"
ChainAbandonNotConfigured = "not_configured"
ChainAbandonServerGone = "server_gone"
ChainAbandonServerStarted = "server_started"
ChainAbandonRestoreBusy = "restore_busy"
)
var chainAbandonText = map[string]string{
ChainAbandonSnapshotFailed: "the safety backup failed, so the world was left as it was",
ChainAbandonNotConfigured: "restore is not configured",
ChainAbandonServerGone: "the server no longer exists",
ChainAbandonServerStarted: "the server was started before the restore could run",
ChainAbandonRestoreBusy: "another restore is running on this world",
}
// chainAbandonMessage words a reason code; an unknown code (a newer felis-api
// wrote it) is shown as is.
func chainAbandonMessage(reason string) string {
if s, ok := chainAbandonText[reason]; ok {
return s
}
return reason
}
// RestoreChains reads the pending chains and records how each was settled
// (maintenance.ThenRestoreStarted or ThenRestoreAbandoned, with the reason code
// of a chain given up).
type RestoreChains interface {
PendingRestoreChains(ctx context.Context) ([]RestoreChain, error)
SettleRestoreChain(ctx context.Context, job, state, reason string) error
}
// snapshotFirst reports whether a restore can start with a safety snapshot:
// the Backuper can carry a restore and something settles the chain.
func (a *API) snapshotFirst() (RestoreSnapshotter, bool) {
if a.RestoreChains == nil || a.Restorer == nil {
return nil, false
}
s, ok := a.Backuper.(RestoreSnapshotter)
return s, ok
}
// SettleRestoreChains advances every pending chain once: it starts the restore
// behind a snapshot that succeeded and gives up the one behind a snapshot that
// failed. A chain whose restore could not be created this time stays pending and
// is retried on the next call. cmd/felis runs it on a short interval.
func (a *API) SettleRestoreChains(ctx context.Context) error {
if a.RestoreChains == nil {
return nil
}
chains, err := a.RestoreChains.PendingRestoreChains(ctx)
if err != nil {
return err
}
var errs []error
for _, c := range chains {
state, reason, err := a.settleRestoreChain(ctx, c)
if err != nil {
errs = append(errs, fmt.Errorf("%s: %w", c.Job, err))
continue
}
if state == "" {
continue
}
if err := a.RestoreChains.SettleRestoreChain(ctx, c.Job, state, reason); err != nil {
errs = append(errs, fmt.Errorf("%s: settle as %s: %w", c.Job, state, err))
continue
}
if state == maintenance.ThenRestoreStarted {
log.Printf("api: safety snapshot %s done; restoring %s from backup %s", c.Job, c.Server, c.BackupID)
} else {
log.Printf("api: restore of %s from backup %s abandoned: %s", c.Server, c.BackupID, chainAbandonMessage(reason))
}
}
return errors.Join(errs...)
}
// settleRestoreChain decides one chain: the state to record ("" to leave it
// pending) and, for an abandoned chain, the reason code.
func (a *API) settleRestoreChain(ctx context.Context, c RestoreChain) (string, string, error) {
switch c.Snapshot {
case ChainSnapshotFailed:
return maintenance.ThenRestoreAbandoned, ChainAbandonSnapshotFailed, nil
case ChainSnapshotSucceeded:
default:
return "", "", nil
}
if a.Restorer == nil {
return maintenance.ThenRestoreAbandoned, ChainAbandonNotConfigured, nil
}
// The pending chain holds the world volume, so nothing should have started
// the server; a start that got past it anyway must not have its world
// replaced underneath.
info, err := a.Cluster.GetServer(ctx, c.Server)
if errors.Is(err, ErrNotFound) {
return maintenance.ThenRestoreAbandoned, ChainAbandonServerGone, nil
}
if err != nil {
return "", "", err
}
if info.Ready || info.DesiredState != string(v1alpha1.DesiredStopped) {
return maintenance.ThenRestoreAbandoned, ChainAbandonServerStarted, nil
}
if err := a.Restorer.Restore(ctx, c.Server, c.BackupRef); err != nil {
if isRestoreInProgress(err) {
return maintenance.ThenRestoreAbandoned, ChainAbandonRestoreBusy, nil
}
if errors.Is(err, ErrNotFound) {
return maintenance.ThenRestoreAbandoned, ChainAbandonServerGone, nil
}
return "", "", err
}
return maintenance.ThenRestoreStarted, "", nil
}
+310
View File
@@ -0,0 +1,310 @@
package api
import (
"context"
"encoding/json"
"errors"
"net/http"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/backupjob"
"felis.lolicon.best/internal/maintenance"
"felis.lolicon.best/internal/restore"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
// fakeSnapshotter is a Backuper that can also chain a restore.
type fakeSnapshotter struct {
fakeBackuper
err error
chained []RestoreChain
}
func (f *fakeSnapshotter) BackupThenRestore(_ context.Context, name, formerOwner, backupID, backupRef string) error {
f.gotFormerOwn = formerOwner
f.chained = append(f.chained, RestoreChain{Server: name, BackupID: backupID, BackupRef: backupRef})
return f.err
}
type settled struct{ job, state, note string }
type fakeChains struct {
pending []RestoreChain
listErr error
settleErr error
settled []settled
}
func (f *fakeChains) PendingRestoreChains(context.Context) ([]RestoreChain, error) {
return f.pending, f.listErr
}
func (f *fakeChains) SettleRestoreChain(_ context.Context, job, state, note string) error {
if f.settleErr != nil {
return f.settleErr
}
f.settled = append(f.settled, settled{job, state, note})
return nil
}
// otherRestore mirrors the executor's conflict error, which api only knows by
// its method.
type otherRestore struct{}
func (otherRestore) Error() string { return "another restore" }
func (otherRestore) RestoreInProgress() bool { return true }
var _ RestoreSnapshotter = (*backupjob.Backuper)(nil)
func TestIsRestoreInProgressMatchesTheExecutor(t *testing.T) {
if !isRestoreInProgress(restore.ErrOtherRestoreRunning) {
t.Fatal("restore.ErrOtherRestoreRunning is not recognised")
}
if isRestoreInProgress(restore.ErrAlreadyExists) || isRestoreInProgress(errors.New("x")) {
t.Fatal("an unrelated error reads as a restore in progress")
}
}
// With the chain wired, a restore starts with a safety snapshot of the world as
// it is: the Backuper gets the restore to carry, the Restorer is not called
// yet, and the 202 says so. safety_snapshot:false restores straight away.
func TestRestoreStartsWithSafetySnapshot(t *testing.T) {
owner := &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}
mk := func() (*API, *fakeSnapshotter, *fakeRestorer) {
repo := newFakeRepo()
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
repo.backups = []fakeBackup{{view: BackupView{ID: "bk1", ServerName: "survival", FormerOwner: "owner1",
Status: "present", Reason: "manual", CreatedAt: time.Unix(1_699_000_000, 0)}, ref: "ref-1"}}
cl := newFakeCluster()
cl.byName["survival"] = &ServerInfo{Name: "survival", Phase: "Stopped", DesiredState: string(v1alpha1.DesiredStopped)}
a := newTestAPI(repo, cl)
snap, restorer := &fakeSnapshotter{}, &fakeRestorer{}
a.Backuper, a.Restorer, a.RestoreChains = snap, restorer, &fakeChains{}
a.External = staticExternal{p: owner}
return a, snap, restorer
}
const path = "/api/v1/servers/survival/restore-backup"
decode := func(t *testing.T, body []byte) map[string]any {
t.Helper()
var m map[string]any
if err := json.Unmarshal(body, &m); err != nil {
t.Fatalf("body: %v", err)
}
return m
}
t.Run("default -> snapshot first", func(t *testing.T) {
a, snap, restorer := mk()
w := do(a.ExternalHandler(), "POST", path, `{"backup_id":"bk1"}`, jsonHeader)
if w.Code != http.StatusAccepted {
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
}
if m := decode(t, w.Body.Bytes()); m["safety_snapshot"] != true || m["status"] != "restoring" || m["backup_id"] != "bk1" {
t.Fatalf("response = %v", m)
}
if len(snap.chained) != 1 || snap.chained[0] != (RestoreChain{Server: "survival", BackupID: "bk1", BackupRef: "ref-1"}) || snap.gotFormerOwn != "owner1" {
t.Fatalf("chained = %+v (owner %q)", snap.chained, snap.gotFormerOwn)
}
if restorer.calls != 0 {
t.Fatal("the restore ran before its snapshot")
}
})
t.Run("safety_snapshot false -> restore now", func(t *testing.T) {
a, snap, restorer := mk()
w := do(a.ExternalHandler(), "POST", path, `{"safety_snapshot":false}`, jsonHeader)
if w.Code != http.StatusAccepted {
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
}
if m := decode(t, w.Body.Bytes()); m["safety_snapshot"] != false {
t.Fatalf("response = %v", m)
}
if restorer.calls != 1 || len(snap.chained) != 0 {
t.Fatalf("restorer calls %d, chained %d", restorer.calls, len(snap.chained))
}
})
t.Run("chain not wired -> restore now", func(t *testing.T) {
a, snap, restorer := mk()
a.RestoreChains = nil
if w := do(a.ExternalHandler(), "POST", path, "", nil); w.Code != http.StatusAccepted {
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
}
if restorer.calls != 1 || len(snap.chained) != 0 {
t.Fatalf("restorer calls %d, chained %d", restorer.calls, len(snap.chained))
}
})
t.Run("another restore running -> 409 restore_in_progress", func(t *testing.T) {
a, _, restorer := mk()
restorer.err = otherRestore{}
w := do(a.ExternalHandler(), "POST", path, `{"safety_snapshot":false}`, jsonHeader)
if w.Code != http.StatusConflict || errCode(w.Body.Bytes()) != "restore_in_progress" {
t.Fatalf("code = %d (%s)", w.Code, w.Body.String())
}
})
}
func TestSettleRestoreChains(t *testing.T) {
mk := func(chains ...RestoreChain) (*API, *fakeChains, *fakeRestorer, *fakeCluster) {
cl := newFakeCluster()
cl.byName["survival"] = &ServerInfo{Name: "survival", Phase: "Stopped", DesiredState: string(v1alpha1.DesiredStopped)}
a := newTestAPI(newFakeRepo(), cl)
fc, restorer := &fakeChains{pending: chains}, &fakeRestorer{}
a.RestoreChains, a.Restorer = fc, restorer
return a, fc, restorer, cl
}
chain := func(snapshot string) RestoreChain {
return RestoreChain{Job: "backup-survival-1", Server: "survival", BackupID: "bk1", BackupRef: "ref-1", Snapshot: snapshot}
}
ctx := context.Background()
t.Run("snapshot succeeded -> restore started", func(t *testing.T) {
a, fc, restorer, _ := mk(chain(ChainSnapshotSucceeded))
if err := a.SettleRestoreChains(ctx); err != nil {
t.Fatal(err)
}
if restorer.calls != 1 || restorer.gotName != "survival" || restorer.gotRef != "ref-1" {
t.Fatalf("restorer = %+v", restorer)
}
if len(fc.settled) != 1 || fc.settled[0] != (settled{"backup-survival-1", maintenance.ThenRestoreStarted, ""}) {
t.Fatalf("settled = %+v", fc.settled)
}
})
t.Run("snapshot running -> left alone", func(t *testing.T) {
a, fc, restorer, _ := mk(chain(ChainSnapshotRunning))
if err := a.SettleRestoreChains(ctx); err != nil {
t.Fatal(err)
}
if restorer.calls != 0 || len(fc.settled) != 0 {
t.Fatalf("restorer %d, settled %+v", restorer.calls, fc.settled)
}
})
for _, tc := range []struct {
name string
setup func(*fakeRestorer, *fakeCluster)
snap string
reason string
}{
{"snapshot failed", func(*fakeRestorer, *fakeCluster) {}, ChainSnapshotFailed, ChainAbandonSnapshotFailed},
{"server started meanwhile", func(_ *fakeRestorer, cl *fakeCluster) {
cl.byName["survival"].DesiredState = string(v1alpha1.DesiredRunning)
}, ChainSnapshotSucceeded, ChainAbandonServerStarted},
{"server gone", func(_ *fakeRestorer, cl *fakeCluster) { delete(cl.byName, "survival") }, ChainSnapshotSucceeded, ChainAbandonServerGone},
{"another restore running", func(r *fakeRestorer, _ *fakeCluster) { r.err = otherRestore{} }, ChainSnapshotSucceeded, ChainAbandonRestoreBusy},
} {
t.Run(tc.name+" -> abandoned", func(t *testing.T) {
a, fc, restorer, cl := mk(chain(tc.snap))
tc.setup(restorer, cl)
if err := a.SettleRestoreChains(ctx); err != nil {
t.Fatal(err)
}
if len(fc.settled) != 1 || fc.settled[0].state != maintenance.ThenRestoreAbandoned || fc.settled[0].note != tc.reason {
t.Fatalf("settled = %+v", fc.settled)
}
if tc.snap == ChainSnapshotFailed && restorer.calls != 0 {
t.Fatal("restored after a failed snapshot")
}
})
}
t.Run("restore error -> stays pending, retried", func(t *testing.T) {
a, fc, restorer, _ := mk(chain(ChainSnapshotSucceeded))
restorer.err = errors.New("apiserver hiccup")
if err := a.SettleRestoreChains(ctx); err == nil {
t.Fatal("the failure was swallowed")
}
if len(fc.settled) != 0 {
t.Fatalf("settled = %+v", fc.settled)
}
restorer.err = nil
if err := a.SettleRestoreChains(ctx); err != nil || len(fc.settled) != 1 || fc.settled[0].state != maintenance.ThenRestoreStarted {
t.Fatalf("retry: err %v, settled %+v", err, fc.settled)
}
})
}
// The cluster side: the chain the backupjob executor renders is found by the
// label selector, its snapshot outcome read from the Job conditions, and
// settling it releases the world volume.
func TestK8sRestoreChains(t *testing.T) {
scheme := runtime.NewScheme()
if err := clientgoscheme.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
chain, err := backupjob.BackupJob(backupjob.JobParams{
Server: "survival", JobName: "backup-survival-aa", WorldPVC: "world-survival-0", BackupPVC: "felis-backups",
Namespace: "minecraft", Image: "felis:1", ConfigSecret: "felis-config", ConfigMount: "/etc/felis",
RestoreRef: "/backups/survival/a.tar.gz", RestoreBackupID: "bk-1",
})
if err != nil {
t.Fatal(err)
}
chain.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}}
plain, err := backupjob.BackupJob(backupjob.JobParams{
Server: "survival", JobName: "backup-survival-bb", WorldPVC: "world-survival-0", BackupPVC: "felis-backups",
Namespace: "minecraft", Image: "felis:1", ConfigSecret: "felis-config", ConfigMount: "/etc/felis",
})
if err != nil {
t.Fatal(err)
}
plain.Status.Conditions = chain.Status.Conditions
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(chain, plain).WithStatusSubresource(&batchv1.Job{}).Build()
k := NewK8sJobStatus(c, "minecraft")
ctx := context.Background()
got, err := k.PendingRestoreChains(ctx)
if err != nil {
t.Fatal(err)
}
want := RestoreChain{Job: "backup-survival-aa", Server: "survival", BackupID: "bk-1",
BackupRef: "/backups/survival/a.tar.gz", Snapshot: ChainSnapshotSucceeded}
if len(got) != 1 || got[0] != want {
t.Fatalf("pending = %+v, want [%+v]", got, want)
}
var jobs batchv1.JobList
if err := c.List(ctx, &jobs); err != nil {
t.Fatal(err)
}
if kind, held := maintenance.Holder("survival", nil, jobs.Items, time.Now()); !held || kind != maintenance.KindRestore {
t.Fatalf("pending chain: Holder = %q, %v", kind, held)
}
if err := k.SettleRestoreChain(ctx, "backup-survival-aa", maintenance.ThenRestoreAbandoned, ChainAbandonServerStarted); err != nil {
t.Fatal(err)
}
var settledJob batchv1.Job
if err := c.Get(ctx, types.NamespacedName{Namespace: "minecraft", Name: "backup-survival-aa"}, &settledJob); err != nil {
t.Fatal(err)
}
if settledJob.Labels[maintenance.LabelThenRestore] != maintenance.ThenRestoreAbandoned ||
settledJob.Annotations[maintenance.AnnotationThenRestoreReason] != ChainAbandonServerStarted ||
settledJob.Annotations[maintenance.AnnotationRestoreRef] == "" {
t.Fatalf("settled job meta: labels %v annotations %v", settledJob.Labels, settledJob.Annotations)
}
if got, _ := k.PendingRestoreChains(ctx); len(got) != 0 {
t.Fatalf("still pending after settling: %+v", got)
}
if err := c.List(ctx, &jobs); err != nil {
t.Fatal(err)
}
if kind, held := maintenance.Holder("survival", nil, jobs.Items, time.Now()); held {
t.Fatalf("settled chain still holds the volume as %q", kind)
}
aj, ok := jobToAsyncJob(&settledJob)
if !ok || aj.ThenRestore != maintenance.ThenRestoreAbandoned || aj.RestoreBackupID != "bk-1" || aj.State != "succeeded" ||
aj.ThenRestoreReason != ChainAbandonServerStarted || aj.Message != chainAbandonText[ChainAbandonServerStarted] {
t.Fatalf("projection = %+v", aj)
}
}
+19 -8
View File
@@ -1,6 +1,9 @@
package api
import "context"
import (
"context"
"errors"
)
// Restorer starts a world restore from a stored backup (spec §466: former_owner
// 3mo 内重新 claim → restore PVC). Restore only STARTS the work: recreating the
@@ -12,15 +15,23 @@ import "context"
// call returns once the restore is enqueued, so the handler answers 202
// (restoring), never claiming the world is already back.
//
// It returns ErrNotFound if the server is unknown to the execution backend; any
// other error is an internal failure (the handler maps it to 500).
// It returns ErrNotFound if the server is unknown to the execution backend, and
// an error whose RestoreInProgress method reports true when a restore of another
// backup is still running on the world (409 restore_in_progress; see
// isRestoreInProgress). Any other error is an internal failure (500).
//
// It is an interface so the handler is tested against a fake (api_test.go). The
// production executor — a restore Job mirroring internal/build's jobspec + weak-SA
// isolation — is integration-only and is a deliberate follow-up: until it is wired
// the API.Restorer is nil and POST /servers/{name}/restore-backup reports 503, so
// the restore authorization boundary is exercised without shipping a stub that
// cannot run in the real cluster topology.
// production executor is internal/restore's weak-SA restore Job; when cmd/felis
// cannot configure it the API.Restorer is nil and POST
// /servers/{name}/restore-backup reports 503.
type Restorer interface {
Restore(ctx context.Context, serverName, backupRef string) error
}
// isRestoreInProgress reports whether err says another restore is still running
// on the world. The executor cannot import this package, so its error is matched
// by method rather than by value.
func isRestoreInProgress(err error) bool {
var busy interface{ RestoreInProgress() bool }
return errors.As(err, &busy) && busy.RestoreInProgress()
}
+18
View File
@@ -187,6 +187,24 @@ func (b *Backuper) Backup(ctx context.Context, serverName, formerOwner string) e
return nil
}
// BackupThenRestore enqueues the safety snapshot in front of a restore: a backup
// Job like Backup's, recorded as a pre_restore backup and labelled with the
// restore to run once it succeeds (backupID, backupRef). felis-api creates that
// restore Job when the snapshot finishes and gives it up if the snapshot fails,
// so the world is never overwritten without a way back; the snapshot Job holds
// the world volume as a restore until then (internal/maintenance).
func (b *Backuper) BackupThenRestore(ctx context.Context, serverName, formerOwner, backupID, backupRef string) error {
p := b.jobParams(serverName, formerOwner)
p.RestoreRef, p.RestoreBackupID = backupRef, backupID
if err := b.Jobs.CreateBackupJob(ctx, p); err != nil {
if errors.Is(err, ErrAlreadyExists) {
return nil // suffix collision — treat as enqueued
}
return err
}
return nil
}
// jobNameSuffix is a short random hex tag that makes each backup Job name unique.
// 32 bits is ample: collisions only matter within a single Job's TTL window across
// a handful of manual backups.
+18
View File
@@ -41,3 +41,21 @@ func TestBackupMintsUniqueJobNamePerCall(t *testing.T) {
t.Errorf("two backups reused the Job name %q — retry would silently no-op", jobs.got[0].JobName)
}
}
func TestBackupThenRestoreChainsTheRestore(t *testing.T) {
jobs := &captureJobs{}
b := &Backuper{Jobs: jobs, Config: Config{Image: "img", BackupPVC: "pvc"}}
if err := b.BackupThenRestore(context.Background(), "survival", "usr-1", "bk-1", "/backups/a.tar.gz"); err != nil {
t.Fatalf("BackupThenRestore: %v", err)
}
if len(jobs.got) != 1 {
t.Fatalf("created %d jobs, want 1", len(jobs.got))
}
p := jobs.got[0]
if p.RestoreRef != "/backups/a.tar.gz" || p.RestoreBackupID != "bk-1" || p.FormerOwner != "usr-1" {
t.Errorf("params = %+v", p)
}
if !strings.HasPrefix(p.JobName, BackupJobName("survival")+"-") {
t.Errorf("JobName %q", p.JobName)
}
}
+39 -5
View File
@@ -20,6 +20,16 @@ const (
managedByValue = "felis-backup"
componentValue = "world-backup"
// The restore chain a safety snapshot carries (internal/maintenance keeps the
// canonical copies; maintenance_test pins these against them).
labelThenRestore = "felis.lolicon.best/then-restore"
thenRestorePending = "pending"
annotationRestoreRef = "felis.lolicon.best/restore-ref"
annotationRestoreBackupID = "felis.lolicon.best/restore-backup-id"
// ReasonPreRestore is the world_backups reason of a safety snapshot.
ReasonPreRestore = "pre_restore"
worldVolume = "world"
backupVolume = "backup"
configVolume = "config"
@@ -55,6 +65,14 @@ type JobParams struct {
RunAsGroup int64
FSGroup int64
// RestoreRef / RestoreBackupID make this backup the safety snapshot in front
// of a restore: the Job is labelled as a pending chain and names the backup
// felis-api restores once the snapshot succeeds (internal/maintenance). The
// snapshot is recorded as ReasonPreRestore, and its prune spares the backup
// the restore will extract.
RestoreRef string
RestoreBackupID string
TTLAfterFinished time.Duration
}
@@ -106,6 +124,9 @@ func BackupJob(p JobParams) (*batchv1.Job, error) {
if p.ConfigSecret == "" {
return nil, fmt.Errorf("backup: config secret name is required")
}
if p.RestoreRef != "" && p.RestoreBackupID == "" {
return nil, fmt.Errorf("backup: a chained restore needs the backup id")
}
limits, err := resourceLimits(p.CPULimit, p.MemLimit)
if err != nil {
return nil, err
@@ -129,6 +150,9 @@ func BackupJob(p JobParams) (*batchv1.Job, error) {
if p.FormerOwner != "" {
args = append(args, "--former-owner", p.FormerOwner)
}
if p.RestoreRef != "" {
args = append(args, "--reason", ReasonPreRestore, "--protect", p.RestoreBackupID)
}
container := corev1.Container{
Name: "backup",
@@ -167,12 +191,22 @@ func BackupJob(p JobParams) (*batchv1.Job, error) {
if name == "" {
name = BackupJobName(p.Server)
}
meta := metav1.ObjectMeta{
Name: name,
Namespace: p.Namespace,
Labels: backupLabels(p),
}
if p.RestoreRef != "" {
// On the Job only: felis-api settles the chain by patching this label, and
// the pods never need it.
meta.Labels[labelThenRestore] = thenRestorePending
meta.Annotations = map[string]string{
annotationRestoreRef: p.RestoreRef,
annotationRestoreBackupID: p.RestoreBackupID,
}
}
job := &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: p.Namespace,
Labels: backupLabels(p),
},
ObjectMeta: meta,
Spec: batchv1.JobSpec{
// One shot: a wedged archive must not loop. The TTL GCs the finished Job
// so a later backup of the same server is not blocked forever by a stale
+39
View File
@@ -199,6 +199,44 @@ func TestBackupJobArgsCarryServerAndOwner(t *testing.T) {
}
}
// A safety snapshot records itself as pre_restore, spares the backup the chained
// restore extracts from its prune, and carries the chain on the Job (only there:
// felis-api settles it by patching the Job's label).
func TestBackupJobCarriesTheRestoreChain(t *testing.T) {
p := sampleJobParams()
p.RestoreRef, p.RestoreBackupID = "/backups/survival/a.tar.gz", "bk-1"
job, err := BackupJob(p)
if err != nil {
t.Fatalf("BackupJob: %v", err)
}
args := job.Spec.Template.Spec.Containers[0].Args
if !argsContain(args, "--reason", ReasonPreRestore) || !argsContain(args, "--protect", "bk-1") {
t.Errorf("args = %v, want --reason %s --protect bk-1", args, ReasonPreRestore)
}
if job.Labels[labelThenRestore] != thenRestorePending {
t.Errorf("job labels = %v, want %s=%s", job.Labels, labelThenRestore, thenRestorePending)
}
if _, ok := job.Spec.Template.Labels[labelThenRestore]; ok {
t.Errorf("pod template carries the chain label: %v", job.Spec.Template.Labels)
}
if job.Annotations[annotationRestoreRef] != p.RestoreRef || job.Annotations[annotationRestoreBackupID] != "bk-1" {
t.Errorf("job annotations = %v", job.Annotations)
}
plain, err := BackupJob(sampleJobParams())
if err != nil {
t.Fatalf("BackupJob(plain): %v", err)
}
if _, ok := plain.Labels[labelThenRestore]; ok || len(plain.Annotations) != 0 {
t.Errorf("a plain backup carries a chain: labels %v annotations %v", plain.Labels, plain.Annotations)
}
for _, a := range plain.Spec.Template.Spec.Containers[0].Args {
if a == "--reason" || a == "--protect" {
t.Errorf("a plain backup passes %s: %v", a, plain.Spec.Template.Spec.Containers[0].Args)
}
}
}
func TestBackupJobRejectsMissingInputs(t *testing.T) {
for _, tc := range []struct {
name string
@@ -208,6 +246,7 @@ func TestBackupJobRejectsMissingInputs(t *testing.T) {
{"no world pvc", func(p *JobParams) { p.WorldPVC = "" }},
{"no backup pvc", func(p *JobParams) { p.BackupPVC = "" }},
{"no config secret", func(p *JobParams) { p.ConfigSecret = "" }},
{"chain without backup id", func(p *JobParams) { p.RestoreRef = "/backups/a.tar.gz" }},
} {
t.Run(tc.name, func(t *testing.T) {
p := sampleJobParams()
+43 -4
View File
@@ -20,6 +20,12 @@
// A lock older than Grace with no Job behind it is stale (felis-api died between
// the two writes) and holds nothing.
//
// A restore that starts with a safety snapshot is two Jobs in a row: the backup
// Job carries the restore to run after it (LabelThenRestore), and felis-api
// creates the restore Job once the backup has succeeded. The backup Job keeps
// holding the volume, as a restore, from its creation until felis-api has
// settled what follows it, so nothing can wake the server between the two Jobs.
//
// File reads and listings are not holders. They mount the volume read-only for a
// second or two, and a server starting beside one cannot hurt either side, so
// nobody waits for them.
@@ -50,6 +56,26 @@ const (
// LabelFilesMode is the file-editor operation (list, read, write) a files Job
// performs. Only write holds the volume.
LabelFilesMode = "felis.lolicon.best/files-mode"
// LabelThenRestore marks a backup Job that is the safety snapshot in front of
// a restore. Its value is the chain's state: ThenRestorePending until felis-api
// settles it, then ThenRestoreStarted or ThenRestoreAbandoned. A label, so the
// settling loop finds pending chains with a selector.
LabelThenRestore = "felis.lolicon.best/then-restore"
// AnnotationRestoreRef / AnnotationRestoreBackupID name the backup the chained
// restore extracts: its archive ref and its world_backups id.
AnnotationRestoreRef = "felis.lolicon.best/restore-ref"
AnnotationRestoreBackupID = "felis.lolicon.best/restore-backup-id"
// AnnotationThenRestoreReason is the code saying why a chain was abandoned
// (internal/api defines the codes).
AnnotationThenRestoreReason = "felis.lolicon.best/then-restore-reason"
)
// States of LabelThenRestore.
const (
ThenRestorePending = "pending"
ThenRestoreStarted = "started"
ThenRestoreAbandoned = "abandoned"
)
// Kinds of holder.
@@ -82,6 +108,13 @@ func JobKind(j *batchv1.Job) (string, bool) {
return "", false
}
// RestorePending reports whether j is a safety-snapshot backup Job whose restore
// felis-api has not yet started or abandoned. Such a Job holds the volume as a
// restore whether or not it has finished.
func RestorePending(j *batchv1.Job) bool {
return j.Labels[LabelManagedBy] == "felis-backup" && j.Labels[LabelThenRestore] == ThenRestorePending
}
// JobFinished reports whether a Job has reached a terminal condition. The
// success/failure-target conditions count as terminal: the Job controller sets
// them the moment the outcome is decided, before it finishes tearing the pods
@@ -119,13 +152,19 @@ func parseLock(v string) (string, time.Time, bool) {
}
// Holder reports what, if anything, holds the server's world volume at `now`:
// the first unfinished maintenance Job among jobs, else a lock in annotations
// younger than Grace. jobs may contain unrelated Jobs; only the server's own
// holders count.
// the first unfinished maintenance Job among jobs (or a safety snapshot whose
// restore is still pending), else a lock in annotations younger than Grace. jobs
// may contain unrelated Jobs; only the server's own holders count.
func Holder(server string, annotations map[string]string, jobs []batchv1.Job, now time.Time) (string, bool) {
for i := range jobs {
j := &jobs[i]
if j.Labels[LabelServer] != server || JobFinished(j) {
if j.Labels[LabelServer] != server {
continue
}
if RestorePending(j) {
return KindRestore, true
}
if JobFinished(j) {
continue
}
if kind, ok := JobKind(j); ok {
+46
View File
@@ -125,6 +125,52 @@ func TestHolderFromJobs(t *testing.T) {
}
}
// A safety snapshot holds the volume as a restore from its creation until
// felis-api settles the chain, finished or not, so nothing wakes the server
// between the backup Job and the restore Job it is followed by.
func TestHolderFromRestoreChain(t *testing.T) {
j, err := backupjob.BackupJob(backupjob.JobParams{
Server: "survival", WorldPVC: "world-survival-0", BackupPVC: "felis-backups",
Namespace: "minecraft", Image: "felis:1", ConfigSecret: "felis-config", ConfigMount: "/etc/felis",
RestoreRef: "/backups/survival/a.tar.gz", RestoreBackupID: "bk-1",
})
if err != nil {
t.Fatalf("BackupJob: %v", err)
}
if j.Annotations[AnnotationRestoreRef] != "/backups/survival/a.tar.gz" || j.Annotations[AnnotationRestoreBackupID] != "bk-1" {
t.Fatalf("chain annotations = %v", j.Annotations)
}
if !RestorePending(j) {
t.Fatalf("a fresh safety snapshot is not a pending chain: labels %v", j.Labels)
}
for _, tc := range []struct {
name string
job batchv1.Job
held bool
kind string
}{
{"running", *j, true, KindRestore},
{"succeeded, restore not started yet", finished(*j, batchv1.JobComplete), true, KindRestore},
{"failed, not settled yet", finished(*j, batchv1.JobFailed), true, KindRestore},
{"restore started", settled(finished(*j, batchv1.JobComplete), ThenRestoreStarted), false, ""},
{"abandoned", settled(finished(*j, batchv1.JobFailed), ThenRestoreAbandoned), false, ""},
{"abandoned while running", settled(*j, ThenRestoreAbandoned), true, KindBackup},
} {
kind, held := Holder("survival", nil, []batchv1.Job{tc.job}, now)
if held != tc.held || kind != tc.kind {
t.Errorf("%s: Holder = %q, %v; want %q, %v", tc.name, kind, held, tc.kind, tc.held)
}
}
if plain := backupJob(t, "survival"); RestorePending(&plain) {
t.Fatal("a plain backup is a pending chain")
}
}
func settled(j batchv1.Job, state string) batchv1.Job {
j.Labels = map[string]string{LabelServer: j.Labels[LabelServer], LabelManagedBy: j.Labels[LabelManagedBy], LabelThenRestore: state}
return j
}
func TestHolderFromLock(t *testing.T) {
for _, tc := range []struct {
name string
+13 -2
View File
@@ -187,13 +187,24 @@ func TestManualBackupRationing(t *testing.T) {
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)
excess, err := st.ExcessBackups(ctx, name, "manual", 5, "")
if err != nil {
t.Fatalf("ExcessManualBackups: %v", err)
t.Fatalf("ExcessBackups: %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)
}
// The backup a chained restore will extract is never pruned.
excess, err = st.ExcessBackups(ctx, name, "manual", 5, manual[0])
if err != nil {
t.Fatalf("ExcessBackups(protect): %v", err)
}
if len(excess) != 1 || excess[0].ID != manual[1] {
t.Fatalf("excess with %s protected = %+v; want only %s", manual[0], excess, manual[1])
}
if excess, err = st.ExcessBackups(ctx, name, "pre_restore", 0, ""); err != nil || len(excess) != 0 {
t.Fatalf("pre_restore excess = %+v, %v; the manual ones are not its to prune", excess, err)
}
all, err := st.EvictableBackups(ctx)
if err != nil {
+4 -2
View File
@@ -104,8 +104,10 @@ func APIMinecraftRole(p Params) *rbacv1.Role {
// nothing in felis-api lists or deletes PVCs.
rule([]string{groupCore}, []string{"persistentvolumeclaims"}, []string{"get"}),
// list backs GET /servers/{name}/jobs — the async status outlet reads the
// backup/restore Jobs back by the server label (read-only).
rule([]string{groupBatch}, []string{"jobs"}, []string{"create", "get", "delete", "list"}),
// backup/restore Jobs back by the server label — and finds the pending
// restore chains; patch settles a chain by relabelling its safety-snapshot
// Job (internal/api restorechain.go).
rule([]string{groupBatch}, []string{"jobs"}, []string{"create", "get", "delete", "list", "patch"}),
// Read-side console (spec §8 读=pods/log follow): list pods to find the
// server's running pod, then read its log subresource. Two separate rules so
// the verbs stay tight — list on pods, get on pods/log, and nothing else.
+1 -1
View File
@@ -73,7 +73,7 @@ func TestAPIRole_CreatesJobsInBothNamespaces(t *testing.T) {
if mc.Namespace != "minecraft" {
t.Errorf("felis-api minecraft Role namespace = %q, want minecraft", mc.Namespace)
}
for _, v := range []string{"create", "get", "delete", "list"} {
for _, v := range []string{"create", "get", "delete", "list", "patch"} {
if !hasRule(mc, "batch", "jobs", v) {
t.Errorf("felis-api (minecraft) must have batch/jobs:%s for the restore-Job lifecycle", v)
}
+8 -6
View File
@@ -113,13 +113,15 @@ func (s *PGStore) EvictableBackups(ctx context.Context) ([]StoredBackup, error)
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) {
// ExcessBackups lists server's present backups of one reason beyond the newest
// keep, oldest first: what the backup Job removes after adding one. protect, when
// set, is a backup id left out of the list whatever its age (the one a chained
// restore is about to extract).
func (s *PGStore) ExcessBackups(ctx context.Context, server, reason string, keep int, protect string) ([]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)
WHERE server_name = $1 AND status = 'present' AND reason = $2 AND id <> $4
ORDER BY created_at DESC OFFSET $3`
out, err := s.queryBackups(ctx, q, server, reason, keep, protect)
for i, j := 0, len(out)-1; i < j; i, j = i+1, j-1 {
out[i], out[j] = out[j], out[i]
}
+9 -3
View File
@@ -20,6 +20,11 @@ const (
managedByValue = "felis-restore"
componentValue = "world-restore"
// AnnotationBackupRef records which archive the Job extracts, so a second
// restore that collides with it on the name can tell a duplicate of the same
// request from a request for a different backup (K8sJobs.CreateRestoreJob).
AnnotationBackupRef = "felis.lolicon.best/backup-ref"
worldVolume = "world"
backupVolume = "backup"
felisBinaryPath = "/usr/local/bin/felis"
@@ -144,9 +149,10 @@ func RestoreJob(p JobParams) (*batchv1.Job, error) {
job := &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
Name: RestoreJobName(p.Server),
Namespace: p.Namespace,
Labels: restoreLabels(p),
Name: RestoreJobName(p.Server),
Namespace: p.Namespace,
Labels: restoreLabels(p),
Annotations: map[string]string{AnnotationBackupRef: p.BackupRef},
},
Spec: batchv1.JobSpec{
// One shot: a bad archive must not loop. The TTL GCs the finished Job
+3
View File
@@ -224,6 +224,9 @@ func TestRestoreJobInvokesFelisRestoreWithParams(t *testing.T) {
if !argPairPresent(c.Args, "--ref", p.BackupRef) {
t.Errorf("args must carry --ref %q, got %v", p.BackupRef, c.Args)
}
if got := job.Annotations[AnnotationBackupRef]; got != p.BackupRef {
t.Errorf("backup-ref annotation = %q, want %q", got, p.BackupRef)
}
if !argPairPresent(c.Args, "--archive-store", p.ArchiveStore) {
t.Errorf("args must carry --archive-store %q, got %v", p.ArchiveStore, c.Args)
}
+26 -3
View File
@@ -6,7 +6,9 @@ import (
"time"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
)
@@ -37,8 +39,10 @@ func NewK8sJobs(c client.Client) *K8sJobs {
// of the same server collides on Create. The collision is answered by the state
// of the Job already holding the name:
//
// - still running (or not yet started): ErrAlreadyExists, which the Restorer
// treats as success — the idempotent coalesce.
// - still running (or not yet started) on the same archive: ErrAlreadyExists,
// which the Restorer treats as success — the idempotent coalesce.
// - still running on another archive: ErrOtherRestoreRunning. A Job from
// before the ref annotation cannot be compared and coalesces as before.
// - finished (succeeded OR failed): the finished Job is deleted and replaced,
// so the caller's retry enqueues for real. Without this, the deterministic
// name plus the ten-minute TTL would swallow the retry — most importantly
@@ -68,9 +72,14 @@ func (k *K8sJobs) CreateRestoreJob(ctx context.Context, p JobParams) error {
return getErr
}
if !restoreJobFinished(&existing) {
if ref, ok := existing.Annotations[AnnotationBackupRef]; ok && ref != p.BackupRef {
return ErrOtherRestoreRunning
}
return ErrAlreadyExists
}
if deleteErr := k.c.Delete(ctx, &existing); deleteErr != nil && !apierrors.IsNotFound(deleteErr) {
// Background propagation: a Job deleted with the API's default policy
// orphans its pods, which then outlive it for good.
if deleteErr := k.c.Delete(ctx, &existing, client.PropagationPolicy(metav1.DeletePropagationBackground)); deleteErr != nil && !apierrors.IsNotFound(deleteErr) {
return deleteErr
}
// The API server keeps the object until its job-tracking finalizer has run,
@@ -128,7 +137,21 @@ func (k *K8sJobs) recreate(ctx context.Context, job *batchv1.Job) error {
// that is merely created-but-not-started (no active pods yet, no completions)
// counts as in flight, not finished, so a duplicate enqueue during startup still
// coalesces.
//
// A terminal condition wins over pods still shutting down: the world-volume lock
// (internal/maintenance.JobFinished) already lets the next restore in at that
// point, and answering it with the coalesce would be a 202 for a restore that
// never runs.
func restoreJobFinished(job *batchv1.Job) bool {
for _, c := range job.Status.Conditions {
if c.Status != corev1.ConditionTrue {
continue
}
switch c.Type {
case batchv1.JobComplete, batchv1.JobFailed, batchv1.JobSuccessCriteriaMet, batchv1.JobFailureTarget:
return true
}
}
if job.Status.Active > 0 {
return false
}
+36
View File
@@ -103,6 +103,12 @@ func TestRestoreJobFinished(t *testing.T) {
Conditions: []batchv1.JobCondition{{Type: batchv1.JobFailed, Status: "True"}},
}}, true},
{"failed-and-some-active", batchv1.Job{Status: batchv1.JobStatus{Active: 1, Failed: 1}}, false},
// The failure is decided while the pod is still being torn down: the
// world-volume lock already admits the next restore, so this must too.
{"failure-target-pod-terminating", batchv1.Job{Status: batchv1.JobStatus{
Active: 1,
Conditions: []batchv1.JobCondition{{Type: batchv1.JobFailureTarget, Status: "True"}},
}}, true},
}
for _, tc := range cases {
if got := restoreJobFinished(&tc.job); got != tc.want {
@@ -110,3 +116,33 @@ func TestRestoreJobFinished(t *testing.T) {
}
}
}
// A running restore of a different archive refuses the request instead of
// absorbing it: the caller would otherwise get a 202 naming the backup it chose
// while another one is extracted.
func TestCreateRestoreJobRefusesAnotherArchiveInFlight(t *testing.T) {
inFlight := &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
Name: "restore-survival", Namespace: "minecraft",
Annotations: map[string]string{AnnotationBackupRef: "/backups/other.tar.gz"},
},
Status: batchv1.JobStatus{Active: 1},
}
c := fake.NewClientBuilder().WithScheme(testScheme(t)).WithObjects(inFlight).Build()
err := NewK8sJobs(c).CreateRestoreJob(context.Background(), testParams())
if !errors.Is(err, ErrOtherRestoreRunning) {
t.Fatalf("CreateRestoreJob = %v, want ErrOtherRestoreRunning", err)
}
if err := (&Restorer{Jobs: NewK8sJobs(c), Config: Config{Image: "felis:test", BackupPVC: "felis-backups"}}).
Restore(context.Background(), "survival", "/backups/x.tar.gz"); !errors.Is(err, ErrOtherRestoreRunning) {
t.Fatalf("Restore = %v, want ErrOtherRestoreRunning", err)
}
// The same archive still coalesces.
p := testParams()
p.BackupRef = "/backups/other.tar.gz"
if err := NewK8sJobs(c).CreateRestoreJob(context.Background(), p); !errors.Is(err, ErrAlreadyExists) {
t.Fatalf("same archive: CreateRestoreJob = %v, want ErrAlreadyExists", err)
}
}
+30 -10
View File
@@ -39,6 +39,22 @@ import (
// as success — see Restore.
var ErrAlreadyExists = errors.New("restore: job already exists")
// ErrOtherRestoreRunning is returned when a restore of a DIFFERENT backup is
// still running on the server. The caller's request did not take effect; the
// API answers 409 restore_in_progress rather than a 202 that names the backup
// it asked for.
var ErrOtherRestoreRunning error = otherRestoreRunning{}
type otherRestoreRunning struct{}
func (otherRestoreRunning) Error() string {
return "restore: a restore of another backup is still running"
}
// RestoreInProgress marks the error for internal/api, which cannot import this
// package and recognises it by the method.
func (otherRestoreRunning) RestoreInProgress() bool { return true }
// Jobs is the cluster-side restore lifecycle the Restorer depends on. It is an
// interface so the orchestration is tested against a fake; the controller-runtime
// implementation (K8sJobs) is integration-tested only — it requires a live
@@ -47,7 +63,9 @@ var ErrAlreadyExists = errors.New("restore: job already exists")
// whole contract, exactly matching the asynchronous 202 the handler answers.
type Jobs interface {
// CreateRestoreJob renders and applies the restore Job for p. It returns
// ErrAlreadyExists if a Job of the same (deterministic) name already exists.
// ErrAlreadyExists if an unfinished Job of the same (deterministic) name is
// already restoring p.BackupRef, and ErrOtherRestoreRunning if it is
// restoring another archive.
CreateRestoreJob(ctx context.Context, p JobParams) error
}
@@ -162,16 +180,18 @@ type Restorer struct {
// serverName's world PVC. It returns once the Job is created — the extraction
// runs in the Pod — so the handler's 202 ("restoring") is honest.
//
// It is idempotent: a duplicate enqueue while a restore Job for this server is
// still running is treated as success rather than surfaced as an error.
// It is idempotent: a duplicate enqueue of the same archive while a restore Job
// for this server is still running is treated as success rather than surfaced
// as an error.
//
// The coalescing key is the Job name (RestoreJobName), which depends only on the
// server, NOT on backupRef — so a second request that arrives while one is in
// flight is absorbed regardless of the ref it carries, and if the two refs
// differ the second is silently dropped (the in-flight restore wins). That is
// acceptable here: restore runs only for a Stopped server (handler gate ⑥) and
// the handler always passes the latest backup, which for a stopped server does
// not change, so concurrent requests carry the same ref in practice.
// The Job name (RestoreJobName) depends only on the server, and the handler lets
// the caller pick any of the server's backups, so a collision can carry a
// different ref. The world-volume lock refuses a second restore while the first
// Job runs, which makes this rare (it takes the lock seeing the Job finished
// while its pods are still going), but when it happens the running Job's ref
// annotation decides: the same archive coalesces, another archive is
// ErrOtherRestoreRunning, so no caller is told "restoring" for a backup that is
// not the one being extracted.
//
// A FINISHED Job — succeeded or failed — does not absorb the next request: its
// deterministic name is replaced so the retry enqueues for real (see