fix(reaper): 认领时重置活跃时钟与预警,回收只复用本次认领后的 reaper 归档

This commit is contained in:
Lemon-miaow committed 2026-09-24 19:31:23 +08:00
1 parent 47890ca913
commit 22becd6859
7 files changed
+187 -8

No files matched your search

+7
View File
@@ -13,6 +13,8 @@ graded for how far the in-repo Go test suite proves the behaviour:
- **[GO-TESTED]** — a hermetic `*_test.go` exercises this exact path; the - **[GO-TESTED]** — a hermetic `*_test.go` exercises this exact path; the
string/code is asserted in CI. string/code is asserted in CI.
- **[PG-TESTED]** — a `-tags pgint` test in `internal/pgint` drives it against a
real Postgres with the shipped migrations (CONTRIBUTING.md has the command).
- **[CODE-ONLY]** — the code path and string exist and are real, but no unit - **[CODE-ONLY]** — the code path and string exist and are real, but no unit
test drives them (notably the Velocity Java plugin, which is not compiled or test drives them (notably the Velocity Java plugin, which is not compiled or
tested in this repo). tested in this repo).
@@ -697,6 +699,11 @@ SMTP sink.]
the activity clock never resets** and an actively-played world becomes the activity clock never resets** and an actively-played world becomes
reap-eligible after 15 days. Verify join events are flowing (§3a) — this is the reap-eligible after 15 days. Verify join events are flowing (§3a) — this is the
most important reaper check. [INTEGRATION-ONLY for the live Postgres write.] most important reaper check. [INTEGRATION-ONLY for the live Postgres write.]
A claim also resets it: `ClaimServer` sets `last_active_at` to the claim and
clears `warned_*`, so a world reaped long ago and claimed again gets a full 15
days, and a later reap archives the new owner's world rather than reusing the
previous owner's archive (`FreshBackup` counts only `inactive_15d` archives
taken since the claim). [PG-TESTED `TestReclaimRestartsReaperClock`.]
- **Unowned servers are still reaped.** A server with `owner_id=""` gets **no - **Unowned servers are still reaped.** A server with `owner_id=""` gets **no
pre-deletion warning** (`maybeWarn` skips unowned), but is still reaped at 15d. pre-deletion warning** (`maybeWarn` skips unowned), but is still reaped at 15d.
[GO-TESTED `TestReapUnownedServerStillReaped`.] Claim or exempt servers you [GO-TESTED `TestReapUnownedServerStillReaped`.] Claim or exempt servers you
+8 -1
View File
@@ -417,6 +417,12 @@ func (p *PGRepo) ServerResources(ctx context.Context, name string) (ResourceSpec
// user for two DIFFERENT ownerless servers serialize instead of both passing the // user for two DIFFERENT ownerless servers serialize instead of both passing the
// gate (audit #4); the row is additionally taken FOR UPDATE so concurrent claims // gate (audit #4); the row is additionally taken FOR UPDATE so concurrent claims
// of the SAME server still resolve to exactly one winner. // of the SAME server still resolve to exactly one winner.
//
// A claim starts the reaper's clock afresh: last_active_at moves to the claim
// and the pre-reap warnings clear. A world reaped before keeps its release time
// as last_active_at, so without the reset a new owner who configures it from
// the panel before anyone joins would lose it on the next reaper run, with no
// warning and no archive of their own.
func (p *PGRepo) ClaimServer(ctx context.Context, name, userID string) (bool, error) { func (p *PGRepo) ClaimServer(ctx context.Context, name, userID string) (bool, error) {
tx, err := p.db.BeginTx(ctx, nil) tx, err := p.db.BeginTx(ctx, nil)
if err != nil { if err != nil {
@@ -470,7 +476,8 @@ func (p *PGRepo) ClaimServer(ctx context.Context, name, userID string) (bool, er
} }
res, err := tx.ExecContext(ctx, res, err := tx.ExecContext(ctx,
`UPDATE servers SET owner_id = $2, claimed_at = now() WHERE name = $1 AND owner_id IS NULL AND deleted_at IS NULL`, `UPDATE servers SET owner_id = $2, claimed_at = now(), last_active_at = now(), warned_3d_at = NULL, warned_1d_at = NULL
WHERE name = $1 AND owner_id IS NULL AND deleted_at IS NULL`,
name, userID) name, userID)
if err != nil { if err != nil {
return false, err return false, err
+2
View File
@@ -244,6 +244,8 @@ type Repo interface {
// authoritative gate: a claim that would exceed a cap → ErrQuotaExceeded (403), // authoritative gate: a claim that would exceed a cap → ErrQuotaExceeded (403),
// and two concurrent claims by one user cannot both pass (audit #4). // and two concurrent claims by one user cannot both pass (audit #4).
// QuotaCheck remains the advisory pre-check for the handler's fast-path 403. // QuotaCheck remains the advisory pre-check for the handler's fast-path 403.
// A successful claim resets last_active_at to now and clears warned_*, so the
// reaper counts idleness from the claim.
ClaimServer(ctx context.Context, name, userID string) (bool, error) ClaimServer(ctx context.Context, name, userID string) (bool, error)
// UserInAllowlist reports whether the user's linked UUID is on the server // UserInAllowlist reports whether the user's linked UUID is on the server
// allowlist (spec §9.4). // allowlist (spec §9.4).
+141
View File
@@ -0,0 +1,141 @@
//go:build pgint
package pgint
import (
"context"
"database/sql"
"testing"
"time"
"felis.lolicon.best/internal/backup"
"felis.lolicon.best/internal/reaper"
)
// reclaimCluster knows one server; every other row in the shared schema reads
// as a server whose CRD is gone, which the reaper skips.
type reclaimCluster struct {
name string
deletedPVC []string
}
func (c *reclaimCluster) Inspect(_ context.Context, name string) (reaper.ServerCRD, error) {
if name != c.name {
return reaper.ServerCRD{}, reaper.ErrNotFound
}
return reaper.ServerCRD{PVC: "world-" + name + "-0"}, nil
}
func (c *reclaimCluster) DeletePVC(_ context.Context, pvc string) error {
c.deletedPVC = append(c.deletedPVC, pvc)
return nil
}
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) {
ref := "/archives/" + server + "-new.tar.gz"
a.archived = append(a.archived, ref)
return backup.ArchiveRef(ref), 42, nil
}
func (a *reclaimArchiver) Restore(context.Context, backup.ArchiveRef, string) error { return nil }
func (a *reclaimArchiver) Delete(context.Context, backup.ArchiveRef) error { return nil }
// TestReclaimRestartsReaperClock: a world reaped weeks ago and claimed by a new
// owner is not reaped again on the next run, and when it does go idle the reap
// archives the new owner's world instead of reusing the previous owner's
// archive (data-durability-5).
func TestReclaimRestartsReaperClock(t *testing.T) {
ctx := context.Background()
st := reaper.NewPGStore(db)
name := "reclaim-" + 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)
}
first := newUser(t, "user", "reclaim-a")
if ok, err := repo.ClaimServer(ctx, name, first.ID); err != nil || !ok {
t.Fatalf("first claim = (%v, %v)", ok, err)
}
// The first owner went idle and the reaper took the world 20 days ago.
reapedAt := time.Now().Add(-20 * reaper.Day)
if _, err := db.ExecContext(ctx,
`UPDATE servers SET last_active_at = $2, claimed_at = $3 WHERE name = $1`,
name, reapedAt.Add(-16*reaper.Day), reapedAt.Add(-60*reaper.Day)); err != nil {
t.Fatal(err)
}
if _, err := db.ExecContext(ctx,
`INSERT INTO world_backups (id, server_name, former_owner, backup_ref, size_bytes, reason, status, created_at, expires_at)
VALUES ($1, $2, $3, $4, 7, 'inactive_15d', 'present', $5, $6)`,
"bk-"+suffix(t), name, first.ID, "/archives/"+name+"-old.tar.gz", reapedAt, reapedAt.Add(90*reaper.Day)); err != nil {
t.Fatalf("seed the old archive: %v", err)
}
if err := st.ReleaseWorld(ctx, name, reapedAt); err != nil {
t.Fatalf("ReleaseWorld: %v", err)
}
// Leftover warning stamps must not survive into the next ownership.
if _, err := db.ExecContext(ctx,
`UPDATE servers SET warned_3d_at = $2, warned_1d_at = $2 WHERE name = $1`, name, reapedAt); err != nil {
t.Fatal(err)
}
second := newUser(t, "user", "reclaim-b")
before := time.Now()
if ok, err := repo.ClaimServer(ctx, name, second.ID); err != nil || !ok {
t.Fatalf("second claim = (%v, %v)", ok, err)
}
var lastActive time.Time
var w3, w1 sql.NullTime
if err := db.QueryRowContext(ctx,
`SELECT last_active_at, warned_3d_at, warned_1d_at FROM servers WHERE name = $1`, name).Scan(&lastActive, &w3, &w1); err != nil {
t.Fatal(err)
}
if lastActive.Before(before.Add(-time.Minute)) {
t.Fatalf("last_active_at = %v after the claim at %v: the claim did not restart the clock", lastActive, before)
}
if w3.Valid || w1.Valid {
t.Fatalf("warnings survived the claim: 3d=%v 1d=%v", w3, w1)
}
// The previous owner's archive is not this owner's world, whatever since says.
if f, ok, err := st.FreshBackup(ctx, name, reapedAt.Add(-30*reaper.Day)); err != nil || ok {
t.Fatalf("FreshBackup = (%+v, %v, %v): the old archive counted as the new world's", f, ok, err)
}
cl := &reclaimCluster{name: name}
ar := &reclaimArchiver{}
now := time.Now()
r := &reaper.Reaper{Cfg: reaper.DefaultConfig(), Store: st, Cluster: cl, Archiver: ar,
Now: func() time.Time { return now }}
if _, err := r.RunOnce(ctx); err != nil {
t.Fatalf("RunOnce: %v", err)
}
if len(cl.deletedPVC) != 0 || len(ar.archived) != 0 {
t.Fatalf("the reaper took a world claimed moments ago: deleted %v, archived %v", cl.deletedPVC, ar.archived)
}
// Sixteen idle days after the claim it goes, with a fresh archive.
now = time.Now().Add(16 * reaper.Day)
if _, err := r.RunOnce(ctx); err != nil {
t.Fatalf("RunOnce after 16 days: %v", err)
}
if len(ar.archived) != 1 || len(cl.deletedPVC) != 1 {
t.Fatalf("after 16 idle days: archived %v, deleted %v; want one new archive, then the delete", ar.archived, cl.deletedPVC)
}
var owner sql.NullString
var newest string
if err := db.QueryRowContext(ctx, `SELECT owner_id FROM servers WHERE name = $1`, name).Scan(&owner); err != nil {
t.Fatal(err)
}
if err := db.QueryRowContext(ctx,
`SELECT backup_ref FROM world_backups WHERE server_name = $1 ORDER BY created_at DESC LIMIT 1`, name).Scan(&newest); err != nil {
t.Fatal(err)
}
if owner.Valid || newest != ar.archived[0] {
t.Fatalf("owner = %v, newest backup = %q; want released and %q", owner, newest, ar.archived[0])
}
}
+8 -3
View File
@@ -50,9 +50,14 @@ func (s *PGStore) ListActiveServers(ctx context.Context) ([]Candidate, error) {
} }
func (s *PGStore) FreshBackup(ctx context.Context, server string, since time.Time) (Fresh, bool, error) { func (s *PGStore) FreshBackup(ctx context.Context, server string, since time.Time) (Fresh, bool, error) {
const q = `SELECT backup_ref, offsite_at IS NOT NULL FROM world_backups // Only the reaper's own archives count, and only those taken since the
WHERE server_name = $1 AND status = 'present' AND created_at >= $2 // current owner claimed the server: a manual backup may predate a panel edit
ORDER BY offsite_at IS NOT NULL DESC, created_at DESC LIMIT 1` // 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
WHERE b.server_name = $1 AND b.status = 'present' AND b.reason = 'inactive_15d' AND b.created_at >= $2
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 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.Ref, &f.Offsite); {
case err == sql.ErrNoRows: case err == sql.ErrNoRows:
+3 -2
View File
@@ -168,8 +168,9 @@ type Store interface {
// ListActiveServers returns every servers row with deleted_at IS NULL. // ListActiveServers returns every servers row with deleted_at IS NULL.
ListActiveServers(ctx context.Context) ([]Candidate, error) ListActiveServers(ctx context.Context) ([]Candidate, error)
// FreshBackup reports an existing present backup for server whose world is // FreshBackup reports an existing present reaper archive (reason
// still current — created at or after since (the world's last_active_at), // 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 // 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 // a DeletePVC failure, and across the wait for the off-site copy: the retry
// reuses the archive instead of writing a duplicate. // reuses the archive instead of writing a duplicate.
+18 -2
View File
@@ -104,6 +104,7 @@ func (c *fakeCluster) Stop(_ context.Context, name string) error {
type fakeBackup struct { type fakeBackup struct {
id, server, ref string id, server, ref string
reason string
size int64 size int64
status string // present | deleted status string // present | deleted
createdAt, expires time.Time createdAt, expires time.Time
@@ -137,7 +138,7 @@ func (s *fakeStore) ListActiveServers(context.Context) ([]Candidate, error) {
func (s *fakeStore) FreshBackup(_ context.Context, server string, since time.Time) (Fresh, bool, error) { func (s *fakeStore) FreshBackup(_ context.Context, server string, since time.Time) (Fresh, bool, error) {
var found *fakeBackup var found *fakeBackup
for _, b := range s.backups { for _, b := range s.backups {
if b.server == server && b.status == "present" && !b.createdAt.Before(since) { if b.server == server && b.status == "present" && b.reason == ReasonInactive && !b.createdAt.Before(since) {
if found == nil || (b.offsite && !found.offsite) { if found == nil || (b.offsite && !found.offsite) {
found = b found = b
} }
@@ -154,7 +155,7 @@ func (s *fakeStore) InsertBackup(_ context.Context, rec BackupRecord) error {
return s.insertErr return s.insertErr
} }
s.backups = append(s.backups, &fakeBackup{ s.backups = append(s.backups, &fakeBackup{
id: rec.ID, server: rec.ServerName, ref: rec.BackupRef, size: rec.SizeBytes, 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,
}) })
s.rec.add("insert") s.rec.add("insert")
@@ -433,6 +434,21 @@ func TestReapDeletePVCFailureIsIdempotent(t *testing.T) {
} }
} }
// A manual backup taken after the last join is not the reaper's archive: the
// owner may have edited the world from the panel since, which does not move
// last_active_at. The reap writes its own archive.
func TestReapDoesNotReuseManualBackup(t *testing.T) {
r, st, _, ar := newReaper(DefaultConfig(),
Candidate{Name: "eta", OwnerID: "user-7", LastActiveAt: idleBy(20 * Day)})
st.backups = []*fakeBackup{
{id: "man", server: "eta", ref: "ref-man", reason: "manual", size: 5, status: "present", createdAt: idleBy(10 * Day), expires: testNow.Add(80 * Day)},
}
sum := mustRun(t, r)
if sum.WorldsReaped != 1 || ar.archives != 1 {
t.Fatalf("summary = %+v, archives = %d: want the reap to archive afresh", sum, ar.archives)
}
}
// With an off-site bucket configured the reaper archives an idle world but // With an off-site bucket configured the reaper archives an idle world but
// keeps it until the archive's copy is confirmed; the next run reuses that // keeps it until the archive's copy is confirmed; the next run reuses that
// archive and deletes the world, without archiving it again. // archive and deletes the world, without archiving it again.