fix(offsite): 复制新归档后补打并上传 DB 包,桶内无记录的归档过最长保留期后清理

This commit is contained in:
Lemon-miaow committed 2026-09-27 09:16:00 +08:00
1 parent c146ce84e2
commit 983727d9c0
10 files changed
+669 -28

No files matched your search

+20
View File
@@ -240,6 +240,26 @@ func (c *fakeCatalog) ExpiredRefs(_ context.Context, now time.Time) ([]string, e
return out, nil
}
func (c *fakeCatalog) NewestOffsite(context.Context) (time.Time, error) {
var at time.Time
for _, r := range c.rows {
if r.offsite.After(at) {
at = r.offsite
}
}
return at, nil
}
func (c *fakeCatalog) KeptRefs(_ context.Context, now time.Time) ([]string, error) {
var out []string
for _, r := range c.rows {
if r.status == "present" || !r.expires.Before(now) {
out = append(out, r.Ref)
}
}
return out, nil
}
func writeFile(t *testing.T, dir, name string, size int) []byte {
t.Helper()
data := make([]byte, size)
+15 -1
View File
@@ -43,8 +43,12 @@ func (c PGCatalog) MarkOffsite(ctx context.Context, id string, at time.Time) err
}
func (c PGCatalog) ExpiredRefs(ctx context.Context, now time.Time) ([]string, error) {
rows, err := c.DB.QueryContext(ctx, `SELECT backup_ref FROM world_backups
return c.refs(ctx, `SELECT backup_ref FROM world_backups
WHERE status = 'deleted' AND expires_at < $1 AND offsite_at IS NOT NULL`, now)
}
func (c PGCatalog) refs(ctx context.Context, q string, args ...any) ([]string, error) {
rows, err := c.DB.QueryContext(ctx, q, args...)
if err != nil {
return nil, err
}
@@ -60,6 +64,16 @@ func (c PGCatalog) ExpiredRefs(ctx context.Context, now time.Time) ([]string, er
return out, rows.Err()
}
func (c PGCatalog) NewestOffsite(ctx context.Context) (time.Time, error) {
var at sql.NullTime
err := c.DB.QueryRowContext(ctx, `SELECT max(offsite_at) FROM world_backups`).Scan(&at)
return at.Time, err
}
func (c PGCatalog) KeptRefs(ctx context.Context, now time.Time) ([]string, error) {
return c.refs(ctx, `SELECT backup_ref FROM world_backups WHERE status = 'present' OR expires_at >= $1`, now)
}
// PendingCount is how many present archives wait for their copy, and since
// when the oldest has waited.
func (c PGCatalog) PendingCount(ctx context.Context) (int, time.Time, error) {
+286
View File
@@ -0,0 +1,286 @@
package offsite
import (
"bytes"
"context"
"errors"
"os"
"path/filepath"
"slices"
"sort"
"strings"
"testing"
"time"
"felis.lolicon.best/internal/dbbackup"
)
// snapshotter stands in for `felis db backup -label offsite`: each call writes
// a bundle stamped with the syncer's clock into its DBDir.
type snapshotter struct {
t *testing.T
s *Syncer
calls int
err error
}
func (sn *snapshotter) take(context.Context) error {
sn.calls++
if sn.err != nil {
return sn.err
}
name := dbbackup.BundleName(sn.s.now(), dbbackup.LabelOffsite)
data := dbBundle(sn.t, sn.s.now(), &dbbackup.Counts{Users: 3, Servers: 2}, 64)
return os.WriteFile(filepath.Join(sn.s.DBDir, name), data, 0o600)
}
func withSnapshot(t *testing.T, s *Syncer) *snapshotter {
sn := &snapshotter{t: t, s: s}
s.Snapshot = sn.take
return sn
}
func dbKeys(b *memBucket) []string {
var keys []string
for k := range b.objs {
if strings.HasPrefix(k, dbDir) {
keys = append(keys, strings.TrimSuffix(strings.TrimPrefix(k, dbDir), objExt))
}
}
sort.Strings(keys)
return keys
}
// TestSyncSnapshotsAfterCopyingWorlds: a pass that copies an archive leaves a
// bundle in the bucket taken after the copy, which a restore takes as the
// latest and whose rows list the archive; the pass counts the bucket's bundles
// once, with that one newest. A pass that copies nothing takes none, and the
// next daily bundle replaces it.
func TestSyncSnapshotsAfterCopyingWorlds(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/a/alpha-1.tar.gz"}, status: "present"},
}}
s, b := newSyncer(t, cat)
sn := withSnapshot(t, s)
writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 100)
daily := "felis-db-20260924T030000Z-daily.tar"
writeFile(t, s.DBDir, daily, 50)
res, err := s.Run(context.Background())
if err != nil {
t.Fatalf("run: %v (%+v)", err, res)
}
snap := dbbackup.BundleName(now, dbbackup.LabelOffsite)
if sn.calls != 1 {
t.Fatalf("snapshots taken = %d, want 1 after copying an archive", sn.calls)
}
if got := dbKeys(b); !slices.Equal(got, []string{daily, snap}) {
t.Fatalf("bucket bundles = %v, want the daily and the snapshot", got)
}
if res.NewestDB != snap || res.RemoteDB != 2 || res.DBUploaded != 2 || res.WorldsUploaded != 1 {
t.Fatalf("result = %+v, want the snapshot newest of 2 bundles, 2 bundles and 1 archive copied", res)
}
name, _, err := ChooseDB(context.Background(), b, s.Key)
if err != nil || name != snap {
t.Fatalf("fetch-db latest = %q, %v; want the snapshot", name, err)
}
if res, err := s.Run(context.Background()); err != nil || sn.calls != 1 || res.NewestDB != snap || res.RemoteDB != 2 {
t.Fatalf("a pass that copied nothing: snapshots %d, result %+v, err %v; want no new snapshot", sn.calls, res, err)
}
next := now.Add(15 * time.Hour)
s.Now = func() time.Time { return next }
nextDaily := dbbackup.BundleName(next.Add(-time.Minute), dbbackup.LabelDaily)
writeFile(t, s.DBDir, nextDaily, 50)
res, err = s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
if got := dbKeys(b); !slices.Equal(got, []string{daily, nextDaily}) || sn.calls != 1 || res.DBPruned != 1 {
t.Fatalf("after the next daily: bucket %v, snapshots %d, result %+v; want the snapshot pruned, both dailies kept", got, sn.calls, res)
}
}
// TestSyncSnapshotOnlyWhenACopyIsNewer: no snapshot while the newest bundle in
// the bucket was taken after the last copy, in the same second included (bundle
// names have one-second resolution); one when a copy is newer, or when the
// bucket has no bundle at all; none on an install that never copied an archive.
func TestSyncSnapshotOnlyWhenACopyIsNewer(t *testing.T) {
stamped := func(at time.Time) string { return dbbackup.BundleName(at, dbbackup.LabelDaily) }
cases := []struct {
name string
copied time.Time
bundle string // in the bucket before the pass; "" for none
want int
wantNew string
}{
{name: "never copied", want: 0},
{name: "bundle after the copy", copied: now.Add(-time.Hour), bundle: stamped(now.Add(-time.Minute)), want: 0},
{name: "same second", copied: now.Add(400 * time.Millisecond), bundle: stamped(now), want: 0},
{name: "copy after the bundle", copied: now.Add(-time.Hour), bundle: stamped(now.Add(-2 * time.Hour)), want: 1},
{name: "no bundle", copied: now.Add(-time.Hour), want: 1},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Ref: "/a/alpha-1.tar.gz"}, status: "present", offsite: c.copied},
}}
if c.copied.IsZero() {
cat.rows[0].status = "deleted"
}
s, b := newSyncer(t, cat)
sn := withSnapshot(t, s)
if c.bundle != "" {
b.objs[DBKey(c.bundle)] = seal(t, []byte("bundle"), s.Key)
}
res, err := s.Run(context.Background())
if err != nil {
t.Fatalf("run: %v", err)
}
if sn.calls != c.want {
t.Fatalf("snapshots taken = %d, want %d (newest %s)", sn.calls, c.want, res.NewestDB)
}
if c.want == 1 && res.NewestDB != dbbackup.BundleName(now, dbbackup.LabelOffsite) {
t.Fatalf("newest bundle = %s, want the snapshot", res.NewestDB)
}
})
}
// A pass that copies no bundles (-db-dir "") takes none either.
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Ref: "/a/alpha-1.tar.gz"}, status: "present", offsite: now.Add(-time.Hour)},
}}
s, _ := newSyncer(t, cat)
s.DBDir = ""
calls := 0
s.Snapshot = func(context.Context) error { calls++; return errors.New("no bundle directory") }
if _, err := s.Run(context.Background()); err != nil || calls != 0 {
t.Fatalf("no bundle directory: snapshots %d, err %v; want none", calls, err)
}
}
// TestSyncSnapshotFailureIsReported: a snapshot that fails fails the pass, which
// the watchdog reports, and the bundles already there still count.
func TestSyncSnapshotFailureIsReported(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/a/alpha-1.tar.gz"}, status: "present"},
}}
s, _ := newSyncer(t, cat)
sn := withSnapshot(t, s)
sn.err = errors.New("pg_dump: connection refused")
writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 100)
daily := "felis-db-20260924T030000Z-daily.tar"
writeFile(t, s.DBDir, daily, 50)
res, err := s.Run(context.Background())
if err == nil || !strings.Contains(err.Error(), "pg_dump: connection refused") {
t.Fatalf("run error = %v, want the snapshot's failure", err)
}
if res.NewestDB != daily || res.RemoteDB != 1 || res.WorldsUploaded != 1 {
t.Fatalf("result = %+v, want the daily still counted and the archive copied", res)
}
}
// TestSyncBucketKeepsDailiesPastSnapshots: snapshots do not count against
// DBKeep, so hourly ones leave the daily restore points alone, and only the
// newest snapshot stays, while it is the newest bundle of all.
func TestSyncBucketKeepsDailiesPastSnapshots(t *testing.T) {
s, b := newSyncer(t, &fakeCatalog{})
for _, name := range []string{
"felis-db-20260922T030000Z-daily.tar",
"felis-db-20260923T030000Z-daily.tar",
"felis-db-20260924T030000Z-daily.tar",
} {
b.objs[DBKey(name)] = []byte("daily")
}
for h := 4; h <= 9; h++ {
name := dbbackup.BundleName(time.Date(2026, 9, 24, h, 0, 0, 0, time.UTC), dbbackup.LabelOffsite)
b.objs[DBKey(name)] = []byte("snapshot")
}
b.objs[DBKey("felis-db-20260923T120000Z-offsite.tar")] = []byte("older than a daily")
res, err := s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
want := []string{"felis-db-20260923T030000Z-daily.tar", "felis-db-20260924T030000Z-daily.tar", "felis-db-20260924T090000Z-offsite.tar"}
if got := dbKeys(b); !slices.Equal(got, want) {
t.Fatalf("bucket bundles = %v, want %v", got, want)
}
if res.RemoteDB != 3 || res.DBPruned != 7 || res.NewestDB != want[2] {
t.Fatalf("result = %+v, want 3 held, 7 pruned, the snapshot newest", res)
}
got := keptBundles([]string{"felis-db-20260925T030000Z-daily.tar", "felis-db-20260924T090000Z-offsite.tar", "felis-db-20260924T030000Z-daily.tar"}, 2)
if len(got) != 2 || got["felis-db-20260924T090000Z-offsite.tar"] {
t.Fatalf("kept = %v, want the two dailies: a snapshot older than a daily lists nothing it does not", got)
}
}
// TestSyncSweepsUnrecordedWorlds: a world object no row keeps goes once it is
// older than OrphanAfter; a younger one, one whose age the bucket does not
// report, one a row keeps however old, and anything that is not an archive
// object stay.
func TestSyncSweepsUnrecordedWorlds(t *testing.T) {
day := 24 * time.Hour
rows := func() []*row {
return []*row{
{WorldBackup: WorldBackup{ID: "kept", Ref: "/a/kept.tar.gz"}, status: "present", offsite: now.Add(-200 * day)},
{WorldBackup: WorldBackup{ID: "evicted", Ref: "/a/evicted.tar.gz"}, status: "deleted", expires: now.Add(day), offsite: now.Add(-100 * day)},
// Copied, but the run died before MarkOffsite, and the row has
// since expired: ExpiredRefs never lists it.
{WorldBackup: WorldBackup{ID: "unmarked", Ref: "/a/unmarked.tar.gz"}, status: "deleted", expires: now.Add(-day)},
}
}
ages := map[string]time.Duration{
"worlds/kept.tar.gz.fenc": 200 * day,
"worlds/evicted.tar.gz.fenc": 100 * day,
"worlds/unmarked.tar.gz.fenc": 95 * day,
"worlds/restored-old.tar.gz.fenc": 91 * day,
"worlds/restored-new.tar.gz.fenc": 10 * day,
"worlds/no-age.tar.gz.fenc": 0,
"worlds/README": 300 * day,
}
setup := func(t *testing.T, cat *fakeCatalog) (*Syncer, *memBucket, *bytes.Buffer) {
s, b := newSyncer(t, cat)
s.OrphanAfter = 90 * day
var log bytes.Buffer
s.Log = &log
b.modified = map[string]time.Time{}
for k, age := range ages {
b.objs[k] = []byte("x")
if age > 0 {
b.modified[k] = now.Add(-age)
}
}
return s, b, &log
}
s, b, log := setup(t, &fakeCatalog{rows: rows()})
res, err := s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
removed := slices.Sorted(slices.Values(b.removed))
if want := []string{"worlds/restored-old.tar.gz.fenc", "worlds/unmarked.tar.gz.fenc"}; !slices.Equal(removed, want) {
t.Fatalf("removed %v, want %v", removed, want)
}
if res.WorldsExpired != 2 || res.RemoteWorlds != len(ages)-2 {
t.Fatalf("result = %+v, want 2 expired, %d left", res, len(ages)-2)
}
if !strings.Contains(log.String(), "2 world archives in the bucket have no backup record") {
t.Errorf("the young unrecorded archives were not reported:\n%s", log.String())
}
// A database not restored yet keeps nothing: nothing is swept on it.
s, b, _ = setup(t, &fakeCatalog{})
if _, err := s.Run(context.Background()); err != nil || len(b.removed) != 0 {
t.Fatalf("an empty catalog: removed %v, err %v; want nothing", b.removed, err)
}
s, b, _ = setup(t, &fakeCatalog{rows: rows()})
s.OrphanAfter = 0
if _, err := s.Run(context.Background()); err != nil || len(b.removed) != 0 {
t.Fatalf("no OrphanAfter: removed %v, err %v; want nothing", b.removed, err)
}
}
+136 -10
View File
@@ -12,7 +12,11 @@
// deletes an idle world only after that (internal/reaper). Remote world
// objects go when their row has expired, so the bucket keeps each archive for
// the same retention the panel promises, including archives evicted early from
// the local disk to make room. Remote bundles are pruned to the newest DBKeep.
// the local disk to make room; an object whose row a restore did not bring
// back goes once it is older than the longest retention. Remote bundles are
// pruned to the newest DBKeep. A pass that copies archives then sends a fresh
// bundle, whose rows list them, so the newest bundle in the bucket lists every
// archive in it and a restore from it can fetch them all.
package offsite
import (
@@ -78,6 +82,12 @@ type Catalog interface {
ExpiredRefs(ctx context.Context, now time.Time) ([]string, error)
// PresentWorlds lists every present archive, for a restore of the volume.
PresentWorlds(ctx context.Context) ([]WorldBackup, error)
// NewestOffsite is the latest offsite_at of any row: when the bucket last
// gained an archive. Zero when no archive was ever copied.
NewestOffsite(ctx context.Context) (time.Time, error)
// KeptRefs lists the backup_ref of every row whose archive the bucket
// keeps: present ones, and the rest until their expires_at.
KeptRefs(ctx context.Context, now time.Time) ([]string, error)
}
// Syncer copies what is missing from the bucket and prunes what has expired.
@@ -89,9 +99,16 @@ type Syncer struct {
// there is none yet (nothing has been archived on this install).
ArchiveDir string
// DBDir holds the database bundles; DBKeep is how many of the newest the
// bucket keeps.
// bucket keeps, not counting the newest Snapshot bundle.
DBDir string
DBKeep int
// Snapshot writes a bundle labelled dbbackup.LabelOffsite into DBDir. A
// pass that leaves the bucket holding an archive newer than its newest
// bundle takes one and sends it (snapshotDB); nil takes none.
Snapshot func(ctx context.Context) error
// OrphanAfter is how old a world object no row keeps (KeptRefs) may get
// before it goes: the longest retention of any archive. Zero keeps them.
OrphanAfter time.Duration
// Images is the platform registry whose user images are copied (images.go);
// nil copies none.
Images ImageSource
@@ -190,10 +207,11 @@ func (s *Syncer) logf(format string, args ...any) {
}
}
// Run does one pass: database bundles, world archives, registry images,
// submission uploads, then expiry. The bundles go first: they are small, and
// every restore starts from one. A failure on one item is recorded and the
// pass carries on; the returned error is non-nil when anything failed.
// Run does one pass: database bundles, world archives, a bundle listing the
// archives just copied, registry images, submission uploads, then expiry. The
// bundles go first: they are small, and every restore starts from one. A
// failure on one item is recorded and the pass carries on; the returned error
// is non-nil when anything failed.
func (s *Syncer) Run(ctx context.Context) (Result, error) {
var res Result
fail := func(format string, args ...any) {
@@ -214,9 +232,11 @@ func (s *Syncer) Run(ctx context.Context) (Result, error) {
}
s.syncDB(ctx, &res, fail)
s.syncWorlds(ctx, remoteWorlds, &res, fail)
s.snapshotDB(ctx, &res, fail)
s.syncImages(ctx, &res, fail)
s.syncUploads(ctx, &res, fail)
s.expireWorlds(ctx, remoteWorlds, &res, fail)
s.sweepWorlds(ctx, remoteWorlds, &res, fail)
for _, size := range remoteWorlds {
res.RemoteWorlds++
@@ -333,10 +353,7 @@ func (s *Syncer) syncDB(ctx context.Context, res *Result, fail func(string, ...a
// Bundle names start with their UTC stamp, so reversed order is newest first.
ranked := slices.Sorted(maps.Keys(names))
slices.Reverse(ranked)
kept := map[string]bool{}
for _, name := range ranked[:min(keep, len(ranked))] {
kept[name] = true
}
kept := keptBundles(ranked, keep)
for _, b := range local {
if !kept[b.Name] {
continue
@@ -367,6 +384,7 @@ func (s *Syncer) syncDB(ctx context.Context, res *Result, fail func(string, ...a
delete(remote, key)
res.DBPruned++
}
res.RemoteDB, res.NewestDB = 0, ""
for _, name := range ranked {
if _, ok := remote[DBKey(name)]; ok {
res.RemoteDB++
@@ -377,6 +395,59 @@ func (s *Syncer) syncDB(ctx context.Context, res *Result, fail func(string, ...a
}
}
// keptBundles picks the bundles the bucket keeps from ranked, newest first:
// the newest keep that are not snapshots, and the newest snapshot while it is
// the newest of all. A pass takes a snapshot whenever it copies archives, up to
// once an hour, so counting them in keep would push the daily bundles, the
// restore points of the last fortnight, out of the bucket within a day; one
// older than the newest daily lists nothing that daily does not.
func keptBundles(ranked []string, keep int) map[string]bool {
kept := map[string]bool{}
n := 0
for i, name := range ranked {
if _, label, _ := dbbackup.ParseBundleName(name); label == dbbackup.LabelOffsite {
if i == 0 {
kept[name] = true
}
continue
}
if n++; n > keep {
break
}
kept[name] = true
}
return kept
}
// snapshotDB takes a bundle and sends it when the bucket holds an archive
// copied after its newest bundle was taken. A restore finds archives through
// the world_backups rows of the bundle it restores: an archive with no row
// there is neither fetched back nor ever expired, and until the next daily
// bundle, up to a day later, every archive copied since the last one was in
// that state. The comparison is by second, the resolution of bundle names; a
// bundle taken in the second of the copy was taken after it.
func (s *Syncer) snapshotDB(ctx context.Context, res *Result, fail func(string, ...any)) {
if s.Snapshot == nil || s.DBDir == "" {
return
}
copied, err := s.Catalog.NewestOffsite(ctx)
if err != nil {
fail("read when an archive was last copied: %v", err)
return
}
if copied.IsZero() {
return
}
if taken, _, ok := dbbackup.ParseBundleName(res.NewestDB); ok && !copied.Truncate(time.Second).After(taken) {
return
}
if err := s.Snapshot(ctx); err != nil {
fail("take a database bundle listing the archives just copied: %v", err)
return
}
s.syncDB(ctx, res, fail)
}
func (s *Syncer) expireWorlds(ctx context.Context, remote map[string]int64, res *Result, fail func(string, ...any)) {
refs, err := s.Catalog.ExpiredRefs(ctx, s.now())
if err != nil {
@@ -400,6 +471,61 @@ func (s *Syncer) expireWorlds(ctx context.Context, remote map[string]int64, res
}
}
// sweepWorlds removes world objects no row keeps once they are older than
// OrphanAfter, the longest any archive is kept: what a restore from an older
// bundle leaves in the bucket (its archives were copied after that bundle),
// and the copy of a row whose MarkOffsite failed before the row expired.
// Younger ones are kept and reported, as the reaper does on the archive
// volume. A catalog that keeps nothing is a database not restored yet, whose
// rows would keep everything, so nothing is swept on it; an object whose age
// the bucket does not report stays.
func (s *Syncer) sweepWorlds(ctx context.Context, remote map[string]int64, res *Result, fail func(string, ...any)) {
if s.OrphanAfter <= 0 {
return
}
refs, err := s.Catalog.KeptRefs(ctx, s.now())
if err != nil {
fail("list the world archives the bucket keeps: %v", err)
return
}
if len(refs) == 0 {
return
}
kept := make(map[string]bool, len(refs))
for _, ref := range refs {
if key, ok := WorldKey(ref); ok {
kept[key] = true
}
}
objs, err := s.Bucket.List(ctx, worldsDir)
if err != nil {
fail("list %s in the bucket: %v", worldsDir, err)
return
}
cutoff := s.now().Add(-s.OrphanAfter)
var young []string
for _, o := range objs {
if _, ok := remote[o.Key]; !ok || kept[o.Key] || !strings.HasSuffix(o.Key, objExt) {
continue
}
if o.Modified.IsZero() || o.Modified.After(cutoff) {
young = append(young, path.Base(o.Key))
continue
}
if err := s.Bucket.Remove(ctx, o.Key); err != nil {
fail("remove %s, which no backup records: %v", o.Key, err)
continue
}
delete(remote, o.Key)
res.WorldsExpired++
s.logf("removed world archive %s: no backup records it, and it was copied %s ago, past the longest retention", path.Base(o.Key), s.now().Sub(o.Modified).Round(time.Hour))
}
if len(young) > 0 {
s.logf("%d world archives in the bucket have no backup record (a database restored from an older bundle?); each goes once it is older than %s: %s",
len(young), s.OrphanAfter, strings.Join(young[:min(len(young), 5)], ", "))
}
}
// putFile encrypts the file at p into key.
func (s *Syncer) putFile(ctx context.Context, key, p string, size int64) error {
f, err := os.Open(p)