From e06680b79794f0cd5d5a1dd953f04758c0de371a Mon Sep 17 00:00:00 2001 From: kisfenyo Date: Thu, 8 Oct 2026 14:38:39 +0200 Subject: [PATCH] Operator actions in the report reply (D1, R-314/R-279/R-177, decision 185) The hub may ask the running controller for a CLOSED list of actions, carried in the report ACK (operator_actions) and answered on the next report (operator_action_results): offsite_backup_now, abandon_stop, abandon_extend (1-30 days), run_job (fill-watch, offsite-integrity, offsite-proof, disk-health-check). Unknown action/job/argument -> refused, nothing called. Once per id (in memory; every action is safe to repeat). - internal/report/opactions.go: the executor; results re-sent until the hub stops listing the id. - scheduler.RunNow: refuses unknown / already-running jobs; OnDemand(ctx) makes an operator's offsite-integrity run even when not due. - ExtendAbandon never shortens the countdown and refuses in the hub phase; StopAbandon reports failure when the hub cancel failed (it said success). - Report ACK read cap 4 KiB -> 64 KiB (an ACK over the cap dropped every field in it). Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS --- controller/cmd/controller/main.go | 16 +- controller/cmd/controller/opactions_wiring.go | 60 +++ .../cmd/controller/opactions_wiring_test.go | 169 +++++++++ controller/internal/backup/offbox_abandon.go | 22 ++ .../backup/offbox_abandon_levers_d1_test.go | 85 +++++ controller/internal/report/builder.go | 3 + controller/internal/report/opactions.go | 345 ++++++++++++++++++ controller/internal/report/opactions_test.go | 193 ++++++++++ controller/internal/report/pusher.go | 12 +- controller/internal/report/types.go | 5 + controller/internal/scheduler/runnow_test.go | 76 ++++ controller/internal/scheduler/scheduler.go | 74 +++- 12 files changed, 1055 insertions(+), 5 deletions(-) create mode 100644 controller/cmd/controller/opactions_wiring.go create mode 100644 controller/cmd/controller/opactions_wiring_test.go create mode 100644 controller/internal/backup/offbox_abandon_levers_d1_test.go create mode 100644 controller/internal/report/opactions.go create mode 100644 controller/internal/report/opactions_test.go create mode 100644 controller/internal/scheduler/runnow_test.go diff --git a/controller/cmd/controller/main.go b/controller/cmd/controller/main.go index 2867c89..f506fac 100644 --- a/controller/cmd/controller/main.go +++ b/controller/cmd/controller/main.go @@ -1077,6 +1077,15 @@ func main() { RecordEscrowKeyHash: sett.SetHubEscrowKeySHA256, Logger: logger, } + // R-314/R-279/R-177 (`09` §3 decision 185): the operator's closed list of actions, carried in the + // report reply and run in THIS process (opactions_wiring.go). The result rides the next report; + // a finished result fires an out-of-cycle report through the trigger once it exists. + opActions := report.NewOperatorActions(ctx, operatorActionHandlers(backupMgr, sched, func() { + if f, ok := opActionReportFire.Load().(func()); ok && f != nil { + f() + } + }, logger), logger) + report.SetOperatorActions(opActions) // Wire hub verification: update settings when hub reports customer status hubPusher.OnPushResponse = func(resp *report.PushResponse) { if resp.CustomerBlocked { @@ -1120,6 +1129,8 @@ func main() { // generation). The web gate reads it on the next request — no restart needed. claimSync := &report.ClaimSync{Settings: sett, Logger: logger} claimSync.Reconcile(resp.Claim) + // Decision 185: the operator's pending actions (each acted on once per id). + opActions.Reconcile(resp.OperatorActions) } // Wire hub push status into alert manager for dashboard alerts alertMgr.SetHubPushStatus(func() web.HubPushStatusData { @@ -1496,7 +1507,9 @@ func main() { // collision costs one skipped day, not a missed check, because due-ness makes tomorrow try // again. What would be a real defect is a slot that collides EVERY night; this is not one. sched.Daily("offsite-integrity", "06:00", func(ctx context.Context) error { - runOffsiteIntegrityCheck(ctx, backupMgr, notifier, logger, false) + // Decision 185: an operator's run_job runs the check even when it is not due — the operator + // asked for a check, not for „is a check due". The schedule still checks due-ness. + runOffsiteIntegrityCheck(ctx, backupMgr, notifier, logger, scheduler.OnDemand(ctx)) return nil }) @@ -1811,6 +1824,7 @@ func main() { } reportTrigger = report.NewTrigger(fireReport, logger) go reportTrigger.Run(ctx) + opActionReportFire.Store(reportTrigger.Fire) // Direction-2 immediate-sync (v0.140.0): hold a hanging GET against the hub's wait channel // and fire the trigger the instant operator intent moves — so an operator save round-trips in diff --git a/controller/cmd/controller/opactions_wiring.go b/controller/cmd/controller/opactions_wiring.go new file mode 100644 index 0000000..4532741 --- /dev/null +++ b/controller/cmd/controller/opactions_wiring.go @@ -0,0 +1,60 @@ +package main + +import ( + "context" + "errors" + "fmt" + "log" + "sync/atomic" + "time" + + "gitea.dooplex.hu/admin/felhom-controller/internal/backup" + "gitea.dooplex.hu/admin/felhom-controller/internal/report" + "gitea.dooplex.hu/admin/felhom-controller/internal/scheduler" +) + +// opActionReportFire holds the out-of-cycle report trigger once main has built it (func()). The +// executor is built earlier than the trigger, so it reads it through this holder. +var opActionReportFire atomic.Value + +// operatorActionHandlers wires the four operator actions (R-314/R-279/R-177, `09` §3 decision 185) +// to the running process. ONE wiring site, called from main and from the wiring test, so the test +// exercises the doors production uses (the seam-built-but-never-wired trap). +// +// offsite_backup_now makes the SAME four checks as the household's „back up now" +// (offboxRunHandler) and then runs the same call, so an operator press can never do what the +// household's own press cannot. +func operatorActionHandlers(backupMgr *backup.Manager, sched *scheduler.Scheduler, resultReady func(), logger *log.Logger) report.OperatorActionHandlers { + h := report.OperatorActionHandlers{ResultReady: resultReady} + if sched != nil { + h.RunJob = sched.RunNow + } + if backupMgr == nil { + return h + } + h.OffsiteBackupNow = func() (string, func(ctx context.Context) error) { + switch { + case !backupMgr.OffboxConfigured(): + return "no off-site target is set on this box", nil + case !backupMgr.OffboxRunnable(): + return "the off-site key step (escrow) is not done yet — the run is refused until it is", nil + case backupMgr.OffboxOrphaned(): + return "the off-site store is orphaned (written with an older key) — a run cannot write to it", nil + case backupMgr.IsRunning(): + return "a backup run is already in flight", nil + } + return "", func(ctx context.Context) error { + if logger != nil { + logger.Printf("[INFO] [opaction] off-site backup started at the operator's request") + } + err := backupMgr.RunOffboxBackupWithProgress(ctx) + if errors.Is(err, backup.ErrOffboxRunInFlight) { + return fmt.Errorf("%w: %v", report.ErrActionRefused, err) + } + return err + } + } + h.AbandonStop = backupMgr.StopAbandon + h.AbandonExtend = func(days int) (time.Time, error) { return backupMgr.ExtendAbandon(days) } + return h +} diff --git a/controller/cmd/controller/opactions_wiring_test.go b/controller/cmd/controller/opactions_wiring_test.go new file mode 100644 index 0000000..11c7178 --- /dev/null +++ b/controller/cmd/controller/opactions_wiring_test.go @@ -0,0 +1,169 @@ +package main + +import ( + "context" + "go/ast" + "go/parser" + "go/token" + "io" + "log" + "path/filepath" + "testing" + "time" + + "gitea.dooplex.hu/admin/felhom-controller/internal/backup" + "gitea.dooplex.hu/admin/felhom-controller/internal/config" + "gitea.dooplex.hu/admin/felhom-controller/internal/report" + "gitea.dooplex.hu/admin/felhom-controller/internal/scheduler" + "gitea.dooplex.hu/admin/felhom-controller/internal/settings" +) + +// `09` §3 decision 185 (D1), the design's red test (1): a reply carrying abandon_stop during a +// countdown stops it IN THE RUNNING MANAGER — and settings.json on disk agrees. This is the reason +// the action rides the report reply instead of the CLI lever: the CLI runs a SECOND process, whose +// write the running controller can overwrite (the lost update, main.go's R-241 block). +func TestD1_AbandonStopActsOnTheRunningManagerAndDisk(t *testing.T) { + logger := log.New(io.Discard, "", 0) + dir := t.TempDir() + settingsPath := filepath.Join(dir, "settings.json") + sett, err := settings.Load(settingsPath, logger) + if err != nil { + t.Fatal(err) + } + cfg := &config.Config{} + cfg.Paths.DataDir = dir + cfg.Paths.SystemDataPath = filepath.Join(dir, "sys") + mgr := backup.NewManager(cfg, sett, logger) + // No test reaches a real ssh or restic (R-488 / R-650). + mgr.SetOffboxSSH(func(context.Context, string, string, int, string, string, string) ([]byte, error) { + t.Fatal("abandon_stop must issue no remote command") + return nil, nil + }) + mgr.SetOffboxRunner(func(context.Context, []string, ...string) ([]byte, error) { + t.Fatal("abandon_stop must run no restic") + return nil, nil + }) + due := time.Now().Add(10 * 24 * time.Hour).UTC() + if err := sett.SetOffboxTarget(&settings.OffboxTarget{ + Enabled: true, Host: "nas.local", Port: 22, User: "felhom", RepoPath: "/srv/repo", Schedule: "daily", + EscrowState: "escrowed", + AbandonRepoPath: "/srv/repo.orphaned-20261001", + AbandonStartedAt: time.Now().Add(-4 * 24 * time.Hour).UTC().Format(time.RFC3339), + AbandonAt: due.Format(time.RFC3339), + }); err != nil { + t.Fatal(err) + } + if !mgr.AbandonStatus().Active { + t.Fatal("fixture: a countdown must be running") + } + + o := report.NewOperatorActions(context.Background(), operatorActionHandlers(mgr, nil, nil, logger), logger) + o.Reconcile([]report.OperatorAction{{ID: 1, Action: report.OpAbandonStop}}) + + if mgr.AbandonStatus().Active { + t.Fatal("the running manager still counts down after abandon_stop") + } + onDisk, err := settings.Load(settingsPath, logger) + if err != nil { + t.Fatal(err) + } + if tg := onDisk.GetOffboxTarget(); tg == nil || tg.AbandonAt != "" || tg.AbandonRepoPath == "" { + t.Fatalf("settings.json disagrees: %+v (want no due date, the set-aside path kept)", tg) + } + rs := o.Results() + if len(rs) != 1 || rs[0].Outcome != report.OutcomeDone { + t.Fatalf("result = %+v", rs) + } + + // Control: the same countdown, a second press → nothing to stop, reported as such (not „done"). + o.Reconcile([]report.OperatorAction{{ID: 2, Action: report.OpAbandonStop}}) + if r := o.Results(); len(r) != 1 || r[0].ID != 2 || r[0].Outcome != report.OutcomeFailed { + t.Fatalf("second stop: %+v", r) + } +} + +// offsite_backup_now refuses exactly where the household's „back up now" refuses: no target here. +func TestD1_OffsiteBackupNowRefusesWithoutATarget(t *testing.T) { + logger := log.New(io.Discard, "", 0) + dir := t.TempDir() + sett, err := settings.Load(filepath.Join(dir, "settings.json"), logger) + if err != nil { + t.Fatal(err) + } + cfg := &config.Config{} + cfg.Paths.DataDir = dir + mgr := backup.NewManager(cfg, sett, logger) + mgr.SetOffboxRunner(func(context.Context, []string, ...string) ([]byte, error) { + t.Fatal("no restic may run") + return nil, nil + }) + o := report.NewOperatorActions(context.Background(), operatorActionHandlers(mgr, nil, nil, logger), logger) + o.Reconcile([]report.OperatorAction{{ID: 9, Action: report.OpOffsiteBackupNow}}) + if r := o.Results(); len(r) != 1 || r[0].Outcome != report.OutcomeRefused { + t.Fatalf("result = %+v", r) + } +} + +// run_job reaches the scheduler's RunNow: a registered job runs; a job not registered on this box is +// refused. +func TestD1_RunJobUsesTheScheduler(t *testing.T) { + logger := log.New(io.Discard, "", 0) + sched := scheduler.New(logger) + ran := make(chan bool, 1) + sched.Daily("fill-watch", "03:30", func(ctx context.Context) error { ran <- scheduler.OnDemand(ctx); return nil }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + sched.Start(ctx) + defer sched.Stop() + o := report.NewOperatorActions(ctx, operatorActionHandlers(nil, sched, nil, logger), logger) + o.Reconcile([]report.OperatorAction{{ID: 1, Action: report.OpRunJob, Arg: "fill-watch"}, {ID: 2, Action: report.OpRunJob, Arg: "offsite-proof"}}) + select { + case od := <-ran: + if !od { + t.Fatal("the job did not see itself as on-demand") + } + case <-time.After(5 * time.Second): + t.Fatal("fill-watch never ran") + } + deadline := time.Now().Add(5 * time.Second) + for len(o.Results()) < 2 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + got := map[int64]string{} + for _, r := range o.Results() { + got[r.ID] = r.Outcome + } + if got[1] != report.OutcomeDone || got[2] != report.OutcomeRefused { + t.Fatalf("outcomes = %v, want 1:done 2:refused (offsite-proof not registered here)", got) + } +} + +// main.go hands the reply's list to the executor and the trigger to its holder — one call each +// (the seam-built-but-never-wired trap: a green executor that main never calls proves nothing). +func TestD1_MainWiresTheExecutor(t *testing.T) { + f, err := parser.ParseFile(token.NewFileSet(), "main.go", nil, 0) + if err != nil { + t.Fatal(err) + } + found := map[string]int{} + ast.Inspect(f, func(n ast.Node) bool { + call, ok := n.(*ast.CallExpr) + if !ok { + return true + } + switch fn := call.Fun.(type) { + case *ast.SelectorExpr: + if x, ok := fn.X.(*ast.Ident); ok { + found[x.Name+"."+fn.Sel.Name]++ + } + case *ast.Ident: + found[fn.Name]++ + } + return true + }) + for _, want := range []string{"opActions.Reconcile", "report.SetOperatorActions", "operatorActionHandlers", "opActionReportFire.Store"} { + if found[want] != 1 { + t.Errorf("main.go calls %s %d time(s), want exactly 1", want, found[want]) + } + } +} diff --git a/controller/internal/backup/offbox_abandon.go b/controller/internal/backup/offbox_abandon.go index 5837b12..63f8d93 100644 --- a/controller/internal/backup/offbox_abandon.go +++ b/controller/internal/backup/offbox_abandon.go @@ -356,7 +356,21 @@ func (m *Manager) ExtendAbandon(days int) (time.Time, error) { } return time.Time{}, fmt.Errorf("no abandonment countdown is running on this box — nothing to extend") } + // Decision 74: once the deletion is the HUB's, a box-side date changes the page and not the hub's + // schedule — the hub would still delete on its own date. Stop it instead (StopAbandon cancels it + // at the hub). Pinned by TestD1_ExtendRefusedWhileTheDeletionIsTheHubs. + if st.HubPending { + return time.Time{}, fmt.Errorf("the deletion is already pending at the hub (due %s) — it cannot be extended from the box; stop it instead", + st.HubDueAt.Format("2006-01-02")) + } due := m.abandonNow().UTC().AddDate(0, 0, days) + // `09` §3 decision 185: an extension never moves the deletion EARLIER. „N days from now" can land + // before today's due date (day 1 of 14, extend by 3); that would shorten the countdown, which no + // lever may do. Pinned by TestD1_ExtendNeverShortensTheCountdown. + if !due.After(st.DueAt) { + return time.Time{}, fmt.Errorf("%d day(s) from now (%s) is not later than the current deletion date %s — that would shorten the countdown", + days, due.Format("2006-01-02"), st.DueAt.Format("2006-01-02")) + } if err := m.settings.UpdateOffboxStatus(func(o *settings.OffboxTarget) { o.AbandonAt = due.Format(time.RFC3339) }); err != nil { @@ -380,5 +394,13 @@ func (m *Manager) StopAbandon() error { return fmt.Errorf("no abandonment countdown is running on this box — nothing to stop") } m.CancelAbandon("stopped by an operator") + // CancelAbandon keeps the countdown when the hub's pending deletion could not be cancelled (so the + // sweep retries). Report that as a failure — an operator told „stopped" while the hub still deletes + // on schedule is the comfort this project keeps removing. Pinned by + // TestD1_StopReportsFailureWhenTheHubCancelFails. + if after := m.AbandonStatus(); after.Active { + return fmt.Errorf("the countdown could not be stopped (the hub's pending deletion was not cancelled; see the log) — it is still due %s", + after.DueAt.Format("2006-01-02")) + } return nil } diff --git a/controller/internal/backup/offbox_abandon_levers_d1_test.go b/controller/internal/backup/offbox_abandon_levers_d1_test.go new file mode 100644 index 0000000..9bcdedd --- /dev/null +++ b/controller/internal/backup/offbox_abandon_levers_d1_test.go @@ -0,0 +1,85 @@ +package backup + +import ( + "context" + "errors" + "testing" + "time" + + "gitea.dooplex.hu/admin/felhom-controller/internal/settings" +) + +// `09` §3 decision 185 (D1): the operator's abandon_extend and abandon_stop now arrive from the hub, +// with no shell and no human reading the CLI's printout. Three guarantees the CLI lever did not give, +// each pinned here because the hub action relies on it: +// 1. an „extend" never moves the deletion EARLIER (decision 185: no action may shorten a countdown); +// 2. an „extend" is refused while the deletion is the HUB's (decision 74) — the box's own date would +// change on the page while the hub still deleted on its date; +// 3. a „stop" that could not cancel the hub's pending deletion reports a failure, never success. + +func TestD1_ExtendNeverShortensTheCountdown(t *testing.T) { + start := time.Date(2026, 8, 7, 12, 0, 0, 0, time.UTC) + m, sett, _ := abandonFixture(t, start) + if err := m.ResetOrphanedRepo(context.Background()); err != nil { + t.Fatal(err) + } + before := m.AbandonStatus().DueAt // start + 14 days + // Day 1: „extend by 3 days" would mean day 4 — ten days EARLIER than today's due date. + m.SetOffboxClock(func() time.Time { return start.AddDate(0, 0, 1) }) + if _, err := m.ExtendAbandon(3); err == nil { + t.Fatal("an extension that moves the deletion earlier was accepted") + } + if got := m.AbandonStatus().DueAt; !got.Equal(before) { + t.Fatalf("a refused extension changed the due date: %v -> %v", before, got) + } + if got := sett.GetOffboxTarget().AbandonAt; got != before.Format(time.RFC3339) { + t.Fatalf("settings.json due date moved: %q", got) + } + // Control: a real extension (later than today's due date) is accepted. + if due, err := m.ExtendAbandon(20); err != nil || !due.After(before) { + t.Fatalf("a real extension: due=%v err=%v", due, err) + } +} + +func TestD1_ExtendRefusedWhileTheDeletionIsTheHubs(t *testing.T) { + m, sett := newOffboxManager(t) + pinTarget(t, sett) + hubDue := time.Now().Add(5 * 24 * time.Hour).UTC() + if err := sett.UpdateOffboxStatus(func(o *settings.OffboxTarget) { + o.AbandonRepoPath = "/home/felhom-repo.orphaned-20260901" + o.AbandonHubDueAt = hubDue.Format(time.RFC3339) + }); err != nil { + t.Fatal(err) + } + if !m.AbandonStatus().HubPending { + t.Fatal("fixture: expected a hub-phase deletion") + } + if _, err := m.ExtendAbandon(30); err == nil { + t.Fatal("an extension during the hub phase was accepted — the hub would still delete on its own date") + } + if got := sett.GetOffboxTarget().AbandonAt; got != "" { + t.Fatalf("the refused extension wrote a box-side date %q the page would show", got) + } +} + +type failingCancelAbandon struct{ fakeAbandon } + +func (f *failingCancelAbandon) Cancel(context.Context) error { return errors.New("hub unreachable") } + +func TestD1_StopReportsFailureWhenTheHubCancelFails(t *testing.T) { + m, sett := newOffboxManager(t) + pinTarget(t, sett) + if err := sett.UpdateOffboxStatus(func(o *settings.OffboxTarget) { + o.AbandonRepoPath = "/home/felhom-repo.orphaned-20260901" + o.AbandonHubDueAt = time.Now().Add(5 * 24 * time.Hour).UTC().Format(time.RFC3339) + }); err != nil { + t.Fatal(err) + } + m.SetOffsiteAbandonClient(&failingCancelAbandon{}) + if err := m.StopAbandon(); err == nil { + t.Fatal("StopAbandon returned success while the hub's deletion is still pending") + } + if !m.AbandonStatus().Active { + t.Fatal("fixture: the countdown must still be running after a failed cancel") + } +} diff --git a/controller/internal/report/builder.go b/controller/internal/report/builder.go index 4e45abd..0754b5b 100644 --- a/controller/internal/report/builder.go +++ b/controller/internal/report/builder.go @@ -173,6 +173,9 @@ func BuildReport( // consume-once ACK-flag pattern (selftail.go). nil in the steady state. r.ControllerLogTail = buildControllerLogTail(logger) + // Operator actions (decision 185): results the hub is still waiting for (opactions.go). + r.OperatorActionResults = pendingOperatorActionResults() + // Geo-restriction status — ALWAYS present (even when never configured) so the hub // always renders the section. A nil pointer (omitempty) made the hub hide the whole // section for a never-configured controller; a present-but-disabled report renders diff --git a/controller/internal/report/opactions.go b/controller/internal/report/opactions.go new file mode 100644 index 0000000..a3e174a --- /dev/null +++ b/controller/internal/report/opactions.go @@ -0,0 +1,345 @@ +package report + +import ( + "context" + "errors" + "fmt" + "log" + "sort" + "strconv" + "sync" + "time" +) + +// Operator actions (R-314 / R-279 / R-177, `09` §3 decision 185 — D1, design option A of +// felhom.eu/documentation/audits/day-2026-10-08/design-R-314-279-177.md). +// +// THE DOOR IS THE REPORT REPLY, the shape selftail.go proved for log pulls: the operator presses a +// button on the hub, the hub stores a row and wakes this box's wait channel, the out-of-cycle report's +// reply carries `operator_actions: [{id, action, arg}]`, and the RUNNING controller acts on it — inside +// the process that owns the state, so there is no second process and no lost update (the CLI levers' +// problem, main.go). The result rides the next report as `operator_action_results: [{id, outcome, +// message}]`; the hub lists an action until its result has arrived, then stops. +// +// THE LIST IS CLOSED. Four actions, and run_job names one of four jobs. Nothing here deletes data, +// starts a countdown or shortens one — `03` §4 asks for a signing key only to destroy or overwrite the +// only copy, and none of these does. An unknown action, an unknown job or an argument out of range is +// answered `refused` and NOTHING is called. Pinned by TestOpActions_ClosedList (which fails when an +// entry is added, so a new action is a decision, not an edit) and TestOpActions_UnknownIsRefused. +// +// ONCE PER ID. An id is acted on the first time a reply carries it; every later reply that still lists +// it (the result has not reached the hub yet) only re-sends the result. The record is in memory: after a +// controller restart an action the hub still lists runs again. That is safe by construction — an +// off-site run or a check repeated adds nothing harmful, a second stop finds nothing to stop, and a +// second extension counts from the new „now" and can never shorten (ExtendAbandon refuses that). + +const ( + OpOffsiteBackupNow = "offsite_backup_now" + OpAbandonStop = "abandon_stop" + OpAbandonExtend = "abandon_extend" + OpRunJob = "run_job" + + OutcomeDone = "done" + OutcomeRefused = "refused" + OutcomeFailed = "failed" + + // AbandonExtendMinDays / AbandonExtendMaxDays bound abandon_extend's argument (decision 185). + AbandonExtendMinDays = 1 + AbandonExtendMaxDays = 30 + + // opActionsPerReply caps how many NEW actions one reply may start — a defence against a hub bug + // flooding the box, never reached by a person pressing buttons. + opActionsPerReply = 10 + // opActionForgetAfter drops a finished, no-longer-listed id from memory. + opActionForgetAfter = 24 * time.Hour + // offsiteRunTimeout matches the household's own „back up now" (offbox_handlers.go). + offsiteRunTimeout = 3 * time.Hour +) + +// operatorActionNames is THE closed list. Adding an entry fails TestOpActions_ClosedList by design. +var operatorActionNames = []string{OpOffsiteBackupNow, OpAbandonStop, OpAbandonExtend, OpRunJob} + +// operatorJobNames is the fixed set run_job may name — the scheduler's own job names (main.go). +var operatorJobNames = []string{"fill-watch", "offsite-integrity", "offsite-proof", "disk-health-check"} + +// OperatorActionNames returns the closed action list (sorted copy). +func OperatorActionNames() []string { return sortedCopy(operatorActionNames) } + +// OperatorJobNames returns the fixed run_job set (sorted copy). +func OperatorJobNames() []string { return sortedCopy(operatorJobNames) } + +func sortedCopy(in []string) []string { + out := append([]string(nil), in...) + sort.Strings(out) + return out +} + +func contains(list []string, v string) bool { + for _, x := range list { + if x == v { + return true + } + } + return false +} + +// ErrActionRefused, wrapped by a handler, turns its error into a `refused` outcome (the action could +// not START — e.g. an off-site run already in flight) instead of `failed` (it started and went wrong). +var ErrActionRefused = errors.New("refused") + +// OperatorAction is one entry of the reply's operator_actions list. +type OperatorAction struct { + ID int64 `json:"id"` + Action string `json:"action"` + Arg string `json:"arg,omitempty"` +} + +// OperatorActionResult is one entry of the report's operator_action_results list. +type OperatorActionResult struct { + ID int64 `json:"id"` + Outcome string `json:"outcome"` // done | refused | failed + Message string `json:"message"` +} + +// OperatorActionHandlers are the four doors, wired by main.go. A nil handler answers `refused` +// („not available on this box"). +type OperatorActionHandlers struct { + // OffsiteBackupNow returns a refusal sentence when the run may not start (the household's + // „back up now" checks), else the run itself, which the executor starts in its own goroutine. + OffsiteBackupNow func() (refusal string, run func(ctx context.Context) error) + AbandonStop func() error + AbandonExtend func(days int) (time.Time, error) + // RunJob is Scheduler.RunNow: an error = not started (unknown/running); the channel = the job's end. + RunJob func(name string) (<-chan error, error) + // ResultReady sends an out-of-cycle report so a result reaches the hub in seconds (optional). + ResultReady func() +} + +type opEntry struct { + finished bool + result OperatorActionResult + finishedAt time.Time +} + +// OperatorActions is the executor. One per process. +type OperatorActions struct { + ctx context.Context + h OperatorActionHandlers + logger *log.Logger + now func() time.Time + + mu sync.Mutex + seen map[int64]*opEntry + listed map[int64]bool // ids the LAST reply carried — results are sent only for these + running sync.WaitGroup // tests wait on async actions +} + +// NewOperatorActions builds the executor. ctx bounds the asynchronous actions. +func NewOperatorActions(ctx context.Context, h OperatorActionHandlers, logger *log.Logger) *OperatorActions { + return &OperatorActions{ctx: ctx, h: h, logger: logger, now: time.Now, + seen: map[int64]*opEntry{}, listed: map[int64]bool{}} +} + +func (o *OperatorActions) logf(format string, args ...interface{}) { + if o.logger != nil { + o.logger.Printf(format, args...) + } +} + +// Reconcile takes one reply's operator_actions list: new ids are acted on (once), every listed id is +// remembered so its result keeps riding the reports until the hub stops listing it. +func (o *OperatorActions) Reconcile(list []OperatorAction) { + o.mu.Lock() + o.listed = make(map[int64]bool, len(list)) + var fresh []OperatorAction + for _, a := range list { + o.listed[a.ID] = true + if _, ok := o.seen[a.ID]; ok { + continue // once per id: already acted on (or acting) — only its result is re-sent + } + if len(fresh) >= opActionsPerReply { + continue // not marked seen: the next reply offers it again + } + o.seen[a.ID] = &opEntry{} + fresh = append(fresh, a) + } + // Forget finished ids the hub no longer lists, after a day (memory bound; the hub never re-lists + // an id whose result it holds). + cutoff := o.now().Add(-opActionForgetAfter) + for id, e := range o.seen { + if e.finished && !o.listed[id] && e.finishedAt.Before(cutoff) { + delete(o.seen, id) + } + } + o.mu.Unlock() + + if len(list) > 0 { + o.logf("[DEBUG] [opaction] reply listed %d action(s), %d new", len(list), len(fresh)) + } + for _, a := range fresh { + o.dispatch(a) + } +} + +// Results returns the finished results the hub is still waiting for (ids the last reply listed). +func (o *OperatorActions) Results() []OperatorActionResult { + o.mu.Lock() + defer o.mu.Unlock() + var out []OperatorActionResult + for id, e := range o.seen { + if e.finished && o.listed[id] { + out = append(out, e.result) + } + } + sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID }) + return out +} + +func (o *OperatorActions) finish(a OperatorAction, outcome, msg string) { + o.mu.Lock() + if e, ok := o.seen[a.ID]; ok { + e.finished = true + e.finishedAt = o.now() + e.result = OperatorActionResult{ID: a.ID, Outcome: outcome, Message: msg} + } + o.mu.Unlock() + o.logf("[INFO] [opaction] operator action #%d %s: %s — %s", a.ID, a.Action, outcome, msg) + if o.h.ResultReady != nil { + o.h.ResultReady() + } +} + +func outcomeOf(err error) string { + switch { + case err == nil: + return OutcomeDone + case errors.Is(err, ErrActionRefused): + return OutcomeRefused + default: + return OutcomeFailed + } +} + +// dispatch validates against the closed list and acts. Validation happens BEFORE any handler is +// touched: a refusal calls nothing. +func (o *OperatorActions) dispatch(a OperatorAction) { + // Customer-visible transparency (the box's own log, like an operator log pull): never silent. + o.logf("[INFO] [opaction] operator action #%d received: %s", a.ID, a.Action) + switch a.Action { + case OpOffsiteBackupNow: + if a.Arg != "" { + o.finish(a, OutcomeRefused, "this action takes no argument") + return + } + if o.h.OffsiteBackupNow == nil { + o.finish(a, OutcomeRefused, "off-site backup is not available on this box") + return + } + refusal, run := o.h.OffsiteBackupNow() + if refusal != "" || run == nil { + if refusal == "" { + refusal = "the off-site backup could not be started" + } + o.finish(a, OutcomeRefused, refusal) + return + } + o.running.Add(1) + go func() { + defer o.running.Done() + ctx, cancel := context.WithTimeout(o.ctx, offsiteRunTimeout) + defer cancel() + err := run(ctx) + msg := "the off-site backup finished" + if err != nil { + msg = "the off-site backup: " + err.Error() + } + o.finish(a, outcomeOf(err), msg) + }() + case OpAbandonStop: + if a.Arg != "" { + o.finish(a, OutcomeRefused, "this action takes no argument") + return + } + if o.h.AbandonStop == nil { + o.finish(a, OutcomeRefused, "not available on this box") + return + } + if err := o.h.AbandonStop(); err != nil { + o.finish(a, outcomeOf(err), err.Error()) + return + } + o.finish(a, OutcomeDone, "the deletion countdown is stopped; the set-aside history is kept") + case OpAbandonExtend: + days, err := strconv.Atoi(a.Arg) + if err != nil || days < AbandonExtendMinDays || days > AbandonExtendMaxDays { + o.finish(a, OutcomeRefused, fmt.Sprintf("the number of days must be %d-%d", AbandonExtendMinDays, AbandonExtendMaxDays)) + return + } + if o.h.AbandonExtend == nil { + o.finish(a, OutcomeRefused, "not available on this box") + return + } + due, err := o.h.AbandonExtend(days) + if err != nil { + o.finish(a, outcomeOf(err), err.Error()) + return + } + o.finish(a, OutcomeDone, "the set-aside history is now deleted on "+due.Format("2006-01-02")) + case OpRunJob: + if !contains(operatorJobNames, a.Arg) { + o.finish(a, OutcomeRefused, "unknown job name") + return + } + if o.h.RunJob == nil { + o.finish(a, OutcomeRefused, "not available on this box") + return + } + done, err := o.h.RunJob(a.Arg) + if err != nil { + o.finish(a, OutcomeRefused, "the job "+a.Arg+" was not started: "+err.Error()) + return + } + o.running.Add(1) + go func() { + defer o.running.Done() + var jerr error + select { + case jerr = <-done: + case <-o.ctx.Done(): + jerr = o.ctx.Err() + } + msg := "the job " + a.Arg + " ran" + if jerr != nil { + msg = "the job " + a.Arg + ": " + jerr.Error() + } + o.finish(a, outcomeOf(jerr), msg) + }() + default: + o.finish(a, OutcomeRefused, "unknown action") + } +} + +// ── package wiring (the selftail.go shape: BuildReport reads it, main.go sets it) ── + +var ( + activeOpActionsMu sync.Mutex + activeOpActions *OperatorActions +) + +// SetOperatorActions wires the executor whose results BuildReport attaches. nil = none. +func SetOperatorActions(o *OperatorActions) { + activeOpActionsMu.Lock() + defer activeOpActionsMu.Unlock() + activeOpActions = o +} + +// pendingOperatorActionResults is what BuildReport attaches (nil in the steady state). +func pendingOperatorActionResults() []OperatorActionResult { + activeOpActionsMu.Lock() + o := activeOpActions + activeOpActionsMu.Unlock() + if o == nil { + return nil + } + return o.Results() +} diff --git a/controller/internal/report/opactions_test.go b/controller/internal/report/opactions_test.go new file mode 100644 index 0000000..be37bad --- /dev/null +++ b/controller/internal/report/opactions_test.go @@ -0,0 +1,193 @@ +package report + +import ( + "context" + "encoding/json" + "errors" + "reflect" + "sync" + "testing" + "time" +) + +// `09` §3 decision 185 (D1): the operator's actions. See opactions.go's header. + +// The list is CLOSED. This test fails when an action or a job is added, by design: a new entry is an +// operator decision (no action may delete data, start a countdown or shorten one), not an edit. The +// hub pins the same four in its own test (hub internal/store opactions_test.go). +func TestOpActions_ClosedList(t *testing.T) { + if got, want := OperatorActionNames(), []string{"abandon_extend", "abandon_stop", "offsite_backup_now", "run_job"}; !reflect.DeepEqual(got, want) { + t.Fatalf("operator actions = %v, want exactly %v — a new action needs the operator's word (decision 185)", got, want) + } + if got, want := OperatorJobNames(), []string{"disk-health-check", "fill-watch", "offsite-integrity", "offsite-proof"}; !reflect.DeepEqual(got, want) { + t.Fatalf("run_job names = %v, want exactly %v", got, want) + } +} + +// calls records every handler call. +type calls struct { + mu sync.Mutex + list []string +} + +func (c *calls) add(s string) { c.mu.Lock(); c.list = append(c.list, s); c.mu.Unlock() } +func (c *calls) get() []string { + c.mu.Lock() + defer c.mu.Unlock() + return append([]string(nil), c.list...) +} + +func recordingHandlers(c *calls) OperatorActionHandlers { + return OperatorActionHandlers{ + OffsiteBackupNow: func() (string, func(context.Context) error) { + c.add("offsite-check") + return "", func(context.Context) error { c.add("offsite-run"); return nil } + }, + AbandonStop: func() error { c.add("stop"); return nil }, + AbandonExtend: func(d int) (time.Time, error) { + c.add("extend") + return time.Date(2026, 11, 1, 0, 0, 0, 0, time.UTC), nil + }, + RunJob: func(name string) (<-chan error, error) { + c.add("job:" + name) + ch := make(chan error, 1) + ch <- nil + close(ch) + return ch, nil + }, + } +} + +func resultFor(t *testing.T, o *OperatorActions, id int64) OperatorActionResult { + t.Helper() + for _, r := range o.Results() { + if r.ID == id { + return r + } + } + t.Fatalf("no result for #%d in %+v", id, o.Results()) + return OperatorActionResult{} +} + +// Red test (3) of the design: an unknown action, an unknown job, an argument out of range → refused, +// and NOTHING is called. +func TestOpActions_UnknownIsRefusedAndCallsNothing(t *testing.T) { + c := &calls{} + o := NewOperatorActions(context.Background(), recordingHandlers(c), nil) + list := []OperatorAction{ + {ID: 1, Action: "delete_everything"}, + {ID: 2, Action: OpRunJob, Arg: "offsite-abandon-sweep"}, // a real job, NOT on the list: it deletes + {ID: 3, Action: OpRunJob, Arg: ""}, + {ID: 4, Action: OpAbandonExtend, Arg: "0"}, + {ID: 5, Action: OpAbandonExtend, Arg: "31"}, + {ID: 6, Action: OpAbandonExtend, Arg: "-3"}, + {ID: 7, Action: OpAbandonExtend, Arg: "7 "}, + {ID: 8, Action: OpAbandonStop, Arg: "x"}, + {ID: 9, Action: OpOffsiteBackupNow, Arg: "x"}, + } + o.Reconcile(list) + o.running.Wait() + if got := c.get(); len(got) != 0 { + t.Fatalf("a refused action called %v", got) + } + for _, a := range list { + if r := resultFor(t, o, a.ID); r.Outcome != OutcomeRefused || r.Message == "" { + t.Errorf("#%d %s(%q): %+v, want refused with a reason", a.ID, a.Action, a.Arg, r) + } + } +} + +// Red test (2): the same id delivered twice → the action runs ONCE; its result is re-sent while listed. +func TestOpActions_OncePerID(t *testing.T) { + c := &calls{} + o := NewOperatorActions(context.Background(), recordingHandlers(c), nil) + a := OperatorAction{ID: 42, Action: OpAbandonStop} + o.Reconcile([]OperatorAction{a}) + o.Reconcile([]OperatorAction{a}) + o.Reconcile([]OperatorAction{a}) + if got := c.get(); len(got) != 1 || got[0] != "stop" { + t.Fatalf("calls = %v, want exactly one stop", got) + } + if r := resultFor(t, o, 42); r.Outcome != OutcomeDone { + t.Fatalf("result = %+v", r) + } + // The hub has the result and stops listing the id → it is no longer sent, and is not run again. + o.Reconcile(nil) + if rs := o.Results(); len(rs) != 0 { + t.Fatalf("a result the hub no longer waits for is still sent: %+v", rs) + } + o.Reconcile([]OperatorAction{a}) + if got := c.get(); len(got) != 1 { + t.Fatalf("re-listed id ran again: %v", got) + } +} + +func TestOpActions_EachDoorAndItsOutcome(t *testing.T) { + c := &calls{} + h := recordingHandlers(c) + fired := 0 + var fmu sync.Mutex + h.ResultReady = func() { fmu.Lock(); fired++; fmu.Unlock() } + o := NewOperatorActions(context.Background(), h, nil) + o.Reconcile([]OperatorAction{ + {ID: 1, Action: OpOffsiteBackupNow}, + {ID: 2, Action: OpAbandonExtend, Arg: "30"}, + {ID: 3, Action: OpRunJob, Arg: "fill-watch"}, + }) + o.running.Wait() + for id := int64(1); id <= 3; id++ { + if r := resultFor(t, o, id); r.Outcome != OutcomeDone { + t.Errorf("#%d: %+v", id, r) + } + } + if r := resultFor(t, o, 2); r.Message != "the set-aside history is now deleted on 2026-11-01" { + t.Errorf("extend message = %q", r.Message) + } + fmu.Lock() + defer fmu.Unlock() + if fired != 3 { + t.Errorf("ResultReady fired %d times, want 3 (one per finished action)", fired) + } +} + +func TestOpActions_RefusalsAndFailuresAreDistinct(t *testing.T) { + h := OperatorActionHandlers{ + OffsiteBackupNow: func() (string, func(context.Context) error) { return "a backup run is already in flight", nil }, + AbandonStop: func() error { return errors.New("no abandonment countdown is running on this box") }, + RunJob: func(string) (<-chan error, error) { return nil, errors.New("the job is already running") }, + } + o := NewOperatorActions(context.Background(), h, nil) + o.Reconcile([]OperatorAction{{ID: 1, Action: OpOffsiteBackupNow}, {ID: 2, Action: OpAbandonStop}, {ID: 3, Action: OpRunJob, Arg: "offsite-proof"}, {ID: 4, Action: OpAbandonExtend, Arg: "5"}}) + o.running.Wait() + want := map[int64]string{1: OutcomeRefused, 2: OutcomeFailed, 3: OutcomeRefused, 4: OutcomeRefused /* nil handler */} + for id, w := range want { + if r := resultFor(t, o, id); r.Outcome != w { + t.Errorf("#%d: outcome %q, want %q (%s)", id, r.Outcome, w, r.Message) + } + } +} + +// The wire: the reply's operator_actions decodes into PushResponse, and the report's +// operator_action_results is what BuildReport attaches. +func TestOpActions_WireShapes(t *testing.T) { + var pr PushResponse + if err := json.Unmarshal([]byte(`{"status":"ok","operator_actions":[{"id":7,"action":"run_job","arg":"fill-watch"}]}`), &pr); err != nil { + t.Fatal(err) + } + if len(pr.OperatorActions) != 1 || pr.OperatorActions[0] != (OperatorAction{ID: 7, Action: "run_job", Arg: "fill-watch"}) { + t.Fatalf("decoded %+v", pr.OperatorActions) + } + c := &calls{} + o := NewOperatorActions(context.Background(), recordingHandlers(c), nil) + SetOperatorActions(o) + defer SetOperatorActions(nil) + o.Reconcile(pr.OperatorActions) + o.running.Wait() + b, _ := json.Marshal(Report{OperatorActionResults: pendingOperatorActionResults()}) + var back struct { + Results []map[string]interface{} `json:"operator_action_results"` + } + if err := json.Unmarshal(b, &back); err != nil || len(back.Results) != 1 || back.Results[0]["outcome"] != "done" || back.Results[0]["id"].(float64) != 7 { + t.Fatalf("report carried %s", b) + } +} diff --git a/controller/internal/report/pusher.go b/controller/internal/report/pusher.go index 31d039e..9e98947 100644 --- a/controller/internal/report/pusher.go +++ b/controller/internal/report/pusher.go @@ -51,8 +51,14 @@ type PushResponse struct { // Claim (v0.122.0, F-4) — the hub's active claim-code state (bcrypt hash + generation) for // the customer-claim gate. nil on an old hub / no claim row → the cache stays as-is. Claim *ClaimStatus `json:"claim"` + // OperatorActions (R-314/R-279/R-177, `09` §3 decision 185) — the operator's pending actions for + // this box, listed until their result arrives (opactions.go). Absent on an old hub = none. + OperatorActions []OperatorAction `json:"operator_actions"` } +// maxPushResponseBytes bounds the report ACK read (see Push). +const maxPushResponseBytes = 64 << 10 + // Pusher sends reports to the central hub. type Pusher struct { hubURL string @@ -126,8 +132,10 @@ func (p *Pusher) Push(report *Report) error { continue } - // Read response body to parse customer_blocked field - respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + // Read response body to parse customer_blocked field. The cap was 4 KiB; an ACK over it fails to + // parse and EVERY field in it is dropped (floor, config version, escrow, operator actions), so it + // is 64 KiB now — far above any real ACK, still a bound. + respBody, _ := io.ReadAll(io.LimitReader(resp.Body, maxPushResponseBytes)) resp.Body.Close() if resp.StatusCode >= 200 && resp.StatusCode < 300 { diff --git a/controller/internal/report/types.go b/controller/internal/report/types.go index ae15cb5..f20ff7c 100644 --- a/controller/internal/report/types.go +++ b/controller/internal/report/types.go @@ -45,6 +45,11 @@ type Report struct { // app-tail flow above is untouched). ControllerLogTail *ControllerLogTail `json:"controller_log_tail,omitempty"` + // OperatorActionResults (R-314/R-279/R-177, `09` §3 decision 185) — the outcome of each operator + // action the reply listed (opactions.go), re-sent until the hub stops listing the id. Additive; + // absent in the steady state. + OperatorActionResults []OperatorActionResult `json:"operator_action_results,omitempty"` + // Claimed (v0.122.0, F-4) — whether the customer has completed the dashboard claim (set // their own password). The hub ingests it SET-ONLY: a later false (wiped settings.json // after DR) never un-claims the customer hub-side. diff --git a/controller/internal/scheduler/runnow_test.go b/controller/internal/scheduler/runnow_test.go new file mode 100644 index 0000000..78e426e --- /dev/null +++ b/controller/internal/scheduler/runnow_test.go @@ -0,0 +1,76 @@ +package scheduler + +import ( + "context" + "errors" + "testing" + "time" +) + +// R-314/R-177 (`09` §3 decision 185): RunNow runs a registered job once, now, and refuses an unknown +// name and a job already running — the operator's run_job action rests on these three answers. +func TestRunNow_RunsOnceAndReportsCompletion(t *testing.T) { + s := discardScheduler() + calls := 0 + var sawOnDemand bool + s.Daily("fill-watch", farFuture(5), func(ctx context.Context) error { + calls++ + sawOnDemand = OnDemand(ctx) + return nil + }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + s.Start(ctx) + defer s.Stop() + + done, err := s.RunNow("fill-watch") + if err != nil { + t.Fatalf("RunNow(fill-watch): %v", err) + } + select { + case jerr := <-done: + if jerr != nil { + t.Fatalf("job error = %v", jerr) + } + case <-time.After(5 * time.Second): + t.Fatal("RunNow never completed") + } + if calls != 1 || !sawOnDemand { + t.Fatalf("calls=%d onDemand=%v, want 1 and true", calls, sawOnDemand) + } + if j := s.GetJobs()[0]; j.Running || j.LastRun.IsZero() { + t.Fatalf("after RunNow: Running=%v LastRun=%v — want false and set", j.Running, j.LastRun) + } +} + +func TestRunNow_RefusesUnknownRunningAndUnstarted(t *testing.T) { + s := discardScheduler() + release := make(chan struct{}) + s.Daily("offsite-integrity", farFuture(5), func(ctx context.Context) error { <-release; return nil }) + if _, err := s.RunNow("offsite-integrity"); !errors.Is(err, ErrNotStarted) { + t.Fatalf("before Start: err=%v, want ErrNotStarted", err) + } + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + s.Start(ctx) + defer s.Stop() + + if _, err := s.RunNow("no-such-job"); !errors.Is(err, ErrUnknownJob) { + t.Fatalf("unknown: err=%v, want ErrUnknownJob", err) + } + done, err := s.RunNow("offsite-integrity") + if err != nil { + t.Fatalf("first RunNow: %v", err) + } + if _, err := s.RunNow("offsite-integrity"); !errors.Is(err, ErrJobRunning) { + t.Fatalf("second RunNow while running: err=%v, want ErrJobRunning", err) + } + close(release) + <-done + // Released: it may run again. + done2, err := s.RunNow("offsite-integrity") + if err != nil { + t.Fatalf("RunNow after completion: %v", err) + } + <-done2 +} diff --git a/controller/internal/scheduler/scheduler.go b/controller/internal/scheduler/scheduler.go index ed34233..728351e 100644 --- a/controller/internal/scheduler/scheduler.go +++ b/controller/internal/scheduler/scheduler.go @@ -2,6 +2,7 @@ package scheduler import ( "context" + "errors" "fmt" "log" "sync" @@ -359,7 +360,74 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) { } job.Running = true s.mu.Unlock() + s.runClaimed(s.ctx, job, quiet) +} +// ErrUnknownJob / ErrJobRunning / ErrNotStarted are RunNow's refusals — values, so a caller can +// branch on them (R-224's lesson: a string is not something a caller can branch on). +var ( + ErrUnknownJob = errors.New("no job with this name is registered on this box") + ErrJobRunning = errors.New("the job is already running") + ErrNotStarted = errors.New("the scheduler is not running") +) + +// onDemandKey marks the context of a RunNow execution (OnDemand reads it). +type onDemandKey struct{} + +// OnDemand reports whether the job is running because someone asked for it now (RunNow), not on +// its schedule. A job whose scheduled run checks due-ness first (the off-site integrity check) uses +// it to run anyway — the operator asked for a check, not for "is a check due". +func OnDemand(ctx context.Context) bool { + v, _ := ctx.Value(onDemandKey{}).(bool) + return v +} + +// RunNow runs a registered job once, now, outside its schedule (R-314/R-177, `09` §3 decision 185: +// the operator's run_job action). It REFUSES an unknown name and a job that is already running — +// never queues, never runs it twice at once (the same Running flag the schedule uses, claimed under +// the mutex, so a scheduled tick and RunNow cannot both start it). On success the job runs in its own +// goroutine and its error (nil = completed) arrives on the returned channel, which is then closed. +func (s *Scheduler) RunNow(name string) (<-chan error, error) { + s.mu.Lock() + if !s.started || s.ctx == nil { + s.mu.Unlock() + return nil, ErrNotStarted + } + var job *Job + for _, j := range s.jobs { + if j.Name == name { + job = j + break + } + } + if job == nil { + s.mu.Unlock() + s.dbg("RunNow %q refused: unknown job", name) + return nil, ErrUnknownJob + } + if job.Running { + s.mu.Unlock() + s.dbg("RunNow %q refused: already running", name) + return nil, ErrJobRunning + } + job.Running = true + ctx := context.WithValue(s.ctx, onDemandKey{}, true) + s.wg.Add(1) + s.mu.Unlock() + + s.logger.Printf("[INFO] [scheduler] Job %s started on demand", name) + done := make(chan error, 1) + go func() { + defer s.wg.Done() + defer close(done) + done <- s.runClaimed(ctx, job, false) + }() + return done, nil +} + +// runClaimed executes a job whose Running flag the caller has already set. It clears the flag, +// records LastRun/LastErr, and returns the job's error. +func (s *Scheduler) runClaimed(ctx context.Context, job *Job, quiet bool) (err error) { defer func() { s.mu.Lock() job.Running = false @@ -369,8 +437,9 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) { // Panic recovery defer func() { if r := recover(); r != nil { + err = fmt.Errorf("panic: %v", r) s.mu.Lock() - job.LastErr = fmt.Errorf("panic: %v", r) + job.LastErr = err job.LastRun = time.Now() s.mu.Unlock() s.logger.Printf("[ERROR] [scheduler] Job %s panicked: %v", job.Name, r) @@ -383,7 +452,7 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) { s.dbg("job %s: execution starting", job.Name) start := time.Now() - err := job.Fn(s.ctx) + err = job.Fn(ctx) elapsed := time.Since(start) s.mu.Lock() @@ -400,6 +469,7 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) { // Routine per-cycle timing line → TRACE (dropped from the ring; the completed/failed lines above // carry the outcome). This was the biggest ring filler under load (fix-6). s.trace("job %s: finished in %s (err=%v)", job.Name, elapsed.Round(time.Millisecond), err) + return err } // parseDailyTime parses "HH:MM" and returns hour and minute.