From 9d1f001dcc0626ad9e0486546989333377e23e41 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 6 Sep 2026 14:08:44 +0200 Subject: [PATCH] apply: a scheduled step is a container run on a cadence (ADR 0053) The recurring twin of run-once, one modifier over: a container marked schedule: "" is run to completion on its cadence, not started as a service and not run once as a gate. The gating rule is deliberately reversed. Installing a schedule records it as present state and reports the node current at once (applySchedule) -- it never runs the container and does not gate what follows. A Scheduler, held for the life of the daemon and re-established from each applied declaration (the declaration is the source of truth, ADR 0018), fires the container off an injected clock. A run that exits non-zero is logged and never fails the apply or flips the node's state, because it happens outside the apply and the store entirely. Runs never stack: a run still going when the next is due is skipped, not started as a second copy. No new host shape and no new action -- schedule is a string on the container the host already has, and the host process runs the container itself rather than installing a system timer (the rejected option 1). A minimal five-field cron (declaration/cron.go) validates on arrival and computes the next due minute; time is injected so the scheduler is tested without the wall clock. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- cmd/mesh-host/main.go | 39 +- internal/apply/apply.go | 83 +++- internal/apply/schedule.go | 246 ++++++++++++ internal/apply/schedule_test.go | 366 ++++++++++++++++++ internal/declaration/cron.go | 182 +++++++++ internal/declaration/cron_test.go | 100 +++++ internal/declaration/declaration.go | 29 ++ .../declaration/schedule_declaration_test.go | 53 +++ 8 files changed, 1076 insertions(+), 22 deletions(-) create mode 100644 internal/apply/schedule.go create mode 100644 internal/apply/schedule_test.go create mode 100644 internal/declaration/cron.go create mode 100644 internal/declaration/cron_test.go create mode 100644 internal/declaration/schedule_declaration_test.go diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 77b0558..2095c26 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -594,16 +594,24 @@ func runLink(ctx context.Context, opts options) error { fmt.Printf("node %s, linking to %s\n", mine.Node, mine.Membership.Broker) - apply := func(ctx context.Context, raw, signature []byte) link.Report { - return applyAndKeep(ctx, opts, raw, &store.Declared{Declaration: raw, Signature: signature}) - } say := func(line string) { fmt.Println(line) } + // One scheduler for the life of the process, re-established from each applied declaration + // (novox/hq ADR 0053). It fires scheduled steps on their cadence, surviving across applies and + // across the reconcile loop; a host restart rebuilds it from the declaration the node kept, the + // first time either path applies. Its own loop is the thing on the clock — no system timer. + sched := apply.NewScheduler(apply.SystemClock(), apply.ExecRunner, say) + go sched.Run(ctx) + + applier := func(ctx context.Context, raw, signature []byte) link.Report { + return applyAndKeep(ctx, opts, raw, &store.Declared{Declaration: raw, Signature: signature}, sched) + } + // Two things at once, and the second is what makes disconnection ordinary. The link brings // new declarations; this holds the machine in the last one whether the link is up or not. A // laptop shut for a week comes back and reconciles — it does not come back and ask what it is // (novox/hq ADR 0004). - go holdTheMachine(ctx, opts, mine, say) + go holdTheMachine(ctx, opts, mine, say, sched) return link.HoldRoused(ctx, link.Membership{ Node: mine.Node, @@ -611,7 +619,7 @@ func runLink(ctx context.Context, opts options) error { Fingerprint: mine.Membership.Fingerprint, Password: mine.Membership.Password, Signer: mine.Membership.Signer, - }, apply, say, opts.timeout, rousedBySignal(ctx)) + }, applier, say, opts.timeout, rousedBySignal(ctx)) } // rousedBySignal is the machine telling this process that its link is probably stale. @@ -658,7 +666,8 @@ func rousedBySignal(ctx context.Context) link.Roused { // changed it, and then for ever. const ReconcileEvery = 5 * time.Minute -func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, say link.Announce) { +func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, say link.Announce, + sched *apply.Scheduler) { ticker := time.NewTicker(ReconcileEvery) defer ticker.Stop() @@ -680,7 +689,7 @@ func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, s continue } - report := applyDeclared(ctx, opts, declared) + report := applyDeclared(ctx, opts, declared, sched) switch { case report.Refused != "": say("what this node was last told no longer applies: " + report.Refused) @@ -695,13 +704,14 @@ func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, s // Signature checking happens before this is called, in the link. By the time anything here runs, // the question "is this from the mesh I joined" is settled — which is why this can treat the // bytes as instructions. -func applyDeclared(ctx context.Context, opts options, raw []byte) link.Report { - return applyAndKeep(ctx, opts, raw, nil) +func applyDeclared(ctx context.Context, opts options, raw []byte, sched *apply.Scheduler) link.Report { + return applyAndKeep(ctx, opts, raw, nil, sched) } // applyAndKeep applies a declaration and, when it came from the mesh, keeps it so this node can // go on obeying it while disconnected. -func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.Declared) link.Report { +func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.Declared, + sched *apply.Scheduler) link.Report { declared, err := declaration.Parse(raw) if err != nil { return link.Report{Refused: err.Error()} @@ -736,6 +746,15 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D saveErr.Error()} } + // Re-establish the scheduled steps from the declaration just applied (novox/hq ADR 0053). Done + // from the parsed declaration, which is the source of truth (ADR 0018): a schedule newly declared + // is armed, one whose image, environment or cadence changed is re-armed, and one no longer + // declared is forgotten — and after a host restart the first apply rebuilds them all. A nil + // scheduler is the one-shot CLI path, which exits rather than staying up to fire anything. + if sched != nil { + sched.Sync(declared) + } + report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)} for _, change := range outcome.Outcomes { report.Applied = append(report.Applied, change.ID) diff --git a/internal/apply/apply.go b/internal/apply/apply.go index 838b2ed..ad21d29 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -874,6 +874,12 @@ func containerSpec(r *declaration.Container) string { for _, a := range r.Args { b.WriteString("arg " + a + "\n") } + // The cadence is part of what was declared, so a changed schedule is a changed spec — the marker + // moves and the install is reported "updated" and re-established. Added only when present, so no + // ordinary container's or run-once step's digest moves for a field it does not set. + if r.Schedule != "" { + b.WriteString("schedule " + r.Schedule + "\n") + } return fmt.Sprintf("%x", sha256.Sum256([]byte(b.String()))) } @@ -942,6 +948,16 @@ func applyContainer(ctx context.Context, r *declaration.Container, run Runner, c out := begin(r) want := containerSpec(r) + // A scheduled step is state that is present, not a container to start (novox/hq ADR 0053). + // Installing it records the schedule and reports the node current at once — the deliberate + // inversion of run-once, which gates. The recurring run is fired by the host's Scheduler off the + // clock, re-established from this declaration each apply, and NEVER here — so installing does not + // run the container and does not even need a runtime present. Checked before the runtime probe + // for exactly that reason. + if r.Schedule != "" { + return applySchedule(r, want, previous) + } + cri, err := containerRuntime(ctx, run) if err != nil { return out, fmt.Errorf("%w, so nothing can be said about %q", err, r.Name) @@ -1069,6 +1085,29 @@ func applyRunOnce(ctx context.Context, r *declaration.Container, run Runner, cri // Run in the foreground so the runtime waits for the container and hands back its exit code. // No --detach and no --restart: a step that is restarted is not a step. + args := foregroundRunArgs(r, want) + + if _, err := run(ctx, cri, args...); err != nil { + // A non-zero exit or a runtime that could not start it. Either way the step did not make + // the machine ready, so the apply must not go on to the container that needs it. + return out, fmt.Errorf("run-once step %s did not complete: %w", r.Name, err) + } + + // It completed. Remove the exited container so a later apply is not confused by a stopped one; + // the record that it ran is the digest below, which the caller persists after the fact. + _, _ = run(ctx, cri, "rm", "-f", r.Name) + + out.Action = "created" + out.Detail = "run-once step completed" + out.wrote = want + return out, nil +} + +// foregroundRunArgs builds a `docker run` that runs a container to completion and hands back its +// exit code — no --detach, no --restart, because a step that is restarted is not a step. Shared by +// a run-once step (novox/hq ADR 0052) and by one fire of a scheduled step (ADR 0053), which are the +// same "run the container and let it exit" up to how often it happens. +func foregroundRunArgs(r *declaration.Container, want string) []string { args := []string{"run", "--name", r.Name} for _, file := range r.EnvFile { args = append(args, "--env-file", file) @@ -1088,20 +1127,40 @@ func applyRunOnce(ctx context.Context, r *declaration.Container, run Runner, cri } args = append(args, r.Image) args = append(args, r.Args...) + return args +} - if _, err := run(ctx, cri, args...); err != nil { - // A non-zero exit or a runtime that could not start it. Either way the step did not make - // the machine ready, so the apply must not go on to the container that needs it. - return out, fmt.Errorf("run-once step %s did not complete: %w", r.Name, err) - } - - // It completed. Remove the exited container so a later apply is not confused by a stopped one; - // the record that it ran is the digest below, which the caller persists after the fact. - _, _ = run(ctx, cri, "rm", "-f", r.Name) - - out.Action = "created" - out.Detail = "run-once step completed" +// applySchedule installs a scheduled step: it records the schedule as present and reports the node +// current, without running anything (novox/hq ADR 0053). +// +// This is the deliberate inversion of run-once. A run-once step gates the apply — it runs to +// completion here and a non-zero exit halts everything after it — because "seed the store before the +// broker starts" is a precondition of convergence. A scheduled step is the opposite: it runs *after* +// the machine is up, on its own clock, and a single failed run is an ordinary operational event. So +// installing it is pure state: the schedule is present, like a running service, and the apply is +// current at once. The recurring run is fired by the host's Scheduler off the clock (see +// schedule.go), re-established from the applied declaration each pass because the declaration is the +// source of truth (ADR 0018) — never from here, and never persisted beyond what the mesh already +// owns. +// +// The marker is the declaration's digest, exactly as for run-once, so a re-apply of the same +// declaration reports the schedule unchanged and a changed image, environment or cadence reports it +// re-installed. A failed run touches none of this: it happens entirely in the Scheduler, outside the +// apply and the store, which is why a routine job's failure can never flip the node's state. +func applySchedule(r *declaration.Container, want string, previous store.Applied) (Outcome, error) { + out := begin(r) out.wrote = want + switch { + case previous.Wrote == want: + out.Action = "unchanged" + out.Detail = "scheduled step; already installed for this declaration" + case previous.Wrote != "": + out.Action = "updated" + out.Detail = "scheduled step re-installed; its image, environment or cadence changed" + default: + out.Action = "created" + out.Detail = "scheduled step installed; the host runs it on its cadence" + } return out, nil } diff --git a/internal/apply/schedule.go b/internal/apply/schedule.go new file mode 100644 index 0000000..3357128 --- /dev/null +++ b/internal/apply/schedule.go @@ -0,0 +1,246 @@ +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()) + } + } +} diff --git a/internal/apply/schedule_test.go b/internal/apply/schedule_test.go new file mode 100644 index 0000000..fb30900 --- /dev/null +++ b/internal/apply/schedule_test.go @@ -0,0 +1,366 @@ +package apply + +import ( + "context" + "errors" + "strings" + "sync" + "testing" + "time" + + "github.com/novox/mesh-host/internal/declaration" + "github.com/novox/mesh-host/internal/store" +) + +// A scheduled step is the recurring twin of run-once, with the gating rule deliberately reversed +// (novox/hq ADR 0053). These tests defend that, each named for the claim it holds up: installing it +// does not run it and the node is current at once; the host fires it when the cron is due and not +// before; a failed run is recorded and does not fail anything; and runs never stack. Time is +// injected — nothing here waits on the wall clock. + +// fixedClock is a clock a test sets by hand. +type fixedClock struct { + mu sync.Mutex + now time.Time +} + +func (c *fixedClock) Now() time.Time { + c.mu.Lock() + defer c.mu.Unlock() + return c.now +} + +func (c *fixedClock) set(t time.Time) { + c.mu.Lock() + c.now = t + c.mu.Unlock() +} + +// recordingRun remembers every command it was asked to run, guarded for the goroutines fire() uses. +type recordingRun struct { + mu sync.Mutex + runs int // how many `docker run` (a fire) happened + all [][]string +} + +func (r *recordingRun) run(ctx context.Context, name string, args ...string) (string, error) { + r.mu.Lock() + defer r.mu.Unlock() + r.all = append(r.all, append([]string{name}, args...)) + switch args[0] { + case "info": + return "27.0\n", nil // a container runtime answers + case "run": + r.runs++ + return "", nil + } + return "", nil +} + +func (r *recordingRun) fireCount() int { + r.mu.Lock() + defer r.mu.Unlock() + return r.runs +} + +// specOf returns the mesh-host.spec label value from a docker run argument list. +func specOf(args []string) string { + for _, a := range args { + if strings.HasPrefix(a, specLabel+"=") { + return strings.TrimPrefix(a, specLabel+"=") + } + } + return "" +} + +func scheduledContainer(t *testing.T, schedule string) *declaration.Container { + t.Helper() + d := parseTrusted(t, `{"declaration":1,"resources":[ + {"id":"sync","type":"container","name":"sync","image":"`+pinned+`","schedule":"`+schedule+`"} + ]}`) + return d.Resources[0].(*declaration.Container) +} + +// --- Requirement 4: installing a schedule does not run it, and the node is current at once. --- + +func TestInstallingAScheduleDoesNotRunItAndReportsCurrent(t *testing.T) { + // The deliberate inversion of run-once: the schedule is state that is present, so the apply is + // current as soon as it is recorded — nothing is run, and no runtime is even probed. + var calls []string + run := func(ctx context.Context, name string, args ...string) (string, error) { + calls = append(calls, name+" "+strings.Join(args, " ")) + return "", errors.New("installing a schedule must not run any command") + } + d := parseTrusted(t, `{"declaration":1,"resources":[ + {"id":"sync","type":"container","name":"sync","image":"`+pinned+`","schedule":"0 3 * * *"} + ]}`) + + report, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil, nil) + if err != nil { + t.Fatalf("installing a schedule failed the apply: %v", err) + } + if len(calls) != 0 { + t.Errorf("installing a schedule ran commands, so it did more than record state: %v", calls) + } + if report.Outcomes[0].Action != "created" { + t.Errorf("an installed schedule was not reported created: %+v", report.Outcomes[0]) + } + // Recorded as applied — the node is current, exactly as it is for a running service. + applied, ok := state.Find("sync") + if !ok { + t.Fatal("an installed schedule was not recorded as applied") + } + if applied.Wrote == "" { + t.Error("an installed schedule recorded no marker, so a re-apply cannot tell it is unchanged") + } + + // And a re-apply of the same declaration is unchanged and still runs nothing. + report2, _, err := Apply(context.Background(), archHost(t), d, state, store.OriginCarried, run, nil, nil) + if err != nil { + t.Fatal(err) + } + if report2.Changed() { + t.Errorf("re-installing the same schedule reported a change: %+v", report2.Outcomes) + } +} + +func TestAScheduledStepDoesNotGateWhatFollows(t *testing.T) { + // A run-once step gates: a failure halts what follows. A scheduled step is downstream of + // convergence, so it never gates — a container declared after it is applied normally. + d := parseTrusted(t, `{"declaration":1,"resources":[ + {"id":"sync","type":"container","name":"sync","image":"`+pinned+`","schedule":"0 3 * * *"}, + {"id":"web","type":"container","name":"web","image":"`+pinned+`"} + ]}`) + + var startedWeb bool + specs := map[string]string{} // name -> spec label, so read-back sees a running container + run := func(ctx context.Context, name string, args ...string) (string, error) { + switch args[0] { + case "info": + return "27.0\n", nil + case "inspect": + // The container name is the last argument to `docker inspect --format ... `. + target := args[len(args)-1] + if spec, up := specs[target]; up { + return "true\t" + spec + "\n", nil + } + return "false\t\n", errors.New("no such container") + case "run": + n := nameOf(args) + if n == "web" { + startedWeb = true + } + specs[n] = specOf(args) + return "deadbeef\n", nil + } + return "", nil + } + + if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil, nil); err != nil { + t.Fatalf("a declaration with a scheduled step failed to apply: %v", err) + } + if !startedWeb { + t.Error("a container declared after a scheduled step was not started — the schedule gated, and it must not") + } +} + +// --- Requirement 5: the host fires the container when the cron is due, and not before. --- + +func TestAScheduledStepRunsWhenDueAndNotBefore(t *testing.T) { + clock := &fixedClock{now: time.Date(2026, 9, 7, 12, 0, 30, 0, time.UTC)} + rec := &recordingRun{} + s := NewScheduler(clock, rec.run, nil) + + d := &declaration.Declaration{Version: 1, Resources: []declaration.Resource{ + scheduledContainer(t, "* * * * *"), // every minute; next due 12:01:00 + }} + s.Sync(d) + + // Not yet due: 12:00:45 is before 12:01:00, so nothing runs. + s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 0, 45, 0, time.UTC)) + s.Wait() + if n := rec.fireCount(); n != 0 { + t.Fatalf("a scheduled step ran before it was due (%d run(s))", n) + } + + // Due: 12:01:05 is at or after 12:01:00, so it runs once, to completion. + s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC)) + s.Wait() + if n := rec.fireCount(); n != 1 { + t.Fatalf("a scheduled step due did not run exactly once (%d run(s))", n) + } + + // And it ran to completion, not detached and not with a restart policy — a step, not a service. + rec.mu.Lock() + var runArgs []string + for _, c := range rec.all { + if len(c) > 1 && c[1] == "run" { + runArgs = c + } + } + rec.mu.Unlock() + joined := strings.Join(runArgs, " ") + if strings.Contains(joined, "--detach") { + t.Errorf("a scheduled step was detached, so its exit could not be observed: %s", joined) + } + if strings.Contains(joined, "--restart") { + t.Errorf("a scheduled step was given a restart policy, which makes it a service: %s", joined) + } + if !strings.Contains(joined, pinned) { + t.Errorf("a scheduled step ran the wrong image: %s", joined) + } +} + +// --- Requirement 6: a non-zero run is recorded and does not fail the apply or the node. --- + +func TestAFailedRunIsRecordedAndDoesNotFailAnything(t *testing.T) { + clock := &fixedClock{now: time.Date(2026, 9, 7, 12, 0, 30, 0, time.UTC)} + + var logged []string + var logMu sync.Mutex + log := func(line string) { + logMu.Lock() + logged = append(logged, line) + logMu.Unlock() + } + + run := func(ctx context.Context, name string, args ...string) (string, error) { + switch args[0] { + case "info": + return "27.0\n", nil + case "run": + return "", errors.New("exit status 1") // the run fails, every time + } + return "", nil + } + s := NewScheduler(clock, run, log) + s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ + scheduledContainer(t, "* * * * *"), + }}) + + // Fire a run that exits non-zero. Advance returns nothing — there is no error to fail an apply, + // because the run happens outside any apply and outside the store. + s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC)) + s.Wait() + + logMu.Lock() + defer logMu.Unlock() + var recorded bool + for _, line := range logged { + if strings.Contains(line, "non-zero") { + recorded = true + } + } + if !recorded { + t.Errorf("a failed run was not recorded against the module: %v", logged) + } + // The scheduler is still healthy: the job's run flag was cleared, so the next occurrence can run. + s.mu.Lock() + stillRunning := s.jobs["sync"].running + s.mu.Unlock() + if stillRunning { + t.Error("a failed run left the job marked running, which would block every future run") + } +} + +// --- Requirement 7: runs do not stack. --- + +func TestASlowRunSkipsTheNextDueRunRatherThanStacking(t *testing.T) { + clock := &fixedClock{now: time.Date(2026, 9, 7, 12, 0, 30, 0, time.UTC)} + + started := make(chan struct{}) // fire() signals it entered the runner + release := make(chan struct{}) // the test lets the run finish + var enters int + var mu sync.Mutex + run := func(ctx context.Context, name string, args ...string) (string, error) { + if args[0] == "info" { + return "27.0\n", nil + } + if args[0] == "run" { + mu.Lock() + enters++ + first := enters == 1 + mu.Unlock() + if first { + close(started) + <-release // the first run outlasts the next due time + } + return "", nil + } + return "", nil + } + + var logged []string + var logMu sync.Mutex + log := func(line string) { + logMu.Lock() + logged = append(logged, line) + logMu.Unlock() + } + + s := NewScheduler(clock, run, log) + s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ + scheduledContainer(t, "* * * * *"), // every minute + }}) + + // First occurrence: 12:01 is due — starts a run that blocks in the runner. + s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC)) + <-started // the run is now in flight and will not return until released + + // Second occurrence: 12:02 is due while the first run is still going. It must be skipped, not + // started as a second concurrent copy. + s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 2, 5, 0, time.UTC)) + + // Let the first run finish and settle. + close(release) + s.Wait() + + mu.Lock() + total := enters + mu.Unlock() + if total != 1 { + t.Fatalf("a second run was started while the first was still going (%d run(s)) — runs stacked", total) + } + + logMu.Lock() + defer logMu.Unlock() + var skipped bool + for _, line := range logged { + if strings.Contains(line, "skipped") { + skipped = true + } + } + if !skipped { + t.Errorf("the overrun run was not logged as skipped: %v", logged) + } +} + +// --- Sync re-establishes schedules from the declaration (novox/hq ADR 0018, 0053). --- + +func TestSyncForgetsAScheduleTheDeclarationNoLongerNames(t *testing.T) { + clock := &fixedClock{now: time.Date(2026, 9, 7, 12, 0, 30, 0, time.UTC)} + rec := &recordingRun{} + s := NewScheduler(clock, rec.run, nil) + + s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ + scheduledContainer(t, "* * * * *"), + }}) + s.mu.Lock() + have := len(s.jobs) + s.mu.Unlock() + if have != 1 { + t.Fatalf("a declared schedule was not established: %d job(s)", have) + } + + // A later declaration no longer names it — the schedule is dropped, and a due tick runs nothing. + s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ + parseTrusted(t, `{"declaration":1,"resources":[ + {"id":"web","type":"container","name":"web","image":"`+pinned+`"} + ]}`).Resources[0], + }}) + s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 5, 5, 0, time.UTC)) + s.Wait() + if n := rec.fireCount(); n != 0 { + t.Errorf("a schedule the declaration no longer names still fired (%d run(s))", n) + } +} diff --git a/internal/declaration/cron.go b/internal/declaration/cron.go new file mode 100644 index 0000000..69b97df --- /dev/null +++ b/internal/declaration/cron.go @@ -0,0 +1,182 @@ +package declaration + +// A minimal five-field cron, enough to say when a scheduled step is due (novox/hq ADR 0053). +// +// **Written rather than pulled in.** The host depends on nothing it does not have to +// (novox/hq ADR 0005), and a scheduled step needs exactly two questions answered — *is this minute +// a match* and *when is the next one* — over the ordinary five fields (minute, hour, day-of-month, +// month, day-of-week) with `*`, lists (`,`), ranges (`-`) and steps (`/`). That is small enough to +// keep in the vocabulary the host already owns, and a cron library would be a dependency carried +// for the parts of it nobody here uses. +// +// **The host both validates and evaluates.** The control plane refuses a malformed cron near its +// author (novox/hq ADR 0053), and the host refuses it again on arrival for the same near-versus-far +// reason every other field is checked here: a declaration the host does not fully understand is +// refused whole rather than half-applied. The evaluation half — `Next` — is what the scheduler +// fires on. + +import ( + "fmt" + "strconv" + "strings" + "time" +) + +// Cron is a parsed five-field expression, each field held as a bitset of the values it permits. +// +// The two day fields carry a flag for whether they were written as a bare `*`, because standard +// cron gives them a special rule: when day-of-month and day-of-week are *both* restricted, a day +// matches if *either* does; when one is `*`, only the other is consulted. Getting that wrong is the +// classic cron surprise, so the flag is kept rather than rediscovered. +type Cron struct { + minute uint64 + hour uint64 + dom uint64 + month uint64 + dow uint64 + domStar bool + dowStar bool +} + +// ParseCron reads a five-field cron expression, or says why it is not one. +func ParseCron(expr string) (*Cron, error) { + fields := strings.Fields(expr) + if len(fields) != 5 { + return nil, fmt.Errorf( + "a schedule is a five-field cron expression (minute hour day-of-month month "+ + "day-of-week), and %q has %d field(s)", expr, len(fields)) + } + + minute, _, err := parseCronField(fields[0], 0, 59) + if err != nil { + return nil, fmt.Errorf("schedule minute field: %w", err) + } + hour, _, err := parseCronField(fields[1], 0, 23) + if err != nil { + return nil, fmt.Errorf("schedule hour field: %w", err) + } + dom, domStar, err := parseCronField(fields[2], 1, 31) + if err != nil { + return nil, fmt.Errorf("schedule day-of-month field: %w", err) + } + month, _, err := parseCronField(fields[3], 1, 12) + if err != nil { + return nil, fmt.Errorf("schedule month field: %w", err) + } + // Day-of-week is 0-6 with Sunday at 0, and 7 is also accepted for Sunday — the convention every + // cron keeps, so a crontab copied from elsewhere is not refused for saying 7. + dow, dowStar, err := parseCronField(fields[4], 0, 7) + if err != nil { + return nil, fmt.Errorf("schedule day-of-week field: %w", err) + } + if dow&(1<<7) != 0 { + dow |= 1 << 0 + dow &^= 1 << 7 + } + + return &Cron{ + minute: minute, hour: hour, dom: dom, month: month, dow: dow, + domStar: domStar, dowStar: dowStar, + }, nil +} + +// parseCronField turns one field into the bitset of values it permits, and reports whether it was a +// bare `*` (which the day fields treat specially). +func parseCronField(spec string, min, max int) (uint64, bool, error) { + if spec == "" { + return 0, false, fmt.Errorf("is empty") + } + star := spec == "*" + + var bits uint64 + for _, part := range strings.Split(spec, ",") { + if part == "" { + return 0, false, fmt.Errorf("%q has an empty element between commas", spec) + } + + // A step may follow either `*` or a range: `*/15`, `0-30/5`. + step := 1 + rangePart := part + if slash := strings.IndexByte(part, '/'); slash >= 0 { + rangePart = part[:slash] + n, err := strconv.Atoi(part[slash+1:]) + if err != nil || n < 1 { + return 0, false, fmt.Errorf("step in %q is not a positive number", part) + } + step = n + } + + lo, hi := min, max + switch { + case rangePart == "*": + // Full range, already set. + case strings.IndexByte(rangePart, '-') >= 0: + dash := strings.IndexByte(rangePart, '-') + a, errA := strconv.Atoi(rangePart[:dash]) + b, errB := strconv.Atoi(rangePart[dash+1:]) + if errA != nil || errB != nil { + return 0, false, fmt.Errorf("range %q is not two numbers", rangePart) + } + lo, hi = a, b + default: + n, err := strconv.Atoi(rangePart) + if err != nil { + return 0, false, fmt.Errorf("%q is not a number", rangePart) + } + lo, hi = n, n + } + + if lo < min || hi > max || lo > hi { + return 0, false, fmt.Errorf( + "%q is outside the allowed range %d-%d", part, min, max) + } + for v := lo; v <= hi; v += step { + bits |= 1 << uint(v) + } + } + return bits, star, nil +} + +// Matches reports whether a scheduled step is due at this minute. +func (c *Cron) Matches(t time.Time) bool { + if c.minute&(1< 0 { + problems = append(problems, where+": a scheduled container cannot also declare restart-on; "+ + "it runs to completion on its cadence rather than staying running to be restarted") + } + if c.Schedule != "" { + if _, err := ParseCron(c.Schedule); err != nil { + problems = append(problems, where+": "+err.Error()) + } + } return append(problems, checkImage(where, c.Image)...) } diff --git a/internal/declaration/schedule_declaration_test.go b/internal/declaration/schedule_declaration_test.go new file mode 100644 index 0000000..d87d1fb --- /dev/null +++ b/internal/declaration/schedule_declaration_test.go @@ -0,0 +1,53 @@ +package declaration + +import ( + "strings" + "testing" +) + +// The host refuses a schedule it does not fully understand, on arrival, for the same reason the +// control plane refuses it near its author (novox/hq ADR 0053): a declaration half-understood is +// refused whole rather than half-applied. + +func schedulePinned() string { return "registry.example/runtime@sha256:" + strings.Repeat("a", 64) } + +func TestAScheduledContainerCarriesTheCronField(t *testing.T) { + d, err := Parse([]byte(`{"declaration":1,"resources":[ + {"id":"sync","type":"container","name":"sync","image":"` + schedulePinned() + `","schedule":"0 3 * * *"} + ]}`)) + if err != nil { + t.Fatalf("a valid scheduled container was refused: %v", err) + } + c, ok := d.Resources[0].(*Container) + if !ok { + t.Fatalf("the scheduled resource is not a container: %T", d.Resources[0]) + } + if c.Schedule != "0 3 * * *" { + t.Errorf("the schedule did not survive parsing: %q", c.Schedule) + } +} + +func TestAMalformedScheduleIsRefused(t *testing.T) { + _, err := Parse([]byte(`{"declaration":1,"resources":[ + {"id":"sync","type":"container","name":"sync","image":"` + schedulePinned() + `","schedule":"every night"} + ]}`)) + if err == nil { + t.Fatal("a container with a malformed schedule was accepted") + } + if !strings.Contains(err.Error(), "cron") && !strings.Contains(err.Error(), "field") { + t.Errorf("refused for the wrong reason: %v", err) + } +} + +func TestAContainerCannotBeBothRunOnceAndScheduled(t *testing.T) { + // A container runs once and gates, or on a cadence, or stays up — never two (novox/hq ADR 0053). + _, err := Parse([]byte(`{"declaration":1,"resources":[ + {"id":"sync","type":"container","name":"sync","image":"` + schedulePinned() + `","run-once":true,"schedule":"0 3 * * *"} + ]}`)) + if err == nil { + t.Fatal("a container that was both run-once and scheduled was accepted") + } + if !strings.Contains(err.Error(), "run-once") && !strings.Contains(err.Error(), "scheduled") { + t.Errorf("refused for the wrong reason: %v", err) + } +}