The recurring twin of run-once, one modifier over: a container marked schedule: "<cron>" 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
247 lines
9.6 KiB
Go
247 lines
9.6 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. 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())
|
|
}
|
|
}
|
|
}
|