ec9b997151
gates / gates (push) Successful in 1m0s
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
287 lines
11 KiB
Go
287 lines
11 KiB
Go
package quiesce
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"io"
|
|
"log"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// R-921 (measured 2026-10-08 night on demo-hp): the controller stopped every app for the off-site tier,
|
|
// the agent refused the backup (`a heavy operation is already in flight`), and the apps were down about a
|
|
// minute for nothing. The rule this file pins: CHECK FIRST, STOP SECOND — and when the check could not
|
|
// see the refusal coming, RESUME AT ONCE.
|
|
//
|
|
// The agent's local API serves no read of its host-wide heavy-operation gate (felhom-agent
|
|
// internal/backup.InFlight.Busy() is served by no endpoint). What it DOES serve is each tier's job phase
|
|
// (GET /backup/status?target=…), and a job in flight on ANOTHER tier is the agent's first refusal reason
|
|
// (localapi handleBackup → otherTierInFlight → HTTP 409). So the pre-check covers that reason; the other
|
|
// (the gate held by the night OS step, a restore-test, fstrim) can only be met by the immediate resume.
|
|
|
|
// r921Backend is a tiered fake agent that records every call into a shared, ordered event log with the
|
|
// fake stacks, so a test can say "nothing happened between the refusal and the first restart".
|
|
type r921Backend struct {
|
|
mu sync.Mutex
|
|
events *[]string
|
|
tiers []BackupTier
|
|
due map[string]bool
|
|
inFlight []string // targets whose job the agent holds in flight (running|snapshotted)
|
|
probeErr error
|
|
busyOn map[string]bool // StartBackupFor answers ErrTierBusy for these targets
|
|
starts []string
|
|
}
|
|
|
|
func (b *r921Backend) log(e string) {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
*b.events = append(*b.events, e)
|
|
}
|
|
|
|
func (b *r921Backend) Tiers(context.Context) ([]BackupTier, error) { return b.tiers, nil }
|
|
func (b *r921Backend) DueFor(_ context.Context, target string) (bool, *int64, string, error) {
|
|
return b.due[target], nil, "", nil
|
|
}
|
|
func (b *r921Backend) StartBackupFor(_ context.Context, target string) (string, error) {
|
|
b.mu.Lock()
|
|
b.starts = append(b.starts, target)
|
|
b.mu.Unlock()
|
|
if b.busyOn[target] {
|
|
b.log("start:" + target + "=BUSY")
|
|
return "", busyErr()
|
|
}
|
|
b.log("start:" + target)
|
|
return "job-" + target, nil
|
|
}
|
|
func (b *r921Backend) BackupStatusFor(_ context.Context, target string) (string, error) {
|
|
b.log("status:" + target)
|
|
return phaseDone, nil
|
|
}
|
|
func (b *r921Backend) BackupJobsInFlight(context.Context) ([]string, error) {
|
|
b.log("probe")
|
|
if b.probeErr != nil {
|
|
return nil, b.probeErr
|
|
}
|
|
return append([]string(nil), b.inFlight...), nil
|
|
}
|
|
func (b *r921Backend) Due(context.Context) (bool, *int64, error) { return false, nil, nil }
|
|
func (b *r921Backend) StartBackup(context.Context) (string, error) { return "", errors.New("unused") }
|
|
func (b *r921Backend) BackupStatus(context.Context) (string, error) {
|
|
return phaseDone, nil
|
|
}
|
|
|
|
// r921Stacks logs stops and starts into the same event log.
|
|
type r921Stacks struct {
|
|
mu sync.Mutex
|
|
events *[]string
|
|
running []string
|
|
}
|
|
|
|
func (s *r921Stacks) RunningAppStacks() []string { return append([]string(nil), s.running...) }
|
|
func (s *r921Stacks) StopStack(n string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
*s.events = append(*s.events, "stop:"+n)
|
|
return nil
|
|
}
|
|
func (s *r921Stacks) StartStack(n string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
*s.events = append(*s.events, "startstack:"+n)
|
|
return nil
|
|
}
|
|
|
|
func r921Setup(t *testing.T) (*r921Backend, *r921Stacks, *Loop, *[]string, *strings.Builder) {
|
|
t.Helper()
|
|
var events []string
|
|
be := &r921Backend{
|
|
events: &events,
|
|
tiers: []BackupTier{{Target: "local", Primary: true}, {Target: "felhom-pbs"}},
|
|
due: map[string]bool{},
|
|
busyOn: map[string]bool{},
|
|
}
|
|
st := &r921Stacks{events: &events, running: []string{"immich", "paperless-ngx", "vaultwarden"}}
|
|
var logs strings.Builder
|
|
l := New(Options{
|
|
Backend: be,
|
|
Stacks: st,
|
|
MarkerPath: filepath.Join(t.TempDir(), "quiesce-state.json"),
|
|
StatusPoll: time.Millisecond,
|
|
MaxQuiesce: 30 * time.Second,
|
|
Logger: log.New(&logs, "", 0),
|
|
})
|
|
return be, st, l, &events, &logs
|
|
}
|
|
|
|
func countPrefix(events []string, prefix string) int {
|
|
n := 0
|
|
for _, e := range events {
|
|
if strings.HasPrefix(e, prefix) {
|
|
n++
|
|
}
|
|
}
|
|
return n
|
|
}
|
|
|
|
// RED TEST. Another tier's job is in flight on the agent (a first off-site upload that outlived the
|
|
// quiesce bound, or a controller restarted mid-upload) and the local tier is due. The agent WILL refuse
|
|
// the local backup (otherTierInFlight → 409). Today's code stops every app, asks, is refused, restarts
|
|
// them: a stop for nothing. Asserted on the CONSEQUENCE — no app was stopped and no backup was requested
|
|
// — not on the mechanism. It must also leave the tier due (no breaker, no contention verdict): the next
|
|
// poll asks again, which costs the household nothing.
|
|
func TestR921_NoStopWhileAnotherTiersJobIsInFlight(t *testing.T) {
|
|
be, _, l, events, logs := r921Setup(t)
|
|
be.due["local"] = true
|
|
be.inFlight = []string{"felhom-pbs"}
|
|
be.busyOn["local"] = true // what the real agent answers while another tier's job is in flight
|
|
|
|
if err := l.runOnce(context.Background()); err != nil {
|
|
t.Fatalf("runOnce: %v", err)
|
|
}
|
|
if n := countPrefix(*events, "stop:"); n != 0 {
|
|
t.Fatalf("R-921: %d app(s) stopped for a tier the agent refuses (another tier's job is in flight); events=%v", n, *events)
|
|
}
|
|
if len(be.starts) != 0 {
|
|
t.Fatalf("R-921: a backup was requested although another tier's job is in flight: %v", be.starts)
|
|
}
|
|
if got := l.breaker.failuresFor("local"); got != 0 {
|
|
t.Errorf("the pre-check armed the failure breaker (%d) — nothing failed", got)
|
|
}
|
|
if _, blocked := l.contention.blocked("local", l.now()); blocked {
|
|
t.Errorf("the pre-check set a contention verdict — the tier must simply stay due and be asked again next poll")
|
|
}
|
|
if !strings.Contains(logs.String(), "felhom-pbs") || !strings.Contains(logs.String(), "R-921") {
|
|
t.Errorf("the skip is not explained in the log (want the in-flight tier and R-921): %q", logs.String())
|
|
}
|
|
|
|
// The job ends → the next poll backs the local tier up as normal (the skip did not swallow the tier).
|
|
be.inFlight = nil
|
|
be.busyOn["local"] = false
|
|
if err := l.runOnce(context.Background()); err != nil {
|
|
t.Fatalf("runOnce 2: %v", err)
|
|
}
|
|
if len(be.starts) != 1 || be.starts[0] != "local" {
|
|
t.Fatalf("after the job ended the local tier was not backed up: starts=%v", be.starts)
|
|
}
|
|
}
|
|
|
|
// The measured instance and the race, both on the BUSY path: the pre-check cannot see the refusal
|
|
// coming (the night OS step holds the gate — not readable from the agent — or the agent became busy
|
|
// between the check and the start). Then the resume must follow the refusal AT ONCE: the very next
|
|
// event after the refused start is the first app's restart — no other agent call, no further tier, no
|
|
// wait in between. Measured on the ordered event log, not on a clock.
|
|
func TestR921_BusyRefusalResumesAtOnce(t *testing.T) {
|
|
for _, tc := range []struct {
|
|
name string
|
|
due []string
|
|
busy []string
|
|
probe bool // the pre-check answered "free" (the race) vs. a backend without the pre-check
|
|
}{
|
|
{"measured: only the off-site tier due, refused (no pre-check)", []string{"felhom-pbs"}, []string{"felhom-pbs"}, false},
|
|
{"race: pre-check said free, the start was refused", []string{"felhom-pbs"}, []string{"felhom-pbs"}, true},
|
|
} {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
be, st, l, events, _ := r921Setup(t)
|
|
for _, d := range tc.due {
|
|
be.due[d] = true
|
|
}
|
|
for _, b := range tc.busy {
|
|
be.busyOn[b] = true
|
|
}
|
|
var loop *Loop = l
|
|
if !tc.probe {
|
|
// A backend WITHOUT the pre-check surface (an adapter that does not implement it).
|
|
nb := &noProbeBackend{be}
|
|
loop = New(Options{Backend: nb, Stacks: st, MarkerPath: filepath.Join(t.TempDir(), "m.json"),
|
|
StatusPoll: time.Millisecond, MaxQuiesce: 30 * time.Second, Logger: log.New(io.Discard, "", 0)})
|
|
}
|
|
if err := loop.runOnce(context.Background()); err != nil {
|
|
t.Fatalf("runOnce: %v", err)
|
|
}
|
|
ev := *events
|
|
idx := -1
|
|
for i, e := range ev {
|
|
if strings.HasSuffix(e, "=BUSY") {
|
|
idx = i
|
|
}
|
|
}
|
|
if idx < 0 {
|
|
t.Fatalf("no refused start in the events: %v", ev)
|
|
}
|
|
if idx+1 >= len(ev) || !strings.HasPrefix(ev[idx+1], "startstack:") {
|
|
t.Fatalf("the resume did not follow the refusal at once; events after the refusal: %v", ev[idx+1:])
|
|
}
|
|
if got, want := countPrefix(ev[idx+1:], "startstack:"), len(st.running); got != want {
|
|
t.Fatalf("restarted %d app(s) after the refusal, want all %d; events=%v", got, want, ev)
|
|
}
|
|
for _, e := range ev[idx+1:] {
|
|
if !strings.HasPrefix(e, "startstack:") {
|
|
t.Fatalf("an agent call (%q) came after the refusal — the BUSY path must only resume; events=%v", e, ev)
|
|
}
|
|
}
|
|
if _, blocked := loop.contention.blocked("felhom-pbs", loop.now()); !blocked {
|
|
t.Errorf("the refused tier was not deferred — the next poll would stop the apps again")
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// A pre-check that cannot be answered must not suppress a backup: fail toward backing up, exactly as
|
|
// before R-921 (stop, ask, and on refusal resume at once).
|
|
func TestR921_PrecheckErrorFailsTowardBackingUp(t *testing.T) {
|
|
be, _, l, events, logs := r921Setup(t)
|
|
be.due["local"] = true
|
|
be.probeErr = errors.New("agentapi: GET /backup/tiers: connection refused")
|
|
|
|
if err := l.runOnce(context.Background()); err != nil {
|
|
t.Fatalf("runOnce: %v", err)
|
|
}
|
|
if len(be.starts) != 1 || be.starts[0] != "local" {
|
|
t.Fatalf("an unanswerable pre-check suppressed the backup: starts=%v events=%v", be.starts, *events)
|
|
}
|
|
if !strings.Contains(logs.String(), "[WARN]") {
|
|
t.Errorf("the failed pre-check was not logged at WARN: %q", logs.String())
|
|
}
|
|
}
|
|
|
|
// The window's OWN tier already in flight is NOT skipped: the agent answers that start with the running
|
|
// job (202, idempotent), the loop attaches to it and records its outcome — unchanged behaviour, chosen
|
|
// deliberately (only a tier the agent REFUSES is skipped before the stop).
|
|
func TestR921_SameTierInFlightIsNotSkipped(t *testing.T) {
|
|
be, _, l, _, _ := r921Setup(t)
|
|
be.due["felhom-pbs"] = true
|
|
be.inFlight = []string{"felhom-pbs"}
|
|
|
|
if err := l.runOnce(context.Background()); err != nil {
|
|
t.Fatalf("runOnce: %v", err)
|
|
}
|
|
if len(be.starts) != 1 || be.starts[0] != "felhom-pbs" {
|
|
t.Fatalf("the window's own in-flight tier was skipped; starts=%v", be.starts)
|
|
}
|
|
}
|
|
|
|
// noProbeBackend hides BackupJobsInFlight (an adapter without the pre-check surface).
|
|
type noProbeBackend struct{ b *r921Backend }
|
|
|
|
func (n *noProbeBackend) Tiers(ctx context.Context) ([]BackupTier, error) { return n.b.Tiers(ctx) }
|
|
func (n *noProbeBackend) DueFor(ctx context.Context, t string) (bool, *int64, string, error) {
|
|
return n.b.DueFor(ctx, t)
|
|
}
|
|
func (n *noProbeBackend) StartBackupFor(ctx context.Context, t string) (string, error) {
|
|
return n.b.StartBackupFor(ctx, t)
|
|
}
|
|
func (n *noProbeBackend) BackupStatusFor(ctx context.Context, t string) (string, error) {
|
|
return n.b.BackupStatusFor(ctx, t)
|
|
}
|
|
func (n *noProbeBackend) Due(ctx context.Context) (bool, *int64, error) { return n.b.Due(ctx) }
|
|
func (n *noProbeBackend) StartBackup(ctx context.Context) (string, error) {
|
|
return n.b.StartBackup(ctx)
|
|
}
|
|
func (n *noProbeBackend) BackupStatus(ctx context.Context) (string, error) {
|
|
return n.b.BackupStatus(ctx)
|
|
}
|