From 89a013470781a123a0187a93f9d002c9a877a5b1 Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Fri, 25 Sep 2026 03:16:06 +0800 Subject: [PATCH] =?UTF-8?q?feat(backup):=20=E5=BD=92=E6=A1=A3=E5=85=88?= =?UTF-8?q?=E5=86=99=20.partial=20=E5=86=8D=20fsync=20=E6=94=B9=E5=90=8D?= =?UTF-8?q?=E5=B9=B6=E8=AE=B0=E5=BD=95=20sha256=EF=BC=8C=E5=88=A0=E9=99=A4?= =?UTF-8?q?=E5=8E=9F=E4=BB=B6=E5=89=8D=E5=9B=9E=E8=AF=BB=E6=A0=A1=E9=AA=8C?= =?UTF-8?q?=EF=BC=8Creaper=20=E6=8A=BD=E6=A0=B7=E5=B7=A1=E6=A3=80=E4=B8=8E?= =?UTF-8?q?=E6=B8=85=E6=89=AB=E5=8D=8A=E6=88=AA=E5=BD=92=E6=A1=A3=EF=BC=8C?= =?UTF-8?q?=E4=BF=9D=E7=95=99=E6=9D=83=E9=99=90=E4=B8=8E=20mtime=EF=BC=8C?= =?UTF-8?q?=E9=9D=A2=E6=9D=BF=E6=A0=87=E5=87=BA=E6=8D=9F=E5=9D=8F=E4=B8=8E?= =?UTF-8?q?=E5=B7=B2=E6=A0=A1=E9=AA=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/felis/backup.go | 11 +- cmd/felis/reaper.go | 18 +- cmd/felis/reaper_test.go | 4 + cmd/felis/restore_test.go | 4 +- docs/openapi.yaml | 12 +- internal/api/api_test.go | 4 +- internal/api/handlers_backups.go | 7 + internal/api/handlers_backups_test.go | 23 ++ internal/api/pgrepo.go | 24 +- internal/api/repo.go | 9 + internal/backup/archiver.go | 61 ++- internal/backup/tarlocal.go | 359 +++++++++++++++--- internal/backup/tarlocal_test.go | 279 +++++++++++++- internal/pgint/reaper_test.go | 147 ++++++- internal/reaper/pgstore.go | 61 ++- internal/reaper/reaper.go | 192 +++++++++- internal/reaper/reaper_test.go | 261 ++++++++++++- .../0026_world_backups_integrity.sql | 16 + panel/src/i18n/resources/en-US/backups.json | 8 + panel/src/i18n/resources/en-US/errors.json | 1 + panel/src/i18n/resources/zh-CN/backups.json | 7 + panel/src/i18n/resources/zh-CN/errors.json | 1 + panel/src/lib/api.ts | 4 + panel/src/lib/types.ts | 19 +- panel/src/pages/ServerBackups.tsx | 82 +++- 25 files changed, 1489 insertions(+), 125 deletions(-) create mode 100644 internal/store/migrations/0026_world_backups_integrity.sql diff --git a/cmd/felis/backup.go b/cmd/felis/backup.go index e8a1420..65b7b32 100644 --- a/cmd/felis/backup.go +++ b/cmd/felis/backup.go @@ -7,6 +7,7 @@ import ( "flag" "fmt" "io" + "strings" "time" "felis.lolicon.best/internal/backup" @@ -91,11 +92,16 @@ func cmdBackup(args []string, stdout, stderr io.Writer) int { return 1 } - ref, size, err := archiver.Archive(ctx, *server, naming.WorldPVCName(*server)) + a, err := archiver.Archive(ctx, *server, naming.WorldPVCName(*server)) if err != nil { fmt.Fprintf(stderr, "felis backup: archive: %v\n", err) return 1 } + ref, size := a.Ref, a.Size + if len(a.Skipped) > 0 { + fmt.Fprintf(stderr, "felis backup: %d entries are not plain files or directories and are not in the archive: %s\n", + len(a.Skipped), strings.Join(a.Skipped[:min(len(a.Skipped), 10)], ", ")) + } drv, err := store.Open(ctx, cfg.Database.URL) if err != nil { @@ -112,6 +118,9 @@ func cmdBackup(args []string, stdout, stderr io.Writer) int { SizeBytes: size, Reason: *reason, ExpiresAt: time.Now().Add(rcfg.ManualRetention), + + SHA256: a.SHA256, + SkippedEntries: len(a.Skipped), } st := reaper.NewPGStore(drv.DB()) if err := st.InsertBackup(ctx, rec); err != nil { diff --git a/cmd/felis/reaper.go b/cmd/felis/reaper.go index a57b4d8..f82aa30 100644 --- a/cmd/felis/reaper.go +++ b/cmd/felis/reaper.go @@ -138,14 +138,24 @@ func cmdReaper(args []string, stdout, stderr io.Writer) int { // FelisWorldJobFailed rule) reach the operator: a world that cannot be archived // is kept, and without this nobody would learn that it is never reaped. func reportReaperRun(sum reaper.Summary, stdout, stderr io.Writer) int { - fmt.Fprintf(stdout, "felis reaper: evaluated=%d reaped=%d awaiting_offsite=%d warned=%d skipped=%d store_full=%d evicted=%d expired=%d expire_failed=%d\n", + fmt.Fprintf(stdout, "felis reaper: evaluated=%d reaped=%d awaiting_offsite=%d warned=%d skipped=%d store_full=%d evicted=%d expired=%d expire_failed=%d verified=%d corrupt=%d verify_failed=%d swept=%d orphan_archives=%d\n", sum.Evaluated, sum.WorldsReaped, sum.AwaitingOffsite, sum.Warned, sum.Skipped, sum.StoreFull, - sum.EvictedEarly, sum.BackupsExpired, sum.ExpireFailed) + sum.EvictedEarly, sum.BackupsExpired, sum.ExpireFailed, + sum.Verified, sum.Corrupt, sum.VerifyFailed, sum.Swept, sum.OrphanArchives) if !sum.Failed() { return 0 } - fmt.Fprintf(stderr, "felis reaper: %d servers failed (%d kept because the backup store is full) and %d expired backups were not removed; the errors are above, and each is retried next run\n", - sum.Skipped, sum.StoreFull, sum.ExpireFailed) + if sum.Skipped > 0 || sum.ExpireFailed > 0 { + fmt.Fprintf(stderr, "felis reaper: %d servers failed (%d kept because the backup store is full) and %d expired backups were not removed; the errors are above, and each is retried next run\n", + sum.Skipped, sum.StoreFull, sum.ExpireFailed) + } + if sum.Corrupt > 0 { + fmt.Fprintf(stderr, "felis reaper: %d archives did not read back and are marked corrupt; they are no longer offered for restore (the errors are above)\n", sum.Corrupt) + } + if sum.VerifyFailed > 0 || sum.SweepFailed { + fmt.Fprintf(stderr, "felis reaper: %d archives could not be read back and the store sweep completed=%t; both are retried next run\n", + sum.VerifyFailed, !sum.SweepFailed) + } return 1 } diff --git a/cmd/felis/reaper_test.go b/cmd/felis/reaper_test.go index 85e4d86..e4c2c12 100644 --- a/cmd/felis/reaper_test.go +++ b/cmd/felis/reaper_test.go @@ -30,6 +30,10 @@ func TestReportReaperRunFailsTheJob(t *testing.T) { {"server failed", reaper.Summary{Evaluated: 3, Skipped: 1}, 1}, {"store full", reaper.Summary{Evaluated: 3, Skipped: 1, StoreFull: 1}, 1}, {"expiry failed", reaper.Summary{Evaluated: 3, ExpireFailed: 2}, 1}, + {"corrupt archive", reaper.Summary{Evaluated: 3, Verified: 4, Corrupt: 1}, 1}, + {"read-back failed", reaper.Summary{Evaluated: 3, VerifyFailed: 1}, 1}, + {"sweep failed", reaper.Summary{Evaluated: 3, SweepFailed: true}, 1}, + {"orphans kept", reaper.Summary{Evaluated: 3, Swept: 2, OrphanArchives: 1}, 0}, } { var out, errb bytes.Buffer if got := reportReaperRun(tc.sum, &out, &errb); got != tc.want { diff --git a/cmd/felis/restore_test.go b/cmd/felis/restore_test.go index 1f154d7..6147a49 100644 --- a/cmd/felis/restore_test.go +++ b/cmd/felis/restore_test.go @@ -32,11 +32,11 @@ func archiveTempWorld(t *testing.T, files map[string]string) (ref string, backup BackupRoot: backupRoot, Resolve: func(string) (string, error) { return srcDir, nil }, } - got, _, err := ar.Archive(context.Background(), "survival", "world-survival-0") + got, err := ar.Archive(context.Background(), "survival", "world-survival-0") if err != nil { t.Fatalf("Archive: %v", err) } - return string(got), backupRoot + return string(got.Ref), backupRoot } // The restore subcommand must extract the archived world into the target world diff --git a/docs/openapi.yaml b/docs/openapi.yaml index e3023fd..044ae27 100644 --- a/docs/openapi.yaml +++ b/docs/openapi.yaml @@ -368,6 +368,16 @@ components: status: { type: string } created_at: { type: string, format: date-time } expires_at: { type: string, format: date-time } + corrupt: + type: boolean + description: The archive failed a read-back (a checksum, gzip or tar error) and cannot be restored. Omitted when false. + verified_at: + type: string + format: date-time + description: The archive's last read-back that matched. Omitted until the first. + skipped_entries: + type: integer + description: World entries the archive could not hold (symbolic links, devices, sockets). Omitted when zero. Build: type: object @@ -3016,7 +3026,7 @@ paths: application/json: schema: { $ref: '#/components/schemas/Error' } '409': - description: Server is not stopped (not_stopped), a restore, backup or file write already holds its world volume (maintenance_in_progress), or a restore of another backup is still running (restore_in_progress). + description: Server is not stopped (not_stopped), a restore, backup or file write already holds its world volume (maintenance_in_progress), a restore of another backup is still running (restore_in_progress), or the chosen backup failed a read-back (backup_corrupt). content: application/json: schema: { $ref: '#/components/schemas/Error' } diff --git a/internal/api/api_test.go b/internal/api/api_test.go index 5c41abe..ac27b72 100644 --- a/internal/api/api_test.go +++ b/internal/api/api_test.go @@ -806,7 +806,7 @@ func (f *fakeRepo) LatestBackup(_ context.Context, serverName string) (*BackupRe var latest *fakeBackup for i := range f.backups { b := &f.backups[i] - if b.view.Status != "present" || b.view.ServerName != serverName { + if b.view.Status != "present" || b.view.ServerName != serverName || b.view.Corrupt { continue } if latest == nil || b.view.CreatedAt.After(latest.view.CreatedAt) { @@ -830,7 +830,7 @@ func (f *fakeRepo) BackupByID(_ context.Context, id string) (*BackupRecord, erro return &BackupRecord{ ID: b.view.ID, ServerName: b.view.ServerName, FormerOwner: b.view.FormerOwner, BackupRef: b.ref, - SizeBytes: b.view.SizeBytes, + SizeBytes: b.view.SizeBytes, Corrupt: b.view.Corrupt, }, nil } } diff --git a/internal/api/handlers_backups.go b/internal/api/handlers_backups.go index d06d1cc..65b7513 100644 --- a/internal/api/handlers_backups.go +++ b/internal/api/handlers_backups.go @@ -141,6 +141,13 @@ func (a *API) handleRestoreBackup(w http.ResponseWriter, r *http.Request) { writeError(w, r, errForbidden) return } + // The reaper found this archive damaged when it read it back; the restore + // Job would refuse it too, but only after the world was locked for it. + if backup.Corrupt { + writeError(w, r, newError(http.StatusConflict, "backup_corrupt", + "this backup did not read back intact and cannot be restored; pick another")) + return + } } else { backup, err = a.Repo.LatestBackup(r.Context(), name) if err != nil { diff --git a/internal/api/handlers_backups_test.go b/internal/api/handlers_backups_test.go index 63d85e8..b568df6 100644 --- a/internal/api/handlers_backups_test.go +++ b/internal/api/handlers_backups_test.go @@ -414,6 +414,29 @@ func TestRestoreBackup(t *testing.T) { } }) + t.Run("restore by backup_id that failed a read-back -> 409 backup_corrupt", func(t *testing.T) { + api, repo, _, restorer := mkTwo() + repo.backups[1].view.Corrupt = true + api.External = staticExternal{p: owner} + w := do(api.ExternalHandler(), "POST", path, `{"backup_id":"bk2"}`, jsonHeaders) + if w.Code != http.StatusConflict || decodeErr(t, w) != "backup_corrupt" { + t.Fatalf("code = %d body %s", w.Code, w.Body.String()) + } + if restorer.calls != 0 { + t.Fatal("a corrupt backup reached the restorer") + } + }) + + t.Run("no body skips a latest backup that failed a read-back", func(t *testing.T) { + api, repo, _, restorer := mkTwo() + repo.backups[0].view.Corrupt = true + api.External = staticExternal{p: owner} + w := do(api.ExternalHandler(), "POST", path, "", nil) + if w.Code != http.StatusAccepted || restorer.gotRef != "ref-bk2" { + t.Fatalf("code = %d ref %q, want 202 restoring the newest intact backup", w.Code, restorer.gotRef) + } + }) + t.Run("no body -> falls back to LatestBackup (backward compat)", func(t *testing.T) { api, _, _, restorer := mkTwo() api.External = staticExternal{p: owner} diff --git a/internal/api/pgrepo.go b/internal/api/pgrepo.go index a1d9ec7..b717e73 100644 --- a/internal/api/pgrepo.go +++ b/internal/api/pgrepo.go @@ -675,7 +675,7 @@ func (p *PGRepo) SeedServer(ctx context.Context, name, subdomain string, cpuMill // listed — an expired or deleted backup is gone (spec §466). func (p *PGRepo) AllBackups(ctx context.Context) ([]BackupView, error) { const q = `SELECT id, server_name, COALESCE(former_owner, ''), COALESCE(size_bytes, 0), - reason, status, created_at, expires_at + reason, status, created_at, expires_at, corrupt_at IS NOT NULL, verified_at, skipped_entries FROM world_backups WHERE status = 'present' ORDER BY created_at DESC` rows, err := p.db.QueryContext(ctx, q) if err != nil { @@ -689,7 +689,7 @@ func (p *PGRepo) AllBackups(ctx context.Context) ([]BackupView, error) { // former_owner never matches a user id, so orphaned backups stay admin-only. func (p *PGRepo) BackupsForUser(ctx context.Context, userID string) ([]BackupView, error) { const q = `SELECT id, server_name, COALESCE(former_owner, ''), COALESCE(size_bytes, 0), - reason, status, created_at, expires_at + reason, status, created_at, expires_at, corrupt_at IS NOT NULL, verified_at, skipped_entries FROM world_backups WHERE status = 'present' AND former_owner = $1 ORDER BY created_at DESC` rows, err := p.db.QueryContext(ctx, q, userID) if err != nil { @@ -705,21 +705,26 @@ func scanBackupViews(rows *sql.Rows) ([]BackupView, error) { var out []BackupView for rows.Next() { var v BackupView + var verified sql.NullTime if err := rows.Scan(&v.ID, &v.ServerName, &v.FormerOwner, &v.SizeBytes, - &v.Reason, &v.Status, &v.CreatedAt, &v.ExpiresAt); err != nil { + &v.Reason, &v.Status, &v.CreatedAt, &v.ExpiresAt, &v.Corrupt, &verified, &v.SkippedEntries); err != nil { return nil, err } + if verified.Valid { + v.VerifiedAt = &verified.Time + } out = append(out, v) } return out, rows.Err() } -// LatestBackup returns the most recent present backup for a server (spec §466 -// restore), or ErrNotFound. Unlike the list queries this selects backup_ref — the -// caller (the restore handler) hands it to the Restorer and never serializes it. +// LatestBackup returns the most recent present backup for a server that has not +// failed a read-back (spec §466 restore), or ErrNotFound. Unlike the list queries +// this selects backup_ref — the caller (the restore handler) hands it to the +// Restorer and never serializes it. func (p *PGRepo) LatestBackup(ctx context.Context, serverName string) (*BackupRecord, error) { const q = `SELECT id, server_name, COALESCE(former_owner, ''), backup_ref, COALESCE(size_bytes, 0) - FROM world_backups WHERE server_name = $1 AND status = 'present' + FROM world_backups WHERE server_name = $1 AND status = 'present' AND corrupt_at IS NULL ORDER BY created_at DESC LIMIT 1` var b BackupRecord switch err := p.db.QueryRowContext(ctx, q, serverName).Scan( @@ -734,11 +739,12 @@ func (p *PGRepo) LatestBackup(ctx context.Context, serverName string) (*BackupRe // BackupByID returns a single present backup by its id, or ErrNotFound. func (p *PGRepo) BackupByID(ctx context.Context, id string) (*BackupRecord, error) { - const q = `SELECT id, server_name, COALESCE(former_owner, ''), backup_ref, COALESCE(size_bytes, 0) + const q = `SELECT id, server_name, COALESCE(former_owner, ''), backup_ref, COALESCE(size_bytes, 0), + corrupt_at IS NOT NULL FROM world_backups WHERE id = $1 AND status = 'present'` var b BackupRecord switch err := p.db.QueryRowContext(ctx, q, id).Scan( - &b.ID, &b.ServerName, &b.FormerOwner, &b.BackupRef, &b.SizeBytes); { + &b.ID, &b.ServerName, &b.FormerOwner, &b.BackupRef, &b.SizeBytes, &b.Corrupt); { case errors.Is(err, sql.ErrNoRows): return nil, ErrNotFound case err != nil: diff --git a/internal/api/repo.go b/internal/api/repo.go index 41e65ab..4eb19d3 100644 --- a/internal/api/repo.go +++ b/internal/api/repo.go @@ -69,6 +69,13 @@ type BackupView struct { Status string `json:"status"` CreatedAt time.Time `json:"created_at"` ExpiresAt time.Time `json:"expires_at"` + // Corrupt reports that the archive failed a read-back: it cannot be + // restored. VerifiedAt is its last read-back that matched, and + // SkippedEntries the world entries it could not hold (symbolic links, + // devices, sockets). + Corrupt bool `json:"corrupt,omitempty"` + VerifiedAt *time.Time `json:"verified_at,omitempty"` + SkippedEntries int `json:"skipped_entries,omitempty"` } // BackupRecord is the server-side view of a backup used to drive a restore (spec @@ -81,6 +88,8 @@ type BackupRecord struct { FormerOwner string BackupRef string SizeBytes int64 + // Corrupt reports that the archive failed a read-back (see BackupView). + Corrupt bool } // StaffUser is the login-side projection of a users row (spec §B passwordless diff --git a/internal/backup/archiver.go b/internal/backup/archiver.go index 5a3ccb3..9b41ec9 100644 --- a/internal/backup/archiver.go +++ b/internal/backup/archiver.go @@ -4,7 +4,11 @@ // ArchiveRef is deliberately opaque. package backup -import "context" +import ( + "context" + "errors" + "time" +) // ArchiveRef is an opaque handle to a stored world archive. Depending on the // backend it may be a tar path, a VolumeSnapshot name, or a Longhorn backup URL. @@ -15,15 +19,64 @@ type ArchiveRef string // backends (VolumeSnapshot/Longhorn) cannot produce an io.Reader — they create // K8s objects referencing the source PVC (spec §19). type WorldArchiver interface { - // Archive captures the world living on pvc for server and returns an opaque - // ref plus the stored size in bytes. - Archive(ctx context.Context, server, pvc string) (ref ArchiveRef, size int64, err error) + // Archive captures the world living on pvc for server. It returns only once + // the archive is complete and durable: a failed or interrupted Archive leaves + // nothing that could be mistaken for a finished archive. + Archive(ctx context.Context, server, pvc string) (Archived, error) // Restore writes a previously archived world into targetPVC. Restore(ctx context.Context, ref ArchiveRef, targetPVC string) error // Delete removes the archive identified by ref. Delete(ctx context.Context, ref ArchiveRef) error } +// Archived is what one Archive call stored. +type Archived struct { + Ref ArchiveRef + Size int64 // stored bytes + // SHA256 is the hex digest of the stored archive, "" when the backend keeps + // none. world_backups.sha256 records it, so a later Verify tells an archive + // that rotted on disk from a good one. + SHA256 string + // Skipped lists the world entries the archive leaves out (symbolic links, + // devices, sockets, named pipes), relative to the world root. + Skipped []string +} + +// ErrCorrupt marks an archive that cannot be read back in full, or whose bytes +// no longer match the checksum recorded for it. It is the only Verify error that +// condemns the archive: any other (a mount that is not there, a cancelled +// context) says nothing about it. +var ErrCorrupt = errors.New("backup: archive corrupt") + +// Verifier is implemented by backends that can read a stored archive back end to +// end. Verify returns the archive's SHA256 as read; want, when not "", is the +// digest it must match (an archive recorded before checksums were kept has +// none, and its read-back establishes one). +type Verifier interface { + Verify(ctx context.Context, ref ArchiveRef, want string) (sha256 string, err error) +} + +// Sweeper is implemented by backends that can find what interrupted archives +// leave behind. +type Sweeper interface { + // Sweep removes every unfinished archive last written before partialBefore. + // A finished archive that live does not claim (its record was never inserted) + // is removed once it was last written before orphanBefore, and reported and + // kept until then. + Sweep(ctx context.Context, live func(ArchiveRef) bool, partialBefore, orphanBefore time.Time) (Swept, error) +} + +// Swept is what one Sweep found. +type Swept struct { + Removed []string // paths removed + // Orphans are finished archives no record claims that are still young enough + // to keep. A backup whose record insert failed looks like this, and so do the + // archives of a store brought back before the database that records them: + // removing them early could throw away the only copy of a world. + Orphans []string + OrphanBytes int64 +} + // PVCResolver maps a PVC name to the local filesystem path where it is mounted. // In production the reaper Job mounts the source/backup PVCs and supplies a // resolver over those mount points; tests supply temp dirs. diff --git a/internal/backup/tarlocal.go b/internal/backup/tarlocal.go index 928c988..d5614b3 100644 --- a/internal/backup/tarlocal.go +++ b/internal/backup/tarlocal.go @@ -4,13 +4,19 @@ import ( "archive/tar" "compress/gzip" "context" + "crypto/sha256" + "encoding/hex" + "errors" "fmt" "io" "io/fs" "os" "path" "path/filepath" + "regexp" + "sort" "strings" + "syscall" "time" ) @@ -34,38 +40,218 @@ func (t *TarLocal) now() time.Time { return time.Now() } +// partialSuffix marks an archive still being written. Its file is the final +// name hidden behind a leading dot, so a listing of the store shows finished +// archives only and a leftover still says whose it was. +const partialSuffix = ".partial" + +// archiveName matches the finished archives Archive names (-.tar.gz); partialName matches their unfinished form. Sweep touches +// nothing else in the store. +var ( + archiveName = regexp.MustCompile(`^[^.].*-[0-9]+\.tar\.gz$`) + partialName = regexp.MustCompile(`^\..*-[0-9]+\.tar\.gz\.partial$`) +) + // Archive tars+gzips the world on pvc into BackupRoot and returns the archive -// path as the opaque ref plus its on-disk size. -func (t *TarLocal) Archive(ctx context.Context, server, pvc string) (ArchiveRef, int64, error) { +// path as the opaque ref, with its size and SHA256. +// +// The archive is written under its hidden .partial name, flushed to disk, read +// back in full, and only then renamed to its final name and the rename flushed +// too. A Job killed at its deadline or a node that loses power mid-write leaves +// a .partial for Sweep, never a truncated file under a name that looks finished, +// and the world is deleted only after an archive that has already been read +// back whole. +func (t *TarLocal) Archive(ctx context.Context, server, pvc string) (Archived, error) { srcDir, err := t.Resolve(pvc) if err != nil { - return "", 0, err + return Archived{}, err } if err := os.MkdirAll(t.BackupRoot, 0o750); err != nil { - return "", 0, fmt.Errorf("backup: mkdir backup root: %w", err) + return Archived{}, fmt.Errorf("backup: mkdir backup root: %w", err) } name := fmt.Sprintf("%s-%d.tar.gz", server, t.now().UTC().UnixNano()) dest := filepath.Join(t.BackupRoot, name) + tmp := filepath.Join(t.BackupRoot, "."+name+partialSuffix) - f, err := os.Create(dest) + f, err := os.OpenFile(tmp, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o644) if err != nil { - return "", 0, fmt.Errorf("backup: create archive: %w", err) + return Archived{}, fmt.Errorf("backup: create archive: %w", err) } - if err := writeTarGz(ctx, f, srcDir); err != nil { - f.Close() - os.Remove(dest) - return "", 0, err + h := sha256.New() + st, err := writeTarGz(ctx, io.MultiWriter(f, h), srcDir) + if err == nil { + if err = f.Sync(); err != nil { + err = fmt.Errorf("backup: sync archive: %w", err) + } } - if err := f.Close(); err != nil { - os.Remove(dest) - return "", 0, fmt.Errorf("backup: close archive: %w", err) + if cerr := f.Close(); err == nil && cerr != nil { + err = fmt.Errorf("backup: close archive: %w", cerr) } + if err != nil { + os.Remove(tmp) + return Archived{}, err + } + sum := hex.EncodeToString(h.Sum(nil)) + // Read back what landed: the gzip CRC, every tar header and entry, and the + // entry count must all come out as written. + entries, _, err := verifyArchive(ctx, tmp, sum) + if err == nil && entries != st.entries { + err = fmt.Errorf("%w: %s: %d entries read back, %d written", ErrCorrupt, tmp, entries, st.entries) + } + if err != nil { + os.Remove(tmp) + return Archived{}, err + } + if err := os.Rename(tmp, dest); err != nil { + os.Remove(tmp) + return Archived{}, fmt.Errorf("backup: name archive: %w", err) + } + if err := syncDir(t.BackupRoot); err != nil { + os.Remove(dest) + return Archived{}, fmt.Errorf("backup: sync backup root: %w", err) + } info, err := os.Stat(dest) if err != nil { - return "", 0, fmt.Errorf("backup: stat archive: %w", err) + return Archived{}, fmt.Errorf("backup: stat archive: %w", err) } - return ArchiveRef(dest), info.Size(), nil + return Archived{Ref: ArchiveRef(dest), Size: info.Size(), SHA256: sum, Skipped: st.skipped}, nil +} + +// Verify reads the archive at ref back end to end (see verifyArchive). +func (t *TarLocal) Verify(ctx context.Context, ref ArchiveRef, want string) (string, error) { + _, sum, err := verifyArchive(ctx, string(ref), want) + return sum, err +} + +// verifyArchive reads the archive at p through gzip and tar to the last byte and +// returns its entry count and SHA256. Anything that does not read back — a +// missing file, a torn gzip stream, a CRC or length mismatch, a bad tar header, +// or a digest other than want (when want is not "") — is ErrCorrupt. Only an +// error opening a file that is there, or the context, is not. +func verifyArchive(ctx context.Context, p, want string) (int, string, error) { + f, err := os.Open(p) + if errors.Is(err, fs.ErrNotExist) { + return 0, "", fmt.Errorf("%w: %s is missing", ErrCorrupt, p) + } + if err != nil { + return 0, "", fmt.Errorf("backup: open archive: %w", err) + } + defer f.Close() + h := sha256.New() + raw := io.TeeReader(f, h) + corrupt := func(err error) (int, string, error) { + if ctx.Err() != nil { + return 0, "", ctx.Err() + } + return 0, "", fmt.Errorf("%w: %s: %v", ErrCorrupt, p, err) + } + + gz, err := gzip.NewReader(raw) + if err != nil { + return corrupt(err) + } + tr := tar.NewReader(gz) + entries := 0 + for { + if ctx.Err() != nil { + return 0, "", ctx.Err() + } + _, err := tr.Next() + if err == io.EOF { + break + } + if err != nil { + return corrupt(err) + } + if _, err := io.Copy(io.Discard, tr); err != nil { + return corrupt(err) + } + entries++ + } + // The tar end marker is not the end of the gzip member: reading on to EOF is + // what checks the CRC and length in the gzip trailer. + if _, err := io.Copy(io.Discard, gz); err != nil { + return corrupt(err) + } + if err := gz.Close(); err != nil { + return corrupt(err) + } + if _, err := io.Copy(io.Discard, raw); err != nil { + return corrupt(err) + } + sum := hex.EncodeToString(h.Sum(nil)) + if want != "" && sum != want { + return 0, sum, fmt.Errorf("%w: %s: sha256 %s, recorded %s", ErrCorrupt, p, sum, want) + } + return entries, sum, nil +} + +// Sweep removes the leftovers of archives that never finished (see Sweeper). It +// only looks at the top level of BackupRoot and only at names Archive writes; +// partialBefore must leave room for the longest Archive, and orphanBefore for +// the insert of the record that follows it at the very least. +func (t *TarLocal) Sweep(ctx context.Context, live func(ArchiveRef) bool, partialBefore, orphanBefore time.Time) (Swept, error) { + var out Swept + ents, err := os.ReadDir(t.BackupRoot) + if errors.Is(err, fs.ErrNotExist) { + return out, nil + } + if err != nil { + return out, fmt.Errorf("backup: list backup root: %w", err) + } + var errs []error + for _, e := range ents { + if ctx.Err() != nil { + return out, ctx.Err() + } + name := e.Name() + if !e.Type().IsRegular() { + continue + } + p := filepath.Join(t.BackupRoot, name) + cutoff := partialBefore + switch { + case partialName.MatchString(name): + case archiveName.MatchString(name): + if live(ArchiveRef(p)) { + continue + } + cutoff = orphanBefore + default: + continue + } + info, err := e.Info() + if err != nil || !info.ModTime().Before(partialBefore) { + continue + } + if !info.ModTime().Before(cutoff) { + out.Orphans = append(out.Orphans, p) + out.OrphanBytes += info.Size() + continue + } + if err := os.Remove(p); err != nil && !errors.Is(err, fs.ErrNotExist) { + errs = append(errs, err) + continue + } + out.Removed = append(out.Removed, p) + } + return out, errors.Join(errs...) +} + +// syncDir flushes a change to a directory's entries (a rename) to disk. A +// filesystem that cannot sync a directory has no stronger promise to give. +func syncDir(dir string) error { + d, err := os.Open(dir) + if err != nil { + return err + } + defer d.Close() + if err := d.Sync(); err != nil && !errors.Is(err, syscall.EINVAL) && !errors.Is(err, syscall.ENOTSUP) { + return err + } + return nil } // Restore extracts the archive at ref into the world mount for targetPVC, @@ -74,10 +260,16 @@ func (t *TarLocal) Archive(ctx context.Context, server, pvc string) (ArchiveRef, // griefer's chunks). It extracts over the target, then removes any pre-existing // entry the archive did not contain. // -// The prune runs only after a fully successful extract: a corrupt or truncated -// archive fails before the prune, leaving the target as a (recoverable) partial -// overlay rather than a destroyed world. The archive is retained on restore, so -// such a failure is recoverable by re-running the Job. +// The archive is read back in full before anything is extracted, so a corrupt +// or truncated archive fails with the world untouched. The prune runs only after +// a fully successful extract: an extract that still fails (the volume fills up) +// leaves the target as a recoverable partial overlay rather than a destroyed +// world. The archive is retained on restore, so such a failure is recoverable by +// re-running the Job. +// +// Files and directories get back the permission bits and modification times the +// archive recorded. Ownership is left to the server: every start re-owns the +// world volume to the game uid (felis init-volume). // // A top-level lost+found is never a prune target. It is a filesystem artifact // (root-owned, mode 0700) that the non-root restore Pod cannot delete anyway, @@ -90,17 +282,23 @@ func (t *TarLocal) Restore(ctx context.Context, ref ArchiveRef, targetPVC string if err := os.MkdirAll(dstDir, 0o750); err != nil { return fmt.Errorf("backup: mkdir restore target: %w", err) } + if _, _, err := verifyArchive(ctx, string(ref), ""); err != nil { + return err + } f, err := os.Open(string(ref)) if err != nil { return fmt.Errorf("backup: open archive: %w", err) } defer f.Close() - keep, err := readTarGz(ctx, f, dstDir) + keep, dirs, err := readTarGz(ctx, f, dstDir) if err != nil { return err } - return pruneToManifest(dstDir, keep) + if err := pruneToManifest(dstDir, keep); err != nil { + return err + } + return restoreDirMeta(dirs) } // pruneToManifest removes every entry under dstDir whose archive-relative path @@ -151,7 +349,18 @@ func (t *TarLocal) Delete(_ context.Context, ref ArchiveRef) error { return nil } -func writeTarGz(ctx context.Context, w io.Writer, srcDir string) error { +// tarStats is what writeTarGz put in the archive and what it left out. +type tarStats struct { + entries int + skipped []string +} + +// writeTarGz archives srcDir (not the root entry itself) as gzip+tar. Each entry +// keeps its permission bits, without setuid, setgid and sticky, and its +// modification time. Entries other than regular files and directories are left +// out and listed in the stats. +func writeTarGz(ctx context.Context, w io.Writer, srcDir string) (tarStats, error) { + var st tarStats gz := gzip.NewWriter(w) tw := tar.NewWriter(gz) @@ -173,12 +382,17 @@ func writeTarGz(ctx context.Context, w io.Writer, srcDir string) error { // Normalize to forward slashes so archives are portable. name := filepath.ToSlash(rel) + mode := int64(info.Mode().Perm()) switch { case info.IsDir(): - hdr := &tar.Header{Name: name + "/", Mode: 0o750, Typeflag: tar.TypeDir} - return tw.WriteHeader(hdr) + hdr := &tar.Header{Name: name + "/", Mode: mode, ModTime: info.ModTime(), Typeflag: tar.TypeDir} + if err := tw.WriteHeader(hdr); err != nil { + return err + } + st.entries++ + return nil case info.Mode().IsRegular(): - hdr := &tar.Header{Name: name, Mode: 0o640, Size: info.Size(), Typeflag: tar.TypeReg} + hdr := &tar.Header{Name: name, Mode: mode, ModTime: info.ModTime(), Size: info.Size(), Typeflag: tar.TypeReg} if err := tw.WriteHeader(hdr); err != nil { return err } @@ -187,24 +401,37 @@ func writeTarGz(ctx context.Context, w io.Writer, srcDir string) error { return err } defer src.Close() - _, err = io.Copy(tw, src) - return err + if _, err := io.Copy(tw, src); err != nil { + return err + } + st.entries++ + return nil default: - // Skip symlinks/devices/sockets: a world directory should be plain - // files, and refusing the rest avoids surprising archive contents. + // A symlink could point anywhere on the node, and a device or socket + // has no bytes to keep: a world is plain files. What is left out is + // reported, so a backup never drops something silently. + st.skipped = append(st.skipped, name) return nil } }) if err != nil { - return fmt.Errorf("backup: tar walk: %w", err) + return st, fmt.Errorf("backup: tar walk: %w", err) } if err := tw.Close(); err != nil { - return fmt.Errorf("backup: close tar: %w", err) + return st, fmt.Errorf("backup: close tar: %w", err) } if err := gz.Close(); err != nil { - return fmt.Errorf("backup: close gzip: %w", err) + return st, fmt.Errorf("backup: close gzip: %w", err) } - return nil + return st, nil +} + +// dirMeta is a directory's recorded permission bits and modification time, +// applied once nothing more is written into it. +type dirMeta struct { + path string + mode fs.FileMode + mtime time.Time } // readTarGz extracts the gzip+tar stream into dstDir and returns the keep-set: @@ -213,33 +440,39 @@ func writeTarGz(ctx context.Context, w io.Writer, srcDir string) error { // target files for replace semantics. On any error the keep-set is incomplete // and must not be used to prune (a partial manifest would delete live files the // stream had not yet reached). -func readTarGz(ctx context.Context, r io.Reader, dstDir string) (map[string]struct{}, error) { +// +// Files get their recorded mode and modification time as they are written. A +// directory's are returned instead (restoreDirMeta): extracting and pruning +// inside it would move its mtime again, and a read-only mode would stop the +// extract from writing into it. +func readTarGz(ctx context.Context, r io.Reader, dstDir string) (map[string]struct{}, []dirMeta, error) { gz, err := gzip.NewReader(r) if err != nil { - return nil, fmt.Errorf("backup: open gzip: %w", err) + return nil, nil, fmt.Errorf("backup: open gzip: %w", err) } defer gz.Close() tr := tar.NewReader(gz) keep := make(map[string]struct{}) + var dirs []dirMeta cleanDst := filepath.Clean(dstDir) for { if ctx.Err() != nil { - return nil, ctx.Err() + return nil, nil, ctx.Err() } hdr, err := tr.Next() if err == io.EOF { - return keep, nil + return keep, dirs, nil } if err != nil { - return nil, fmt.Errorf("backup: read tar: %w", err) + return nil, nil, fmt.Errorf("backup: read tar: %w", err) } // Guard against path traversal (zip-slip): the resolved target must stay // within dstDir. target := filepath.Join(cleanDst, filepath.FromSlash(hdr.Name)) if target != cleanDst && !strings.HasPrefix(target, cleanDst+string(os.PathSeparator)) { - return nil, fmt.Errorf("backup: archive entry escapes target: %q", hdr.Name) + return nil, nil, fmt.Errorf("backup: archive entry escapes target: %q", hdr.Name) } // Record this entry and its ancestors in the keep-set. Names are stored @@ -247,25 +480,39 @@ func readTarGz(ctx context.Context, r io.Reader, dstDir string) (map[string]stru // relative paths pruneToManifest derives from the on-disk walk. rememberKept(keep, hdr.Name) + mode := fs.FileMode(hdr.Mode).Perm() switch hdr.Typeflag { case tar.TypeDir: if err := os.MkdirAll(target, 0o750); err != nil { - return nil, err + return nil, nil, err + } + if target != cleanDst { + dirs = append(dirs, dirMeta{path: target, mode: mode, mtime: hdr.ModTime}) } case tar.TypeReg: if err := os.MkdirAll(filepath.Dir(target), 0o750); err != nil { - return nil, err + return nil, nil, err } - out, err := os.OpenFile(target, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o640) + out, err := os.OpenFile(target, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600) if err != nil { - return nil, err + return nil, nil, err } if _, err := io.Copy(out, tr); err != nil { out.Close() - return nil, err + return nil, nil, err } if err := out.Close(); err != nil { - return nil, err + return nil, nil, err + } + // Chmod rather than the create mode: the umask would narrow it, and + // a file that already existed keeps its old mode through O_TRUNC. + if err := os.Chmod(target, mode); err != nil { + return nil, nil, err + } + if recordedTime(hdr.ModTime) { + if err := os.Chtimes(target, hdr.ModTime, hdr.ModTime); err != nil { + return nil, nil, err + } } default: // Ignore entry types tarLocal never writes. @@ -273,6 +520,28 @@ func readTarGz(ctx context.Context, r io.Reader, dstDir string) (map[string]stru } } +// restoreDirMeta applies the directories' recorded modes and times, deepest +// first, so setting a parent's time comes after every change inside it. +func restoreDirMeta(dirs []dirMeta) error { + sort.SliceStable(dirs, func(i, j int) bool { return len(dirs[i].path) > len(dirs[j].path) }) + for _, d := range dirs { + if err := os.Chmod(d.path, d.mode); err != nil { + return fmt.Errorf("backup: restore mode of %s: %w", d.path, err) + } + if recordedTime(d.mtime) { + if err := os.Chtimes(d.path, d.mtime, d.mtime); err != nil { + return fmt.Errorf("backup: restore time of %s: %w", d.path, err) + } + } + } + return nil +} + +// recordedTime reports whether an archive recorded a modification time. Archives +// written before times were kept carry the Unix epoch, which is left alone +// rather than stamped onto every file. +func recordedTime(t time.Time) bool { return t.Unix() > 0 } + // rememberKept adds an archive entry name and every ancestor directory to keep, // normalized to a cleaned forward-slash path with no trailing slash. Adding // ancestors guards against archives that list a file without an explicit entry diff --git a/internal/backup/tarlocal_test.go b/internal/backup/tarlocal_test.go index 6c9bfc3..b9d65c6 100644 --- a/internal/backup/tarlocal_test.go +++ b/internal/backup/tarlocal_test.go @@ -4,9 +4,15 @@ import ( "archive/tar" "compress/gzip" "context" + "crypto/sha256" + "encoding/hex" + "errors" "os" "path/filepath" + "reflect" + "sort" "testing" + "time" "felis.lolicon.best/internal/backup" ) @@ -47,10 +53,11 @@ func TestTarLocalRoundTrip(t *testing.T) { } ctx := context.Background() - ref, size, err := archiver.Archive(ctx, "survival", "src-pvc") + a, err := archiver.Archive(ctx, "survival", "src-pvc") if err != nil { t.Fatalf("Archive: %v", err) } + ref, size := a.Ref, a.Size if size <= 0 { t.Errorf("archive size = %d, want > 0", size) } @@ -139,10 +146,11 @@ func TestTarLocalRestoreReplacesTarget(t *testing.T) { } ctx := context.Background() - ref, _, err := archiver.Archive(ctx, "survival", "src-pvc") + a, err := archiver.Archive(ctx, "survival", "src-pvc") if err != nil { t.Fatalf("Archive: %v", err) } + ref := a.Ref if err := archiver.Restore(ctx, ref, "dst-pvc"); err != nil { t.Fatalf("Restore: %v", err) } @@ -181,7 +189,7 @@ func TestTarLocalUnknownPVC(t *testing.T) { BackupRoot: t.TempDir(), Resolve: backup.StaticResolver(map[string]string{}), } - if _, _, err := archiver.Archive(context.Background(), "x", "missing"); err == nil { + if _, err := archiver.Archive(context.Background(), "x", "missing"); err == nil { t.Fatal("expected error for unknown pvc") } } @@ -222,3 +230,268 @@ func TestTarLocalRejectsZipSlip(t *testing.T) { t.Errorf("zip-slip wrote outside target: %v", err) } } + +// TestTarLocalArchiveIsDurableAndChecksummed pins the write path: the archive +// lands under its final name only, no .partial is left beside it, and the +// SHA256 it reports is the one of the bytes on disk, which Verify reads back. +func TestTarLocalArchiveIsDurableAndChecksummed(t *testing.T) { + src, backupRoot := t.TempDir(), t.TempDir() + writeTree(t, src, map[string]string{"level.dat": "seed", "region/r.0.0.mca": "chunk"}) + archiver := &backup.TarLocal{BackupRoot: backupRoot, Resolve: backup.StaticResolver(map[string]string{"src": src})} + ctx := context.Background() + + a, err := archiver.Archive(ctx, "survival", "src") + if err != nil { + t.Fatalf("Archive: %v", err) + } + ents, err := os.ReadDir(backupRoot) + if err != nil { + t.Fatal(err) + } + if len(ents) != 1 || filepath.Join(backupRoot, ents[0].Name()) != string(a.Ref) { + t.Fatalf("backup root holds %v, want only %s", ents, a.Ref) + } + body, err := os.ReadFile(string(a.Ref)) + if err != nil { + t.Fatal(err) + } + sum := sha256.Sum256(body) + if a.SHA256 != hex.EncodeToString(sum[:]) || a.Size != int64(len(body)) { + t.Errorf("Archived = %+v, want sha256 %x size %d", a, sum, len(body)) + } + got, err := archiver.Verify(ctx, a.Ref, a.SHA256) + if err != nil || got != a.SHA256 { + t.Errorf("Verify = %q, %v; want %q, nil", got, err, a.SHA256) + } + // With no digest recorded, Verify still reads it through and reports one. + if got, err := archiver.Verify(ctx, a.Ref, ""); err != nil || got != a.SHA256 { + t.Errorf("Verify without a digest = %q, %v", got, err) + } +} + +// TestTarLocalFailedArchiveLeavesNothing: an Archive that fails part way (here +// a cancelled context) removes its .partial, so no leftover can pass for an +// archive. +func TestTarLocalFailedArchiveLeavesNothing(t *testing.T) { + src, backupRoot := t.TempDir(), t.TempDir() + writeTree(t, src, map[string]string{"level.dat": "seed"}) + archiver := &backup.TarLocal{BackupRoot: backupRoot, Resolve: backup.StaticResolver(map[string]string{"src": src})} + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if _, err := archiver.Archive(ctx, "survival", "src"); err == nil { + t.Fatal("Archive with a cancelled context succeeded") + } + if ents, _ := os.ReadDir(backupRoot); len(ents) != 0 { + t.Errorf("failed Archive left %v behind", ents) + } +} + +// TestTarLocalKeepsModesAndTimes: a restore gives files and directories back +// their permission bits (setuid, setgid and sticky dropped) and modification +// times, and what the archive cannot hold is reported rather than dropped +// silently. +func TestTarLocalKeepsModesAndTimes(t *testing.T) { + src, dst, backupRoot := t.TempDir(), t.TempDir(), t.TempDir() + writeTree(t, src, map[string]string{ + "start.sh": "#!/bin/sh", + "secret.properties": "rcon", + "private/notes.txt": "n", + "setuid-bin": "x", + }) + mtime := time.Date(2025, 3, 4, 5, 6, 7, 0, time.UTC) + modes := map[string]os.FileMode{ + "start.sh": 0o755, + "secret.properties": 0o600, + "private/notes.txt": 0o640, + "setuid-bin": 0o755 | os.ModeSetuid, + "private": 0o700, + } + for rel, m := range modes { + p := filepath.Join(src, filepath.FromSlash(rel)) + if err := os.Chmod(p, m); err != nil { + t.Fatal(err) + } + } + for _, rel := range []string{"start.sh", "secret.properties", "private/notes.txt", "setuid-bin", "private"} { + if err := os.Chtimes(filepath.Join(src, filepath.FromSlash(rel)), mtime, mtime); err != nil { + t.Fatal(err) + } + } + if err := os.Symlink("/etc/passwd", filepath.Join(src, "link")); err != nil { + t.Fatal(err) + } + + archiver := &backup.TarLocal{BackupRoot: backupRoot, Resolve: backup.StaticResolver(map[string]string{"src": src, "dst": dst})} + ctx := context.Background() + a, err := archiver.Archive(ctx, "survival", "src") + if err != nil { + t.Fatalf("Archive: %v", err) + } + if !reflect.DeepEqual(a.Skipped, []string{"link"}) { + t.Errorf("Skipped = %v, want [link]", a.Skipped) + } + if err := archiver.Restore(ctx, a.Ref, "dst"); err != nil { + t.Fatalf("Restore: %v", err) + } + for rel, m := range modes { + info, err := os.Lstat(filepath.Join(dst, filepath.FromSlash(rel))) + if err != nil { + t.Errorf("%s: %v", rel, err) + continue + } + if got, want := info.Mode().Perm(), m.Perm(); got != want || info.Mode()&os.ModeSetuid != 0 { + t.Errorf("%s restored with mode %v, want %v", rel, info.Mode(), want) + } + if !info.ModTime().Equal(mtime) { + t.Errorf("%s restored with mtime %v, want %v", rel, info.ModTime(), mtime) + } + } + if _, err := os.Lstat(filepath.Join(dst, "link")); !os.IsNotExist(err) { + t.Errorf("symlink came back from the archive: %v", err) + } +} + +// damage rewrites the archive at ref with f applied to its bytes. +func damage(t *testing.T, ref backup.ArchiveRef, f func([]byte) []byte) { + t.Helper() + b, err := os.ReadFile(string(ref)) + if err != nil { + t.Fatal(err) + } + if err := os.WriteFile(string(ref), f(b), 0o644); err != nil { + t.Fatal(err) + } +} + +// TestTarLocalVerifyFindsCorruption: a flipped byte, a truncated file, a digest +// that no longer matches and a missing file all come back as ErrCorrupt. +func TestTarLocalVerifyFindsCorruption(t *testing.T) { + src := t.TempDir() + writeTree(t, src, map[string]string{"level.dat": "seed", "region/r.0.0.mca": string(make([]byte, 64<<10))}) + ctx := context.Background() + fresh := func(t *testing.T) (*backup.TarLocal, backup.Archived) { + archiver := &backup.TarLocal{BackupRoot: t.TempDir(), Resolve: backup.StaticResolver(map[string]string{"src": src})} + a, err := archiver.Archive(ctx, "survival", "src") + if err != nil { + t.Fatalf("Archive: %v", err) + } + return archiver, a + } + cases := map[string]struct { + change func(*testing.T, backup.Archived) + want func(backup.Archived) string + }{ + "flipped byte": { + change: func(t *testing.T, a backup.Archived) { + damage(t, a.Ref, func(b []byte) []byte { b[len(b)/2] ^= 0xff; return b }) + }, + want: func(a backup.Archived) string { return "" }, + }, + "truncated": { + change: func(t *testing.T, a backup.Archived) { + damage(t, a.Ref, func(b []byte) []byte { return b[:len(b)-9] }) + }, + want: func(a backup.Archived) string { return "" }, + }, + "digest mismatch": { + change: func(*testing.T, backup.Archived) {}, + want: func(backup.Archived) string { return hex.EncodeToString(make([]byte, 32)) }, + }, + "missing": { + change: func(t *testing.T, a backup.Archived) { + if err := os.Remove(string(a.Ref)); err != nil { + t.Fatal(err) + } + }, + want: func(a backup.Archived) string { return a.SHA256 }, + }, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + archiver, a := fresh(t) + tc.change(t, a) + if _, err := archiver.Verify(ctx, a.Ref, tc.want(a)); !errors.Is(err, backup.ErrCorrupt) { + t.Errorf("Verify = %v, want ErrCorrupt", err) + } + }) + } +} + +// TestTarLocalRestoreRefusesCorruptArchive: the archive is read back before +// anything is extracted, so a corrupt one leaves the world exactly as it was. +func TestTarLocalRestoreRefusesCorruptArchive(t *testing.T) { + src, dst := t.TempDir(), t.TempDir() + writeTree(t, src, map[string]string{"level.dat": "archived", "region/r.0.0.mca": string(make([]byte, 64<<10))}) + writeTree(t, dst, map[string]string{"level.dat": "live"}) + archiver := &backup.TarLocal{BackupRoot: t.TempDir(), Resolve: backup.StaticResolver(map[string]string{"src": src, "dst": dst})} + ctx := context.Background() + a, err := archiver.Archive(ctx, "survival", "src") + if err != nil { + t.Fatalf("Archive: %v", err) + } + damage(t, a.Ref, func(b []byte) []byte { return b[:len(b)-9] }) + if err := archiver.Restore(ctx, a.Ref, "dst"); !errors.Is(err, backup.ErrCorrupt) { + t.Fatalf("Restore = %v, want ErrCorrupt", err) + } + got, err := os.ReadFile(filepath.Join(dst, "level.dat")) + if err != nil || string(got) != "live" { + t.Errorf("world changed by a refused restore: %q, %v", got, err) + } + if _, err := os.Stat(filepath.Join(dst, "region")); !os.IsNotExist(err) { + t.Errorf("refused restore extracted entries: %v", err) + } +} + +// TestTarLocalSweep removes old .partial files and orphaned archives past the +// orphan cutoff, reports younger orphans, and leaves everything else: fresh +// leftovers (an Archive may still be writing, or its record not yet inserted), +// claimed archives, and files it did not write. +func TestTarLocalSweep(t *testing.T) { + root := t.TempDir() + now := time.Now() + ancient, old, recent := now.Add(-100*24*time.Hour), now.Add(-7*time.Hour), now.Add(-time.Minute) + files := map[string]time.Time{ + ".survival-100.tar.gz.partial": old, + ".survival-200.tar.gz.partial": recent, + "survival-300.tar.gz": ancient, // orphan past the cutoff + "survival-310.tar.gz": old, // orphan, kept and reported + "survival-400.tar.gz": ancient, // claimed + "survival-500.tar.gz": recent, // may be about to be claimed + "notes.txt": ancient, + "felis.dump": ancient, + } + for name, at := range files { + p := filepath.Join(root, name) + if err := os.WriteFile(p, []byte("xyz"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Chtimes(p, at, at); err != nil { + t.Fatal(err) + } + } + archiver := &backup.TarLocal{BackupRoot: root} + live := func(ref backup.ArchiveRef) bool { + return ref == backup.ArchiveRef(filepath.Join(root, "survival-400.tar.gz")) + } + got, err := archiver.Sweep(context.Background(), live, now.Add(-6*time.Hour), now.Add(-90*24*time.Hour)) + if err != nil { + t.Fatalf("Sweep: %v", err) + } + sort.Strings(got.Removed) + want := backup.Swept{ + Removed: []string{filepath.Join(root, ".survival-100.tar.gz.partial"), filepath.Join(root, "survival-300.tar.gz")}, + Orphans: []string{filepath.Join(root, "survival-310.tar.gz")}, + OrphanBytes: 3, + } + if !reflect.DeepEqual(got, want) { + t.Errorf("Sweep = %+v, want %+v", got, want) + } + ents, _ := os.ReadDir(root) + if len(ents) != len(files)-2 { + t.Errorf("%d files left, want %d", len(ents), len(files)-2) + } + // A store that does not exist yet has nothing to sweep. + if _, err := (&backup.TarLocal{BackupRoot: filepath.Join(root, "absent")}).Sweep(context.Background(), live, now, now); err != nil { + t.Errorf("Sweep of a missing root: %v", err) + } +} diff --git a/internal/pgint/reaper_test.go b/internal/pgint/reaper_test.go index 74ab594..88e3db8 100644 --- a/internal/pgint/reaper_test.go +++ b/internal/pgint/reaper_test.go @@ -37,10 +37,10 @@ func (c *reclaimCluster) Stop(context.Context, string) error { return nil } type reclaimArchiver struct{ archived []string } -func (a *reclaimArchiver) Archive(_ context.Context, server, _ string) (backup.ArchiveRef, int64, error) { +func (a *reclaimArchiver) Archive(_ context.Context, server, _ string) (backup.Archived, error) { ref := "/archives/" + server + "-new.tar.gz" a.archived = append(a.archived, ref) - return backup.ArchiveRef(ref), 42, nil + return backup.Archived{Ref: backup.ArchiveRef(ref), Size: 42}, nil } func (a *reclaimArchiver) Restore(context.Context, backup.ArchiveRef, string) error { return nil } @@ -244,3 +244,146 @@ func TestManualBackupRationing(t *testing.T) { t.Fatalf("a request before since still counted: (%v, %v)", at, err) } } + +// TestBackupReadBack pins the SQL behind data-durability-7/15/16: the digest a +// backup is written with survives a later read-back, a corrupt archive is no +// longer offered for reuse, restore or read-back but stays listed, and it is +// the first a keep-N prune removes. +func TestBackupReadBack(t *testing.T) { + ctx := context.Background() + st := reaper.NewPGStore(db) + name := "readback-" + 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) + } + now := time.Now() + sfx := suffix(t) + older, newer := "bk-old-"+sfx, "bk-new-"+sfx + for _, r := range []reaper.BackupRecord{ + {ID: older, ServerName: name, BackupRef: "/archives/" + older + ".tar.gz", SizeBytes: 10, + Reason: "inactive_15d", ExpiresAt: now.Add(90 * reaper.Day), SHA256: "aa", SkippedEntries: 2}, + {ID: newer, ServerName: name, BackupRef: "/archives/" + newer + ".tar.gz", SizeBytes: 10, + Reason: "inactive_15d", ExpiresAt: now.Add(90 * reaper.Day)}, + } { + if err := st.InsertBackup(ctx, r); err != nil { + t.Fatalf("InsertBackup %s: %v", r.ID, err) + } + } + if _, err := db.ExecContext(ctx, + `UPDATE world_backups SET created_at = CASE id WHEN $1 THEN $3::timestamptz ELSE $4::timestamptz END WHERE id IN ($1, $2)`, + older, newer, now.Add(-2*time.Hour), now.Add(-time.Hour)); err != nil { + t.Fatalf("age backups: %v", err) + } + + if f, ok, err := st.FreshBackup(ctx, name, now.Add(-reaper.Day)); err != nil || !ok || f.ID != newer || f.SHA256 != "" { + t.Fatalf("FreshBackup = (%+v, %v, %v); want %s with no digest yet", f, ok, err, newer) + } + + // A read-back fills in a missing digest and never rewrites a recorded one. + checked := now.Add(-time.Minute) + for id, sum := range map[string]string{newer: "bb", older: "zz"} { + if err := st.MarkBackupVerified(ctx, id, sum, checked); err != nil { + t.Fatalf("MarkBackupVerified %s: %v", id, err) + } + } + due := func(before time.Time) map[string]string { + t.Helper() + bs, err := st.BackupsToVerify(ctx, before, 100000) + if err != nil { + t.Fatalf("BackupsToVerify: %v", err) + } + out := map[string]string{} + for _, b := range bs { + if b.ServerName == name { + out[b.ID] = b.SHA256 + } + } + return out + } + if got := due(checked.Add(-time.Second)); len(got) != 0 { + t.Fatalf("due before their read-back = %v; want none", got) + } + if got := due(checked.Add(time.Second)); got[newer] != "bb" || got[older] != "aa" || len(got) != 2 { + t.Fatalf("due after their read-back = %v; want %s=bb and %s=aa", got, newer, older) + } + + corruptAt := now.Add(-30 * time.Second) + if err := st.MarkBackupCorrupt(ctx, newer, corruptAt); err != nil { + t.Fatalf("MarkBackupCorrupt: %v", err) + } + if err := st.MarkBackupCorrupt(ctx, newer, now); err != nil { + t.Fatalf("MarkBackupCorrupt again: %v", err) + } + var first time.Time + if err := db.QueryRowContext(ctx, `SELECT corrupt_at FROM world_backups WHERE id = $1`, newer).Scan(&first); err != nil || + !first.Equal(corruptAt.Truncate(time.Microsecond)) { + t.Fatalf("corrupt_at = (%v, %v); want the first finding %v kept", first, err, corruptAt) + } + + if f, ok, err := st.FreshBackup(ctx, name, now.Add(-reaper.Day)); err != nil || !ok || f.ID != older || f.SHA256 != "aa" { + t.Fatalf("FreshBackup after corruption = (%+v, %v, %v); want the intact %s", f, ok, err, older) + } + if got := due(now.Add(reaper.Day)); len(got) != 1 || got[older] == "" { + t.Fatalf("due after corruption = %v; want only %s", got, older) + } + if b, err := repo.LatestBackup(ctx, name); err != nil || b.ID != older { + t.Fatalf("LatestBackup = (%+v, %v); want %s", b, err, older) + } + if b, err := repo.BackupByID(ctx, newer); err != nil || !b.Corrupt { + t.Fatalf("BackupByID(corrupt) = (%+v, %v); want Corrupt", b, err) + } + if b, err := repo.BackupByID(ctx, older); err != nil || b.Corrupt { + t.Fatalf("BackupByID(intact) = (%+v, %v); want not Corrupt", b, err) + } + + views, err := repo.AllBackups(ctx) + if err != nil { + t.Fatalf("AllBackups: %v", err) + } + seen := 0 + for _, v := range views { + switch v.ID { + case newer: + seen++ + if !v.Corrupt || v.VerifiedAt == nil || v.SkippedEntries != 0 { + t.Fatalf("listed corrupt backup = %+v; want corrupt, verified_at set", v) + } + case older: + seen++ + if v.Corrupt || v.VerifiedAt == nil || v.SkippedEntries != 2 { + t.Fatalf("listed intact backup = %+v; want verified_at set and 2 skipped entries", v) + } + } + } + if seen != 2 { + t.Fatalf("AllBackups listed %d of the 2 backups", seen) + } + + if excess, err := st.ExcessBackups(ctx, name, "inactive_15d", 1, ""); err != nil || len(excess) != 1 || excess[0].ID != newer { + t.Fatalf("ExcessBackups(keep 1) = (%+v, %v); want the corrupt %s pruned first", excess, err, newer) + } + + live := func() map[string]bool { + t.Helper() + refs, err := st.LiveBackupRefs(ctx) + if err != nil { + t.Fatalf("LiveBackupRefs: %v", err) + } + out := map[string]bool{} + for _, r := range refs { + out[r] = true + } + return out + } + if l := live(); !l["/archives/"+older+".tar.gz"] || !l["/archives/"+newer+".tar.gz"] { + t.Fatalf("live refs miss a present backup") + } + if err := st.MarkBackupDeleted(ctx, newer, now); err != nil { + t.Fatalf("MarkBackupDeleted: %v", err) + } + if l := live(); l["/archives/"+newer+".tar.gz"] || !l["/archives/"+older+".tar.gz"] { + t.Fatalf("live refs after deleting %s still claim its archive", newer) + } +} diff --git a/internal/reaper/pgstore.go b/internal/reaper/pgstore.go index 90683de..64c38b3 100644 --- a/internal/reaper/pgstore.go +++ b/internal/reaper/pgstore.go @@ -53,13 +53,14 @@ func (s *PGStore) FreshBackup(ctx context.Context, server string, since time.Tim // Only the reaper's own archives count, and only those taken since the // current owner claimed the server: a manual backup may predate a panel edit // that did not move last_active_at, and an archive from before the claim is - // the previous owner's world. - const q = `SELECT b.backup_ref, b.offsite_at IS NOT NULL FROM world_backups b + // the previous owner's world. One found corrupt is never reused. + const q = `SELECT b.id, b.backup_ref, COALESCE(b.sha256, ''), b.offsite_at IS NOT NULL FROM world_backups b WHERE b.server_name = $1 AND b.status = 'present' AND b.reason = 'inactive_15d' AND b.created_at >= $2 + AND b.corrupt_at IS NULL AND b.created_at >= COALESCE((SELECT s.claimed_at FROM servers s WHERE s.name = $1 AND s.deleted_at IS NULL), '-infinity') ORDER BY b.offsite_at IS NOT NULL DESC, b.created_at DESC LIMIT 1` var f Fresh - switch err := s.db.QueryRowContext(ctx, q, server, since).Scan(&f.Ref, &f.Offsite); { + switch err := s.db.QueryRowContext(ctx, q, server, since).Scan(&f.ID, &f.Ref, &f.SHA256, &f.Offsite); { case err == sql.ErrNoRows: return Fresh{}, false, nil case err != nil: @@ -70,10 +71,11 @@ func (s *PGStore) FreshBackup(ctx context.Context, server string, since time.Tim func (s *PGStore) InsertBackup(ctx context.Context, rec BackupRecord) error { const q = `INSERT INTO world_backups - (id, server_name, former_owner, backup_ref, size_bytes, reason, status, created_at, expires_at) - VALUES ($1, $2, NULLIF($3, ''), $4, $5, $6, 'present', now(), $7)` + (id, server_name, former_owner, backup_ref, size_bytes, reason, status, created_at, expires_at, sha256, skipped_entries) + VALUES ($1, $2, NULLIF($3, ''), $4, $5, $6, 'present', now(), $7, NULLIF($8, ''), $9)` _, err := s.db.ExecContext(ctx, q, - rec.ID, rec.ServerName, rec.FormerOwner, rec.BackupRef, rec.SizeBytes, rec.Reason, rec.ExpiresAt) + rec.ID, rec.ServerName, rec.FormerOwner, rec.BackupRef, rec.SizeBytes, rec.Reason, rec.ExpiresAt, + rec.SHA256, rec.SkippedEntries) return err } @@ -107,7 +109,7 @@ func (s *PGStore) PresentBackupBytes(ctx context.Context) (int64, error) { } func (s *PGStore) EvictableBackups(ctx context.Context) ([]StoredBackup, error) { - const q = `SELECT id, server_name, backup_ref, size_bytes, reason FROM world_backups + const q = `SELECT id, server_name, backup_ref, size_bytes, reason, COALESCE(sha256, '') 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) @@ -118,9 +120,9 @@ func (s *PGStore) EvictableBackups(ctx context.Context) ([]StoredBackup, error) // 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 + const q = `SELECT id, server_name, backup_ref, size_bytes, reason, COALESCE(sha256, '') FROM world_backups WHERE server_name = $1 AND status = 'present' AND reason = $2 AND id <> $4 - ORDER BY created_at DESC OFFSET $3` + ORDER BY corrupt_at IS NULL DESC, 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] @@ -129,7 +131,7 @@ func (s *PGStore) ExcessBackups(ctx context.Context, server, reason string, keep } func (s *PGStore) ListExpiredBackups(ctx context.Context, now time.Time) ([]StoredBackup, error) { - const q = `SELECT id, server_name, backup_ref, size_bytes, reason FROM world_backups + const q = `SELECT id, server_name, backup_ref, size_bytes, reason, COALESCE(sha256, '') FROM world_backups WHERE status = 'present' AND expires_at < $1 ORDER BY expires_at ASC` return s.queryBackups(ctx, q, now) } @@ -143,7 +145,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, &b.Reason); err != nil { + if err := rows.Scan(&b.ID, &b.ServerName, &b.BackupRef, &b.SizeBytes, &b.Reason, &b.SHA256); err != nil { return nil, err } out = append(out, b) @@ -151,6 +153,43 @@ func (s *PGStore) queryBackups(ctx context.Context, q string, args ...any) ([]St return out, rows.Err() } +func (s *PGStore) BackupsToVerify(ctx context.Context, checkedBefore time.Time, limit int) ([]StoredBackup, error) { + const q = `SELECT id, server_name, backup_ref, size_bytes, reason, COALESCE(sha256, '') FROM world_backups + WHERE status = 'present' AND corrupt_at IS NULL AND (verified_at IS NULL OR verified_at < $1) + ORDER BY verified_at NULLS FIRST, created_at LIMIT $2` + return s.queryBackups(ctx, q, checkedBefore, limit) +} + +func (s *PGStore) MarkBackupVerified(ctx context.Context, id, sha256 string, at time.Time) error { + _, err := s.db.ExecContext(ctx, + `UPDATE world_backups SET verified_at = $3, sha256 = COALESCE(sha256, NULLIF($2, '')) WHERE id = $1`, + id, sha256, at) + return err +} + +func (s *PGStore) MarkBackupCorrupt(ctx context.Context, id string, at time.Time) error { + _, err := s.db.ExecContext(ctx, + `UPDATE world_backups SET corrupt_at = $2 WHERE id = $1 AND corrupt_at IS NULL`, id, at) + return err +} + +func (s *PGStore) LiveBackupRefs(ctx context.Context) ([]string, error) { + rows, err := s.db.QueryContext(ctx, `SELECT backup_ref FROM world_backups WHERE status <> 'deleted'`) + if err != nil { + return nil, err + } + defer rows.Close() + var out []string + for rows.Next() { + var ref string + if err := rows.Scan(&ref); err != nil { + return nil, err + } + out = append(out, ref) + } + return out, rows.Err() +} + func (s *PGStore) MarkBackupDeleted(ctx context.Context, id string, at time.Time) error { _, err := s.db.ExecContext(ctx, `UPDATE world_backups SET status = 'deleted', deleted_at = $2 WHERE id = $1`, id, at) diff --git a/internal/reaper/reaper.go b/internal/reaper/reaper.go index 6398c7a..dce4f4d 100644 --- a/internal/reaper/reaper.go +++ b/internal/reaper/reaper.go @@ -97,6 +97,14 @@ type Config struct { // world is deleted on the first run after the copy lands, normally the // next day. RequireOffsite bool + // VerifyEvery is how often each stored archive is read back in full, and + // VerifyPerRun caps how many one run reads (the ones checked longest ago + // first), so bit rot in a backup is found before a restore needs it. + VerifyEvery time.Duration + VerifyPerRun int + // PartialAfter is how old an unfinished archive file must be before it is + // swept: longer than any archive takes to write. + PartialAfter time.Duration } // DefaultConfig is the spec's §24 default window set. @@ -110,6 +118,10 @@ func DefaultConfig() Config { ManualRetention: 30 * Day, ManualKeep: 5, ManualCooldown: 10 * time.Minute, + + VerifyEvery: 7 * Day, + VerifyPerRun: 10, + PartialAfter: 6 * time.Hour, } } @@ -149,11 +161,17 @@ type BackupRecord struct { SizeBytes int64 Reason string ExpiresAt time.Time + // SHA256 is the archive's digest as written ("" when the backend keeps none) + // and SkippedEntries the world entries it could not hold. + SHA256 string + SkippedEntries int } // Fresh is the backup FreshBackup found. type Fresh struct { - Ref string + ID string + Ref string + SHA256 string // Offsite reports that the archive has its off-site copy (offsite_at). Offsite bool } @@ -166,6 +184,7 @@ type StoredBackup struct { BackupRef string SizeBytes int64 Reason string + SHA256 string // "" when none was recorded } // AuditRecord is a reaper-sourced audit_logs entry. The PG binding fills @@ -186,9 +205,10 @@ type Store interface { // FreshBackup reports an existing present reaper archive (reason // inactive_15d) for server whose world is still current — created at or // after since (the world's last_active_at) and after the current claim, - // preferring one already copied off-site. It makes a reap idempotent across - // a DeletePVC failure, and across the wait for the off-site copy: the retry - // reuses the archive instead of writing a duplicate. + // preferring one already copied off-site, and never one found corrupt. It + // makes a reap idempotent across a DeletePVC failure, and across the wait for + // the off-site copy: the retry reuses the archive instead of writing a + // duplicate. FreshBackup(ctx context.Context, server string, since time.Time) (b Fresh, ok bool, err error) // InsertBackup records a world_backups row (status=present). @@ -218,6 +238,19 @@ type Store interface { // MarkBackupDeleted flips a backup to status=deleted, deleted_at=at. MarkBackupDeleted(ctx context.Context, id string, at time.Time) error + // BackupsToVerify lists up to limit present backups not found corrupt and + // not read back since checkedBefore, the ones never read back first, then + // the ones read back longest ago. + BackupsToVerify(ctx context.Context, checkedBefore time.Time, limit int) ([]StoredBackup, error) + // MarkBackupVerified records a read-back that matched at `at`, and the + // archive's digest when none was recorded yet. + MarkBackupVerified(ctx context.Context, id, sha256 string, at time.Time) error + // MarkBackupCorrupt records a read-back that failed: the backup is no + // longer reused for a reap or offered for a restore. + MarkBackupCorrupt(ctx context.Context, id string, at time.Time) error + // LiveBackupRefs lists the backup_ref of every backup not deleted. + LiveBackupRefs(ctx context.Context) ([]string, error) + // Audit appends a reaper-sourced audit_logs row. Audit(ctx context.Context, rec AuditRecord) error } @@ -276,12 +309,27 @@ type Summary struct { // AwaitingOffsite are idle worlds that are archived and kept until the // archive's off-site copy lands. AwaitingOffsite int + // Verified are archives read back in full and found matching; Corrupt are + // the ones that were not (marked, and never reused or restored from); + // VerifyFailed are the ones that could not be read back at all this run. + Verified int + Corrupt int + VerifyFailed int + // Swept are leftover files of interrupted archives removed; OrphanArchives + // are finished archives no backup records, kept for now (see + // backup.Swept); SweepFailed reports that the sweep did not complete. + Swept int + OrphanArchives int + SweepFailed bool } -// Failed reports whether the run left work undone: a server it could not -// process, or an expired backup it could not remove. The world is safe either -// way, but the run did not do its job and whoever operates it must hear. -func (s Summary) Failed() bool { return s.Skipped > 0 || s.ExpireFailed > 0 } +// Failed reports whether the run left work undone or found damage: a server it +// could not process, an expired backup it could not remove, an archive that did +// not read back or could not be read, or a sweep that did not complete. The +// worlds are safe either way, but whoever operates the run must hear. +func (s Summary) Failed() bool { + return s.Skipped > 0 || s.ExpireFailed > 0 || s.Corrupt > 0 || s.VerifyFailed > 0 || s.SweepFailed +} func (r *Reaper) now() time.Time { if r.Now != nil { @@ -340,6 +388,8 @@ func (r *Reaper) RunOnce(ctx context.Context) (Summary, error) { } r.expireBackups(ctx, now, &sum) + r.verifyBackups(ctx, now, &sum) + r.sweepArchives(ctx, now, &sum) return sum, nil } @@ -393,32 +443,49 @@ func (r *Reaper) reap(ctx context.Context, now time.Time, c Candidate, crd Serve if err != nil { return fmt.Errorf("lookup fresh backup: %w", err) } + if ok { + // The archive may have sat on disk for days (a delete that failed, the + // wait for its off-site copy); it is about to become the only copy, so it + // is read back first. One that fails is marked and replaced by a fresh + // archive of the world, which is still there. + good, err := r.verifyOne(ctx, now, fresh.ID, fresh.Ref, fresh.SHA256, sum) + if err != nil { + return fmt.Errorf("read back archive %s: %w", fresh.ID, err) + } + ok = good + } ref, offsite := fresh.Ref, fresh.Offsite if !ok { - aref, size, err := r.Archiver.Archive(ctx, c.Name, crd.PVC) + a, err := r.Archiver.Archive(ctx, c.Name, crd.PVC) if err != nil { // Red line ④: archive failed → the PVC is untouched, the world // survives, and this server is retried next run. return fmt.Errorf("archive: %w", err) } + if len(a.Skipped) > 0 { + r.log().Warn("reaper: archive leaves out entries that are not plain files or directories", + "server", c.Name, "count", len(a.Skipped), "first", a.Skipped[:min(len(a.Skipped), 5)]) + } rec := BackupRecord{ - ID: r.id(), - ServerName: c.Name, - FormerOwner: c.OwnerID, - BackupRef: string(aref), - SizeBytes: size, - Reason: ReasonInactive, - ExpiresAt: now.Add(r.Cfg.Retention), + ID: r.id(), + ServerName: c.Name, + FormerOwner: c.OwnerID, + BackupRef: string(a.Ref), + SizeBytes: a.Size, + Reason: ReasonInactive, + ExpiresAt: now.Add(r.Cfg.Retention), + SHA256: a.SHA256, + SkippedEntries: len(a.Skipped), } if err := r.Store.InsertBackup(ctx, rec); err != nil { // The archive exists but is untracked. Delete the orphan so it does // not leak, then fail without touching the PVC. - if derr := r.Archiver.Delete(ctx, aref); derr != nil { - r.log().Error("reaper: orphan archive cleanup failed", "server", c.Name, "ref", aref, "err", derr) + if derr := r.Archiver.Delete(ctx, a.Ref); derr != nil { + r.log().Error("reaper: orphan archive cleanup failed", "server", c.Name, "ref", a.Ref, "err", derr) } return fmt.Errorf("insert backup: %w", err) } - ref, offsite = string(aref), false + ref, offsite = string(a.Ref), false } // With an off-site bucket configured, the archive on this node's disk is @@ -570,6 +637,93 @@ func (r *Reaper) expireBackups(ctx context.Context, now time.Time, sum *Summary) } } +// verifyBackups is the read-back pass: up to VerifyPerRun present archives not +// read back within VerifyEvery are read end to end and checked against their +// recorded digest. A backend that cannot read its archives back skips it. +func (r *Reaper) verifyBackups(ctx context.Context, now time.Time, sum *Summary) { + if _, ok := r.Archiver.(backup.Verifier); !ok || r.Cfg.VerifyPerRun <= 0 { + return + } + list, err := r.Store.BackupsToVerify(ctx, now.Add(-r.Cfg.VerifyEvery), r.Cfg.VerifyPerRun) + if err != nil { + r.log().Error("reaper: list backups to read back", "err", err) + sum.VerifyFailed++ + return + } + for _, b := range list { + if _, err := r.verifyOne(ctx, now, b.ID, b.BackupRef, b.SHA256, sum); err != nil { + r.log().Error("reaper: read back archive", "id", b.ID, "server", b.ServerName, "err", err) + sum.VerifyFailed++ + } + } +} + +// verifyOne reads one archive back and records the outcome. It reports whether +// the archive is good; an error means it could not be read back at all (the +// store is not mounted, the run was cancelled), which says nothing about the +// archive. A backend that cannot read its archives back reports every archive +// good, as before read-backs existed. +func (r *Reaper) verifyOne(ctx context.Context, now time.Time, id, ref, want string, sum *Summary) (bool, error) { + v, ok := r.Archiver.(backup.Verifier) + if !ok { + return true, nil + } + got, err := v.Verify(ctx, backup.ArchiveRef(ref), want) + switch { + case errors.Is(err, backup.ErrCorrupt): + sum.Corrupt++ + r.log().Error("reaper: archive is corrupt; it will not be reused or offered for restore", "id", id, "ref", ref, "err", err) + if err := r.Store.MarkBackupCorrupt(ctx, id, now); err != nil { + return false, fmt.Errorf("mark corrupt: %w", err) + } + return false, nil + case err != nil: + return false, err + } + if err := r.Store.MarkBackupVerified(ctx, id, got, now); err != nil { + r.log().Error("reaper: record read-back", "id", id, "err", err) + } + sum.Verified++ + return true, nil +} + +// sweepArchives removes what interrupted archives left in the store (see +// backup.Sweeper): unfinished files older than PartialAfter, and finished ones +// no backup records once they are older than Retention, the longest any backup +// is kept. Younger unrecorded archives are reported and kept. +func (r *Reaper) sweepArchives(ctx context.Context, now time.Time, sum *Summary) { + sw, ok := r.Archiver.(backup.Sweeper) + if !ok { + return + } + refs, err := r.Store.LiveBackupRefs(ctx) + if err != nil { + // Without the list every archive would look unclaimed. + r.log().Error("reaper: list recorded archives; sweep skipped", "err", err) + sum.SweepFailed = true + return + } + live := make(map[backup.ArchiveRef]bool, len(refs)) + for _, ref := range refs { + live[backup.ArchiveRef(ref)] = true + } + res, err := sw.Sweep(ctx, func(ref backup.ArchiveRef) bool { return live[ref] }, + now.Add(-r.Cfg.PartialAfter), now.Add(-r.Cfg.Retention)) + if err != nil { + r.log().Error("reaper: sweep archive store", "err", err) + sum.SweepFailed = true + } + for _, p := range res.Removed { + r.log().Info("reaper: removed leftover of an interrupted archive", "path", p) + } + sum.Swept += len(res.Removed) + sum.OrphanArchives += len(res.Orphans) + if len(res.Orphans) > 0 { + r.log().Warn("reaper: archives no backup records are kept until they are older than the retention", + "count", len(res.Orphans), "bytes", res.OrphanBytes, "first", res.Orphans[:min(len(res.Orphans), 5)]) + } +} + // formatRemaining renders an offset as the human-facing time left before reap. func formatRemaining(d time.Duration) string { if d%Day == 0 { diff --git a/internal/reaper/reaper_test.go b/internal/reaper/reaper_test.go index 661bf36..2405daa 100644 --- a/internal/reaper/reaper_test.go +++ b/internal/reaper/reaper_test.go @@ -40,14 +40,48 @@ type fakeArchiver struct { seq int } -func (f *fakeArchiver) Archive(_ context.Context, server, _ string) (backup.ArchiveRef, int64, error) { +func (f *fakeArchiver) Archive(_ context.Context, server, _ string) (backup.Archived, error) { if f.archiveErr != nil { - return "", 0, f.archiveErr + return backup.Archived{}, f.archiveErr } f.archives++ f.seq++ f.rec.add("archive") - return backup.ArchiveRef(fmt.Sprintf("ref-%s-%d", server, f.seq)), 10, nil + ref := fmt.Sprintf("ref-%s-%d", server, f.seq) + return backup.Archived{Ref: backup.ArchiveRef(ref), Size: 10, SHA256: "sha-" + ref}, nil +} + +// checkingArchiver is a fakeArchiver that also reads archives back +// (backup.Verifier) and sweeps its store (backup.Sweeper). +type checkingArchiver struct { + *fakeArchiver + corrupt map[string]bool // refs that no longer read back + verifyErr error // every read-back fails with this + verified []string + + sweepErr error + sweepLive func(backup.ArchiveRef) bool + sweepCutoffs []time.Time + sweepRemoved []string + sweepOrphaned []string +} + +func (c *checkingArchiver) Verify(_ context.Context, ref backup.ArchiveRef, want string) (string, error) { + if c.verifyErr != nil { + return "", c.verifyErr + } + c.verified = append(c.verified, string(ref)) + c.rec.add("verify") + if c.corrupt[string(ref)] { + return "", fmt.Errorf("%w: %s", backup.ErrCorrupt, ref) + } + return "sha-" + string(ref), nil +} + +func (c *checkingArchiver) Sweep(_ context.Context, live func(backup.ArchiveRef) bool, partialBefore, orphanBefore time.Time) (backup.Swept, error) { + c.sweepLive = live + c.sweepCutoffs = []time.Time{partialBefore, orphanBefore} + return backup.Swept{Removed: c.sweepRemoved, Orphans: c.sweepOrphaned, OrphanBytes: int64(len(c.sweepOrphaned))}, c.sweepErr } func (f *fakeArchiver) Restore(context.Context, backup.ArchiveRef, string) error { return nil } @@ -109,6 +143,10 @@ type fakeBackup struct { status string // present | deleted createdAt, expires time.Time offsite bool + sha string + skipped int + verifiedAt time.Time + corruptAt time.Time } type fakeStore struct { @@ -121,6 +159,7 @@ type fakeStore struct { released []string listErr error insertErr error + liveErr error idn int } @@ -138,7 +177,7 @@ func (s *fakeStore) ListActiveServers(context.Context) ([]Candidate, error) { func (s *fakeStore) FreshBackup(_ context.Context, server string, since time.Time) (Fresh, bool, error) { var found *fakeBackup for _, b := range s.backups { - if b.server == server && b.status == "present" && b.reason == ReasonInactive && !b.createdAt.Before(since) { + if b.server == server && b.status == "present" && b.reason == ReasonInactive && !b.createdAt.Before(since) && b.corruptAt.IsZero() { if found == nil || (b.offsite && !found.offsite) { found = b } @@ -147,7 +186,7 @@ func (s *fakeStore) FreshBackup(_ context.Context, server string, since time.Tim if found == nil { return Fresh{}, false, nil } - return Fresh{Ref: found.ref, Offsite: found.offsite}, true, nil + return Fresh{ID: found.id, Ref: found.ref, SHA256: found.sha, Offsite: found.offsite}, true, nil } func (s *fakeStore) InsertBackup(_ context.Context, rec BackupRecord) error { @@ -156,7 +195,7 @@ func (s *fakeStore) InsertBackup(_ context.Context, rec BackupRecord) error { } s.backups = append(s.backups, &fakeBackup{ id: rec.ID, server: rec.ServerName, ref: rec.BackupRef, reason: rec.Reason, size: rec.SizeBytes, - status: "present", createdAt: s.clock, expires: rec.ExpiresAt, + status: "present", createdAt: s.clock, expires: rec.ExpiresAt, sha: rec.SHA256, skipped: rec.SkippedEntries, }) s.rec.add("insert") return nil @@ -233,6 +272,73 @@ func (s *fakeStore) MarkBackupDeleted(_ context.Context, id string, at time.Time return fmt.Errorf("no backup %s", id) } +func (s *fakeStore) BackupsToVerify(_ context.Context, checkedBefore time.Time, limit int) ([]StoredBackup, error) { + var ps []*fakeBackup + for _, b := range s.backups { + if b.status == "present" && b.corruptAt.IsZero() && b.verifiedAt.Before(checkedBefore) { + ps = append(ps, b) + } + } + sort.SliceStable(ps, func(i, j int) bool { + if !ps[i].verifiedAt.Equal(ps[j].verifiedAt) { + return ps[i].verifiedAt.Before(ps[j].verifiedAt) + } + return ps[i].createdAt.Before(ps[j].createdAt) + }) + var out []StoredBackup + for _, b := range ps { + if len(out) == limit { + break + } + out = append(out, StoredBackup{ID: b.id, ServerName: b.server, BackupRef: b.ref, SizeBytes: b.size, Reason: b.reason, SHA256: b.sha}) + } + return out, nil +} + +func (s *fakeStore) find(id string) (*fakeBackup, error) { + for _, b := range s.backups { + if b.id == id { + return b, nil + } + } + return nil, fmt.Errorf("no backup %s", id) +} + +func (s *fakeStore) MarkBackupVerified(_ context.Context, id, sha string, at time.Time) error { + b, err := s.find(id) + if err != nil { + return err + } + b.verifiedAt = at + if b.sha == "" { + b.sha = sha + } + return nil +} + +func (s *fakeStore) MarkBackupCorrupt(_ context.Context, id string, at time.Time) error { + b, err := s.find(id) + if err != nil { + return err + } + b.corruptAt = at + s.rec.add("corrupt") + return nil +} + +func (s *fakeStore) LiveBackupRefs(context.Context) ([]string, error) { + if s.liveErr != nil { + return nil, s.liveErr + } + var out []string + for _, b := range s.backups { + if b.status != "deleted" { + out = append(out, b.ref) + } + } + return out, nil +} + func (s *fakeStore) Audit(_ context.Context, rec AuditRecord) error { s.audits = append(s.audits, rec) s.rec.add("audit:" + rec.Action) @@ -808,3 +914,146 @@ func TestListErrorAbortsBatch(t *testing.T) { t.Fatal("expected a hard error when listing servers fails") } } + +// checking swaps the fixture's archiver for one that reads archives back and +// sweeps. +func checking(r *Reaper, ar *fakeArchiver) *checkingArchiver { + c := &checkingArchiver{fakeArchiver: ar, corrupt: map[string]bool{}} + r.Archiver = c + return c +} + +// An archive left from an earlier run is read back before the world it holds is +// deleted, and the read-back is recorded. +func TestReapReadsBackReusedArchive(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "kilo", OwnerID: "user-1", LastActiveAt: idleBy(20 * Day)}) + ca := checking(r, ar) + st.backups = []*fakeBackup{{id: "old", server: "kilo", ref: "ref-old", reason: ReasonInactive, size: 5, + status: "present", createdAt: idleBy(Day), expires: testNow.Add(89 * Day), sha: "sha-ref-old"}} + + sum := mustRun(t, r) + if sum.WorldsReaped != 1 || ar.archives != 0 || len(cl.deletedPVCs) != 1 { + t.Fatalf("summary %+v archives=%d deleted=%v: want the checked archive reused", sum, ar.archives, cl.deletedPVCs) + } + if st.rec.events[0] != "verify" || st.rec.events[1] != "deletePVC" { + t.Errorf("events = %v, want the read-back before the delete", st.rec.events) + } + if !st.backups[0].verifiedAt.Equal(testNow) || !reflect.DeepEqual(ca.verified, []string{"ref-old"}) { + t.Errorf("read-back not recorded: %+v", st.backups[0]) + } +} + +// A reused archive that no longer reads back is marked corrupt and the world, +// still there, is archived afresh before it is deleted. The run reports it. +func TestReapReplacesCorruptArchive(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "lima", OwnerID: "user-1", LastActiveAt: idleBy(20 * Day)}) + ca := checking(r, ar) + ca.corrupt["ref-rotten"] = true + st.backups = []*fakeBackup{{id: "rotten", server: "lima", ref: "ref-rotten", reason: ReasonInactive, size: 5, + status: "present", createdAt: idleBy(Day), expires: testNow.Add(89 * Day)}} + + sum := mustRun(t, r) + if sum.WorldsReaped != 1 || sum.Corrupt != 1 || !sum.Failed() { + t.Fatalf("summary = %+v, want reaped with one corrupt archive reported", sum) + } + want := []string{"verify", "corrupt", "archive", "insert", "deletePVC"} + if !reflect.DeepEqual(st.rec.events[:len(want)], want) { + t.Errorf("events = %v, want %v first", st.rec.events, want) + } + if st.backups[0].corruptAt.IsZero() || len(st.backups) != 2 || st.backups[1].sha != "sha-ref-lima-1" { + t.Errorf("backups = %+v %+v", st.backups[0], st.backups[len(st.backups)-1]) + } + if len(cl.deletedPVCs) != 1 { + t.Errorf("deleted = %v", cl.deletedPVCs) + } +} + +// When the archive cannot be read back at all (the store is not mounted), the +// world is kept and nothing is marked: the error says nothing about the archive. +func TestReapKeepsWorldWhenReadBackFails(t *testing.T) { + r, st, cl, ar := newReaper(DefaultConfig(), + Candidate{Name: "mike", OwnerID: "user-1", LastActiveAt: idleBy(20 * Day)}) + ca := checking(r, ar) + ca.verifyErr = errors.New("input/output error") + st.backups = []*fakeBackup{{id: "b", server: "mike", ref: "ref-b", reason: ReasonInactive, size: 5, + status: "present", createdAt: idleBy(Day), expires: testNow.Add(89 * Day)}} + + sum := mustRun(t, r) + if sum.WorldsReaped != 0 || sum.Skipped != 1 || cl.deletePVCCalls != 0 || ar.archives != 0 { + t.Fatalf("summary %+v deletes=%d archives=%d: want the world kept", sum, cl.deletePVCCalls, ar.archives) + } + if !st.backups[0].corruptAt.IsZero() { + t.Error("an unreadable store marked the archive corrupt") + } +} + +// The read-back pass takes VerifyPerRun archives, never-checked ones first, and +// records what it finds; a corrupt one is marked and fails the run once. +func TestVerifyPassSamplesArchives(t *testing.T) { + cfg := DefaultConfig() + cfg.VerifyPerRun = 2 + r, st, _, ar := newReaper(cfg) + ca := checking(r, ar) + ca.corrupt["ref-c"] = true + st.backups = []*fakeBackup{ + {id: "a", ref: "ref-a", status: "present", createdAt: idleBy(9 * Day), verifiedAt: idleBy(8 * Day), expires: testNow.Add(30 * Day)}, + {id: "b", ref: "ref-b", status: "present", createdAt: idleBy(9 * Day), verifiedAt: idleBy(Day), expires: testNow.Add(30 * Day)}, + {id: "c", ref: "ref-c", status: "present", createdAt: idleBy(5 * Day), expires: testNow.Add(30 * Day)}, + {id: "d", ref: "ref-d", status: "deleted", createdAt: idleBy(5 * Day), expires: testNow.Add(30 * Day)}, + } + + sum := mustRun(t, r) + if !reflect.DeepEqual(ca.verified, []string{"ref-c", "ref-a"}) { + t.Errorf("read back %v, want [ref-c ref-a]", ca.verified) + } + if sum.Verified != 1 || sum.Corrupt != 1 || !sum.Failed() { + t.Errorf("summary = %+v", sum) + } + if st.backups[2].corruptAt.IsZero() || !st.backups[0].verifiedAt.Equal(testNow) || st.backups[0].sha != "sha-ref-a" { + t.Errorf("outcomes not recorded: %+v %+v", st.backups[0], st.backups[2]) + } + + // The next run skips the corrupt one and what was checked within + // VerifyEvery; a week on, b and a are due again, b first. + ca.verified = nil + if sum := mustRun(t, r); sum.Corrupt != 0 || len(ca.verified) != 0 { + t.Errorf("second run read back %v (%+v), want nothing", ca.verified, sum) + } + r.Now = func() time.Time { return testNow.Add(7*Day + time.Hour) } + if mustRun(t, r); !reflect.DeepEqual(ca.verified, []string{"ref-b", "ref-a"}) { + t.Errorf("a week later read back %v, want [ref-b ref-a]", ca.verified) + } +} + +// The sweep claims every archive a backup records, removes leftovers past +// PartialAfter and orphans past the retention, and never runs without the list +// of recorded archives. +func TestSweepPass(t *testing.T) { + r, st, _, ar := newReaper(DefaultConfig()) + ca := checking(r, ar) + ca.sweepRemoved = []string{"/archives/.x-1.tar.gz.partial"} + ca.sweepOrphaned = []string{"/archives/x-2.tar.gz"} + st.backups = []*fakeBackup{ + {id: "a", ref: "ref-a", status: "present", expires: testNow.Add(Day), verifiedAt: testNow}, + {id: "d", ref: "ref-d", status: "deleted", expires: testNow.Add(Day)}, + } + sum := mustRun(t, r) + if sum.Swept != 1 || sum.OrphanArchives != 1 || sum.Failed() { + t.Errorf("summary = %+v", sum) + } + if !ca.sweepLive("ref-a") || ca.sweepLive("ref-d") { + t.Error("sweep did not get the recorded archives as live") + } + cfg := DefaultConfig() + if want := []time.Time{testNow.Add(-cfg.PartialAfter), testNow.Add(-cfg.Retention)}; !reflect.DeepEqual(ca.sweepCutoffs, want) { + t.Errorf("cutoffs = %v, want %v", ca.sweepCutoffs, want) + } + + ca.sweepLive = nil + st.liveErr = errors.New("db down") + if sum := mustRun(t, r); !sum.SweepFailed || ca.sweepLive != nil { + t.Errorf("sweep ran without the recorded archives: %+v", sum) + } +} diff --git a/internal/store/migrations/0026_world_backups_integrity.sql b/internal/store/migrations/0026_world_backups_integrity.sql new file mode 100644 index 0000000..e55c93a --- /dev/null +++ b/internal/store/migrations/0026_world_backups_integrity.sql @@ -0,0 +1,16 @@ +-- Integrity of world archives. sha256 is the digest of the archive as written +-- (NULL for archives written before digests were kept; the first read-back +-- records one). verified_at is when the reaper last read the archive back in +-- full and found it matching; corrupt_at is when a read-back failed, after +-- which the archive is never reused for a reap or offered for a restore. +-- skipped_entries counts the world entries the archive could not hold +-- (symbolic links, devices, sockets). +ALTER TABLE world_backups + ADD COLUMN sha256 text, + ADD COLUMN verified_at timestamptz, + ADD COLUMN corrupt_at timestamptz, + ADD COLUMN skipped_entries integer NOT NULL DEFAULT 0; + +-- The reaper's read-back work list: the present archives checked longest ago. +CREATE INDEX world_backups_verify_idx ON world_backups (verified_at NULLS FIRST, created_at) + WHERE status = 'present' AND corrupt_at IS NULL; diff --git a/panel/src/i18n/resources/en-US/backups.json b/panel/src/i18n/resources/en-US/backups.json index 10a70ba..16f25d3 100644 --- a/panel/src/i18n/resources/en-US/backups.json +++ b/panel/src/i18n/resources/en-US/backups.json @@ -55,6 +55,14 @@ "restoring": "Restoring…", "stop_timeout": "The server did not stop in time. Close this window and try again shortly.", "expired_cannot_restore": "This backup has expired and can no longer be restored.", + "corrupt_badge": "Corrupt", + "corrupt_hint": "A read-back found this archive no longer matches what was written, so it can't be restored intact. It stays listed until it expires so the loss is visible.", + "corrupt_short": "Corrupt", + "corrupt_cannot_restore": "This backup failed its read-back check and can't be restored — pick another backup.", + "verified_at": "Verified {{when}}", + "skipped_entries_one": "{{count}} entry not archived", + "skipped_entries_other": "{{count}} entries not archived", + "skipped_entries_hint": "The world held symlinks or other entries that aren't plain files, so the archive skipped them and a restore won't bring them back.", "col_created": "Backup Time", "col_size": "Size", "col_expires": "Expires", diff --git a/panel/src/i18n/resources/en-US/errors.json b/panel/src/i18n/resources/en-US/errors.json index 581573b..e8fdba0 100644 --- a/panel/src/i18n/resources/en-US/errors.json +++ b/panel/src/i18n/resources/en-US/errors.json @@ -16,6 +16,7 @@ "not_running": "The server isn't running — wake it before managing access.", "console_unavailable": "Can't reach the server console right now — try again shortly.", "no_backup": "There's no restorable backup for this server yet.", + "backup_corrupt": "This backup failed its read-back check and can't be restored intact — pick another backup.", "not_stopped": "Stop the server completely before restoring — a restore overwrites the live world volume.", "maintenance_in_progress": "This server's world is busy with a restore, backup or file write — try again once it finishes, usually within a minute or two.", "no_world_volume": "This server has no world volume yet — start it once so it is created, then retry.", diff --git a/panel/src/i18n/resources/zh-CN/backups.json b/panel/src/i18n/resources/zh-CN/backups.json index 155e53e..2490bb9 100644 --- a/panel/src/i18n/resources/zh-CN/backups.json +++ b/panel/src/i18n/resources/zh-CN/backups.json @@ -55,6 +55,13 @@ "restoring": "正在恢复备份……", "stop_timeout": "服务器停止超时。请关闭此窗口,稍后重试。", "expired_cannot_restore": "此备份已过期,无法恢复。", + "corrupt_badge": "已损坏", + "corrupt_hint": "回读校验发现这份归档与写入时不一致,已无法完整恢复。它会保留到过期,方便排查。", + "corrupt_short": "已损坏", + "corrupt_cannot_restore": "这份备份回读校验未通过,无法恢复,请选择另一份备份。", + "verified_at": "{{when}}已校验", + "skipped_entries": "{{count}} 个条目未归档", + "skipped_entries_hint": "世界目录里有符号链接等非普通文件,归档时跳过了它们,恢复后这些条目不会存在。", "col_created": "创建时间", "col_size": "大小", "col_expires": "过期时间", diff --git a/panel/src/i18n/resources/zh-CN/errors.json b/panel/src/i18n/resources/zh-CN/errors.json index 327a809..3e16492 100644 --- a/panel/src/i18n/resources/zh-CN/errors.json +++ b/panel/src/i18n/resources/zh-CN/errors.json @@ -16,6 +16,7 @@ "not_running": "服务器未在运行——请先启动它再管理访问权限。", "console_unavailable": "暂时无法连接服务器控制台,请稍后重试。", "no_backup": "这台服务器暂时没有可回档的备份。", + "backup_corrupt": "这份备份回读校验未通过,已无法完整恢复——请选择另一份备份。", "not_stopped": "回档会覆盖世界的实时存储卷,请先把服务器完全停止再回档。", "maintenance_in_progress": "这台服务器的世界正在回档、备份或写入文件——等它完成后再试,通常不超过一两分钟。", "no_world_volume": "这台服务器还没有世界卷——先启动一次让它创建,然后再试。", diff --git a/panel/src/lib/api.ts b/panel/src/lib/api.ts index f281314..3038bbe 100644 --- a/panel/src/lib/api.ts +++ b/panel/src/lib/api.ts @@ -743,6 +743,10 @@ export function humanizeError(e: unknown): string { // the restore subsystem may be unwired (503 restore_unavailable). case "no_backup": return t("no_backup"); + // The backup failed a read-back (sha256 or gzip/tar parse), so the server + // refuses to extract it over the world. + case "backup_corrupt": + return t("backup_corrupt"); case "not_stopped": return t("not_stopped"); // World-volume lock: a restore, backup or file write is running on this diff --git a/panel/src/lib/types.ts b/panel/src/lib/types.ts index f30e3a4..2aacf39 100644 --- a/panel/src/lib/types.ts +++ b/panel/src/lib/types.ts @@ -142,10 +142,11 @@ export interface FleetServer { * by handle; restore resolves the latest present backup server-side. * * Only `status: "present"` rows are ever listed (the query filters them) and the - * list is created_at-descending, so the FIRST row for a given server is exactly - * the one a restore would recover (LatestBackup's WHERE mirrors this) — the UI must - * name that row, not a plausible proxy. `reason`/`status` cross an unvalidated JSON - * boundary; render unknown values tolerantly. */ + * list is created_at-descending, so the first row for a given server that is not + * `corrupt` is exactly the one a restore would recover (LatestBackup's WHERE + * mirrors this) — the UI must name that row, not a plausible proxy. + * `reason`/`status` cross an unvalidated JSON boundary; render unknown values + * tolerantly. */ export interface BackupView { id: string; server_name: string; @@ -156,6 +157,16 @@ export interface BackupView { status: string; created_at: string; expires_at: string; + /** True once the archive failed a read-back (its sha256 or its gzip/tar did + * not check out). It stays listed so the loss is visible, and the API refuses + * to restore it (409 backup_corrupt). */ + corrupt?: boolean; + /** When the reaper last read the archive back intact. Absent until the first + * read-back after it was written. */ + verified_at?: string; + /** How many entries of the world were not plain files or directories + * (symlinks, sockets) and so are not in the archive. */ + skipped_entries?: number; } /** ServerJob is one row of GET /api/v1/servers/{name}/jobs — the observable diff --git a/panel/src/pages/ServerBackups.tsx b/panel/src/pages/ServerBackups.tsx index a6afd47..a739f05 100644 --- a/panel/src/pages/ServerBackups.tsx +++ b/panel/src/pages/ServerBackups.tsx @@ -1,6 +1,7 @@ import { useEffect, useRef, useState } from "react"; import { useParams } from "react-router-dom"; import { + AlertTriangle, Archive, CheckCircle2, Clock, @@ -8,6 +9,7 @@ import { Loader2, RotateCcw, ShieldCheck, + ShieldX, UserMinus, XCircle, } from "lucide-react"; @@ -36,14 +38,12 @@ import { formatBytes, formatRelative, formatAbsolute, isExpired } from "@/lib/fo import { cn } from "@/lib/utils"; import type { BackupView, ServerJob } from "@/lib/types"; -/** LatestBackupCard renders the most-recent backup as the restore card — the one a - * restore actually recovers (the list is created_at-descending and the backend's - * LatestBackup selects the same present row), with the restore action beneath. Older - * archives are shown separately as a compact, read-only history (HistoryRow): a restore - * ALWAYS recovers this latest one, so giving an older backup an action card of its own - * would falsely imply you could restore (or delete) it — the backend offers neither. - * `showOwner` surfaces the former owner (admins list every world's backups; a user only - * ever sees their own). */ +/** BackupRow is one backup in the table, with its own restore action. `isLatest` + * marks the row a restore with no pick recovers: the newest one that is not + * corrupt, which is what the backend's LatestBackup selects. Under the reason it + * shows what the reaper's read-back found — corrupt (restore refused), verified, + * or entries the archive could not hold. `showOwner` surfaces the former owner + * (admins list every world's backups; a user only ever sees their own). */ function BackupRow({ b, isLatest, @@ -93,6 +93,7 @@ function BackupRow({ {reasonLabel} + @@ -119,7 +120,11 @@ function BackupRow({ )} - {!expired ? ( + {b.corrupt ? ( + + {t("corrupt_short")} + + ) : !expired ? ( + {b.corrupt ? ( + + + {t("corrupt_badge")} + + ) : ( + b.verified_at && ( + + + {t("verified_at", { when: formatRelative(b.verified_at, now, locale) })} + + ) + )} + {skipped > 0 && ( + + + {t("skipped_entries", { count: skipped })} + + )} + + ); +} + /** JobStateBadge renders one async Job's state (backup/restore Job history). The * state vocabulary is the API's ("running" | "succeeded" | "failed"); anything * unknown is shown verbatim rather than hidden. */ @@ -265,6 +312,15 @@ function RestoreControls({ const [error, setError] = useState(null); const [safety, setSafety] = useState(true); + if (backup.corrupt) { + if (layout === "row") return null; + return ( +

+ {t("corrupt_cannot_restore")} +

+ ); + } + if (isExpired(backup.expires_at, now)) { if (layout === "row") return null; return ( @@ -525,10 +581,12 @@ export function ServerBackups() { const now = Date.now(); const locale = i18n.language; // The global list, narrowed to this server. Already created_at-descending from the - // API, but re-sorted defensively so all[0] is unambiguously the restore target. + // API, but re-sorted defensively; the newest backup that is not corrupt is the + // one a restore with no pick recovers. const all = (backupsQ.data ?? []) .filter((b) => b.server_name === name) .sort((a, b) => Date.parse(b.created_at) - Date.parse(a.created_at)); + const latestID = all.find((b) => !b.corrupt)?.id; const header = ( - {all.map((b, idx) => ( + {all.map((b) => (