Files
mesh-host/internal/apply/schedule.go
T
jschoubben e5af6bb58e apply: ensure a scheduled container's image is present at apply, without running it
A schedule: container (ADR 0053) is installed as present state and never run at apply — the
Scheduler fires it later on its cadence. But a service or run-once container only gets its
image as a side effect of docker run, so a scheduled step's image was not pulled until its
first scheduled fire: absent from the node right after a successful apply, so the first run
paid the whole pull latency and tooling that expects the image present after apply found it
missing.

applyContainer now probes the runtime and ensures the pinned image present for a scheduled
step before recording it. A new ensureImage helper inspects the image and pulls it only if
absent, then reads back (ADR 0018). Ensuring an image is not running it: no docker run fires
the container, so the no-run invariant of ADR 0053 holds. The runtime probe, previously
skipped for a schedule, now runs because a pull needs it — the schedule.go comment is updated
to match.

Tests: the install-does-not-run test is extended to allow the image-ensure while asserting no
fire and no needless pull; a new test applies a scheduled container whose image is absent and
asserts it is pulled and still not started. go build, go vet, go test ./... all pass.

Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
2026-09-07 02:23:46 +02:00

248 lines
9.7 KiB
Go

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. The Scheduler needs one only to fire a
// run, so detection is deferred to here — separate from the apply, which does its own probe when it
// ensures a scheduled step's image is present (see applyContainer).
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())
}
}
}