package apply // The scheduler that fires scheduled steps on their cadence (novox/hq ADR 0053). // // A scheduled step is downstream of convergence, not a precondition of it. `applyContainer` installs // it as present state and reports the node current at once (see applySchedule); this is what // actually runs it, again and again, on the clock. The two are deliberately apart: a run happens // entirely here, outside the apply and the store, which is the whole reason a routine job's failure // can never fail the apply or flip the node's reported state. // // **It survives across applies and across a host restart.** The daemon holds one Scheduler for the // life of the process and re-establishes it from each applied declaration with Sync — the // declaration is the source of truth (ADR 0018), so nothing about a schedule is persisted that the // mesh does not already own. After a restart the first apply rebuilds every schedule from the // declaration the node kept; there is no separate schedule state to lose or to disagree with. // // **It does not install a system timer.** The rejected option 1 in ADR 0053 was a host command on a // cron/systemd timer, which widens what a compromised control plane can express toward "run this on // the machine, forever." Instead this host process runs the container itself, to completion, on the // cadence — the same `docker run` a run-once step uses, fired repeatedly. No new host action, no new // shape: `schedule` is a string on the container the host already has. import ( "context" "fmt" "sync" "time" "github.com/novox/mesh-host/internal/declaration" ) // Clock is the source of "now", injected so scheduled steps can be tested without waiting on the // wall clock (novox/hq ADR 0053: runs are driven by a controllable clock in tests, never a sleep). type Clock interface { Now() time.Time } type systemClock struct{} func (systemClock) Now() time.Time { return time.Now() } // SystemClock is the real clock, used by the daemon. func SystemClock() Clock { return systemClock{} } // Scheduler holds the node's scheduled steps and fires them when the clock says they are due. // // Safe for concurrent use: the daemon's ticker calls Advance while fired runs finish in their own // goroutines, and Sync may re-establish the set at any time as new declarations arrive. type Scheduler struct { clock Clock run Runner log func(string) mu sync.Mutex cri string // the container runtime, detected once and cached jobs map[string]*scheduledJob wg sync.WaitGroup // in-flight runs, so a caller (and a test) can wait for them } // scheduledJob is one scheduled container and where it is in its cadence. type scheduledJob struct { id string spec string // the declaration's digest, so a changed declaration re-establishes the job cron *declaration.Cron container *declaration.Container next time.Time // the next minute at which it is due running bool // a run is in flight — the next due run is skipped rather than stacked } // NewScheduler builds a scheduler. A nil clock is the system clock; a nil log says nothing. func NewScheduler(clock Clock, run Runner, log func(string)) *Scheduler { if clock == nil { clock = systemClock{} } if log == nil { log = func(string) {} } return &Scheduler{clock: clock, run: run, log: log, jobs: map[string]*scheduledJob{}} } // Sync re-establishes the scheduled steps from a declaration: it adds ones newly declared, re-arms // any whose image, environment or cadence changed, and forgets those the declaration no longer // names. Rebuilt from the declaration each apply because the declaration is the source of truth // (novox/hq ADR 0018) — there is no schedule state kept anywhere else to drift from it. // // A job whose declaration is unchanged keeps its place in the cadence — its next due time and // whether a run is in flight — so an ordinary reconcile every few minutes does not keep resetting // the clock out from under a schedule and prevent it ever firing. func (s *Scheduler) Sync(d *declaration.Declaration) { s.mu.Lock() defer s.mu.Unlock() seen := map[string]bool{} for _, r := range d.Resources { c, ok := r.(*declaration.Container) if !ok || c.Schedule == "" { continue } cron, err := declaration.ParseCron(c.Schedule) if err != nil { // The declaration parser already refused a malformed cron before this runs, so a // schedule that reaches here is well-formed. Guarded rather than trusted: a job silently // dropped would be a schedule that reports installed and never fires. s.log(fmt.Sprintf("scheduled step %s: ignoring an unparseable schedule %q: %v", c.Identity(), c.Schedule, err)) continue } seen[c.Identity()] = true spec := containerSpec(c) if existing := s.jobs[c.Identity()]; existing != nil && existing.spec == spec { // Unchanged: keep where it is in its cadence, refresh the declaration pointer only. existing.container = c continue } // New or changed: arm it for the next due minute after now. next, _ := cron.Next(s.clock.Now()) s.jobs[c.Identity()] = &scheduledJob{ id: c.Identity(), spec: spec, cron: cron, container: c, next: next, } } for id := range s.jobs { if !seen[id] { delete(s.jobs, id) } } } // Advance fires every scheduled step due at or before now, and is the whole of the clock-driven // behaviour — the daemon calls it on a ticker, and a test calls it with a controlled clock, so // nothing here ever waits on the wall clock. // // At most one run per step per call: if the host was asleep and several occurrences came due, they // collapse to a single run rather than a burst of concurrent copies. A step whose previous run is // still going is skipped and the skip logged — never a second copy started, which is the failure // mode that made the old timers dangerous (novox/hq ADR 0053). func (s *Scheduler) Advance(ctx context.Context, now time.Time) { s.mu.Lock() var toFire []*scheduledJob for _, j := range s.jobs { if j.next.IsZero() || now.Before(j.next) { continue } if j.running { s.log(fmt.Sprintf( "scheduled step %s: a previous run was still going when the run due at %s came — "+ "skipped, not stacked", j.id, j.next.Format(time.RFC3339))) // Drop the skipped occurrence and arm the next one after now. j.next, _ = j.cron.Next(now) continue } j.running = true j.next, _ = j.cron.Next(now) s.wg.Add(1) toFire = append(toFire, j) } s.mu.Unlock() // Fired outside the lock and in their own goroutines: a run to completion can take minutes, and // holding the lock — or blocking Advance — would stall every other schedule and the daemon's // ticker behind one slow job. for _, j := range toFire { go s.fire(ctx, j) } } // fire runs one occurrence of a scheduled step to completion, records the outcome, and clears the // running flag so the next occurrence may run. // // A non-zero exit is logged against the module and is otherwise nothing: it does not fail an apply // (there is no apply here), does not halt anything, and does not touch the store — so it cannot flip // the node's reported state. A poll that fails at 03:00 and succeeds at 03:05 is the system working; // a step that fails every time is a loud, repeating log entry, which is the right signal for a broken // recurring job and distinct from a machine that did not converge (novox/hq ADR 0053). func (s *Scheduler) fire(ctx context.Context, j *scheduledJob) { defer s.wg.Done() defer func() { s.mu.Lock() j.running = false s.mu.Unlock() }() cri, err := s.runtime(ctx) if err != nil { s.log(fmt.Sprintf("scheduled step %s: no container runtime to run it: %v", j.id, err)) return } // A container by this name left exited by the previous run would collide with --name. Removing // one that is not there is the state we want, so its error is ignored — the same as run-once. _, _ = s.run(ctx, cri, "rm", "-f", j.container.Name) args := foregroundRunArgs(j.container, j.spec) if _, err := s.run(ctx, cri, args...); err != nil { s.log(fmt.Sprintf( "scheduled step %s: run exited non-zero — recorded, and left for the next run "+ "(it does not fail the node): %v", j.id, err)) return } // Remove the exited container so the next run's --name is free; the step left nothing to inspect. _, _ = s.run(ctx, cri, "rm", "-f", j.container.Name) s.log(fmt.Sprintf("scheduled step %s: run completed", j.id)) } // runtime detects the container runtime once and caches it. A scheduled step does not need one to be // installed (see applySchedule), only to be fired, so detection is deferred to here. func (s *Scheduler) runtime(ctx context.Context) (string, error) { s.mu.Lock() cached := s.cri s.mu.Unlock() if cached != "" { return cached, nil } cri, err := containerRuntime(ctx, s.run) if err != nil { return "", err } s.mu.Lock() s.cri = cri s.mu.Unlock() return cri, nil } // Wait blocks until every in-flight run has finished. For orderly shutdown, and for tests that must // observe a run's effect without racing it. func (s *Scheduler) Wait() { s.wg.Wait() } // Run drives the scheduler off the real clock until the context is cancelled. This is the daemon's // entry point; tests drive Advance directly instead. // // A minute tick because cron resolves to the minute — a step due at 03:00 fires within a minute of // it, which is what a cadence measured in minutes, hours and days asks for. It does not install a // system timer (the rejected option 1): this process is the thing on the clock. func (s *Scheduler) Run(ctx context.Context) { ticker := time.NewTicker(time.Minute) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: s.Advance(ctx, s.clock.Now()) } } }