Files
felhom-controller/controller/internal/quiesce/r921_precheck_test.go
T

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)
}