package backup import ( "context" "log/slog" "time" "gitea.dooplex.hu/admin/felhom-agent/internal/reconcile" ) // RestoreTestRunner is the reconcile-engine seam the scheduler drives (*reconcile.Engine // satisfies it). Kept narrow so the scheduler is unit-testable with a fake. type RestoreTestRunner interface { RunRestoreTest(ctx context.Context, spec reconcile.RestoreTestSpec) reconcile.RestoreTestResult } // CandidatePicker resolves the archive volid to restore-test (newest backup), or "" when // there is none yet (the tick then no-ops). type CandidatePicker func(ctx context.Context) (string, error) // SpecBuilder yields the RestoreTestSpec for ONE run, given the archive that was picked. // // R-85 (1.1): this REPLACES a frozen spec value. It used to be built by an immediately-invoked // function at daemon start, so `storageTier()` and `restoreTaskTimeout()` were evaluated ONCE and // the resulting value reused for every run for the lifetime of the process. Two consequences: // - nothing tier-varying was expressible at all (the offsite tier could never be scheduled), and // - it was a latent staleness bug in its own right — a storage-type or config change did not take // effect until the daemon restarted. // // The archive is passed in because the tier MUST be derived from it (the v0.100.0 rule), never from // the configured target: deriving it from config is what produced the 600 s false failure when a // PBS archive was classified "local" and got the 10-minute local wait. type SpecBuilder func(ctx context.Context, archive string) reconcile.RestoreTestSpec // TierPicker resolves the newest archive on a NAMED tier, or "" when that tier holds none. // (*BackupRunner).PickRestoreCandidateOn satisfies it. "" must NOT be an error — a brand-new // offsite tier legitimately has nothing to restore yet. type TierPicker func(ctx context.Context, target string) (string, error) // Scheduler runs the self-restore-test on an agent-internal cadence. It is the fourth daemon // goroutine; it does real restore→boot→destroy, so it only runs when the cadence is enabled // AND a valid scratch band is configured (validated by the caller before construction). type Scheduler struct { runner RestoreTestRunner pick CandidatePicker store *Store spec SpecBuilder // R-85: evaluated PER RUN, never frozen at construction cadence time.Duration logger *slog.Logger now func() time.Time // R-85 tier rotation. All optional: without them the scheduler behaves exactly as before // (single tier via `pick`), which keeps every existing caller and test working untouched. tiers []string // configured tier target ids, primary first tierPick TierPicker // newest archive on a named tier rtState *RestoreTestState // persisted last-successful-per-tier (drives oldest-first) inFlight *InFlight // shared with the backup path — Scenario F } // SchedulerOptions configures a Scheduler. type SchedulerOptions struct { Runner RestoreTestRunner Pick CandidatePicker Store *Store // Spec builds the run's spec (RestoreStorage, ScratchMin/Max, SourceTier, timeouts) from the // picked archive. Called ONCE PER RUN — see SpecBuilder for why it is not a value. Spec SpecBuilder Cadence time.Duration // 0 → disabled Logger *slog.Logger // R-85 (all optional — omit for the pre-R-85 single-tier behaviour): // Tiers are the configured tier target ids (primary first); TierPick resolves an archive on a // named tier; State persists last-successful-per-tier; InFlight is the shared one-heavy-op gate. Tiers []string TierPick TierPicker State *RestoreTestState InFlight *InFlight } // NewScheduler builds a Scheduler. func NewScheduler(opts SchedulerOptions) *Scheduler { logger := opts.Logger if logger == nil { logger = slog.Default() } return &Scheduler{ runner: opts.Runner, pick: opts.Pick, store: opts.Store, spec: opts.Spec, cadence: opts.Cadence, logger: logger, now: func() time.Time { return time.Now().UTC() }, tiers: append([]string(nil), opts.Tiers...), tierPick: opts.TierPick, rtState: opts.State, inFlight: opts.InFlight, } } // Run fires a restore-test on the cadence until ctx is cancelled. A 0 cadence disables it // (the goroutine just waits for shutdown). It does NOT fire immediately on start (a restore // is heavy; the first runs one interval in) — on-demand runs use the selftest harness. // Returns nil on ctx cancellation. func (s *Scheduler) Run(ctx context.Context) error { if s.cadence <= 0 || s.runner == nil || s.spec == nil || (s.pick == nil && !s.rotating()) { s.logger.Info("backup: restore-test cadence disabled") <-ctx.Done() return nil } s.logger.Info("backup: restore-test scheduler starting", "cadence", s.cadence) t := time.NewTicker(s.cadence) defer t.Stop() for { select { case <-ctx.Done(): s.logger.Info("backup: restore-test scheduler shutting down", "reason", ctx.Err()) return nil case <-t.C: s.tick(ctx) } } } // tick runs one scheduled restore-test: pick a backup → run → record. No-ops cleanly when // no backup exists yet. Deterministic given s.now — tests call it directly. func (s *Scheduler) tick(ctx context.Context) { if s.spec == nil { // Defensive: Run() already refuses to start without a SpecBuilder, but tick is also // reachable directly. Skipping loudly beats panicking the daemon goroutine — a missing // spec must cost a restore-test, never the agent. s.logger.Error("backup: restore-test has no spec builder — skipping (this is a wiring bug)") return } // Scenario F: join the one-heavy-operation-at-a-time gate. A restore-test PULLS a multi-GB // archive over the same tunnel an offsite backup PUSHES one; running both saturates the link and // drives each toward its timeout, which is how a healthy tier gets recorded as failed. DEFER — // never cancel what is already running: a deferred restore-test costs hours of coverage, a // cancelled backup costs the backup. release, busy, ok := s.inFlight.TryAcquire("restore-test") if !ok { s.logger.Info("backup: restore-test deferred — a heavy operation is already in flight", "busy", busy) return } defer release() archive, target, err := s.pickForThisRun(ctx) if err != nil { s.logger.Warn("backup: restore-test could not pick a candidate; skipping", "err", err) return } if archive == "" { s.logger.Info("backup: restore-test skipped; no backup available yet") return } // R-85: build the spec for THIS run, from THIS archive. Never a frozen value. spec := s.spec(ctx, archive) spec.Archive = archive res := s.runner.RunRestoreTest(ctx, spec) if res.Skipped { return // already logged by the engine (no free scratch VMID) } rt := ToHubRestoreTest(res, s.now()) s.store.RecordRestoreTest(rt) // Rotation credit is given ONLY on success. A failing tier must keep sorting first, or a tier // that fails every time would look freshly proven and quietly stop being retried. if rt.Pass && s.rtState != nil && target != "" { if err := s.rtState.RecordSuccess(target, s.now()); err != nil { s.logger.Warn("backup: could not persist the restore-test rotation state", "target", target, "err", err) } } switch { case !rt.Pass: // A failing restore-test is the loudest DR signal there is. s.logger.Error("backup: scheduled restore-test FAILED", "archive", rt.SourceArchive, "err", rt.Error) case len(res.StartWarnings) == 0: s.logger.Info("backup: scheduled restore-test passed", "archive", rt.SourceArchive, "duration_s", rt.DurationSeconds) case res.WarningsRecognized: // Passed; the only warnings are the known-benign (e.g. systemd-nesting) advisory. s.logger.Info("backup: scheduled restore-test passed with warnings (recognized)", "archive", rt.SourceArchive, "duration_s", rt.DurationSeconds, "warnings", res.StartWarnings) default: // Passed liveness, but an UNRECOGNIZED start warning stood out — worth an operator look. s.logger.Warn("backup: scheduled restore-test passed with UNRECOGNIZED warnings", "archive", rt.SourceArchive, "duration_s", rt.DurationSeconds, "warnings", res.StartWarnings) } } // rotating reports whether multi-tier rotation is wired. func (s *Scheduler) rotating() bool { return len(s.tiers) > 0 && s.tierPick != nil } // pickForThisRun chooses the tier and its newest archive. // // OLDEST-FIRST (operator ruling 2026-07-26, Option 1): the tier whose last SUCCESSFUL restore-test // is oldest goes first, never-proven first of all. Self-balancing, no config knob, and it naturally // prioritises a tier that has never been proven — which on this fleet was the offsite tier, unproven // for its entire existence while reporting `applied`. // // A tier with no archives is SKIPPED, not failed, and the next tier is tried. Skipping to a testable // tier is strictly better than burning the whole cadence: a brand-new offsite tier has nothing to // restore yet, and that is normal, not broken. It cannot starve the empty tier either — as soon as // it has an archive it still sorts first, because it is still the least recently proven. // // Returns ("", "", nil) when nothing anywhere is testable. func (s *Scheduler) pickForThisRun(ctx context.Context) (archive, target string, err error) { if !s.rotating() { a, perr := s.pick(ctx) return a, "", perr // pre-R-85 single-tier path; no rotation credit to record } order := s.tiers if s.rtState != nil { order = s.rtState.OldestFirst(s.tiers) } var firstErr error for _, t := range order { a, perr := s.tierPick(ctx, t) if perr != nil { // One tier's storage being unreadable must not block the others. s.logger.Warn("backup: restore-test candidate lookup failed for a tier; trying the next", "target", t, "err", perr) if firstErr == nil { firstErr = perr } continue } if a == "" { s.logger.Debug("backup: restore-test tier has no archive yet; trying the next", "target", t) continue } s.logger.Info("backup: restore-test tier selected (oldest-proven first)", "target", t, "archive", a) return a, t, nil } if firstErr != nil { return "", "", firstErr } return "", "", nil }