Files
mesh-host/internal/apply/schedule.go
T
jschoubben e8c4824ae2 A scheduled step may hold its module's own containers still (hq ADR 0189)
while-stopped names resource ids of the same module's containers; the host
stops them before the run and starts them again after it, in reverse order,
whatever the step did. The restart is deferred before the first stop and runs
on its own context, because the one real risk of this field is a window that
never closes.

Scheduled steps only: at apply the declaration is applied in order and a
run-once step already gates what follows.
2026-10-04 02:34:33 +02:00

348 lines
14 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
// hold is the runtime names of the containers held still for the duration of a run, in the
// order the step named them (novox/hq ADR 0189).
hold []string
}
// heldStillFor is the runtime names of the containers a step holds still, resolved from ids.
func heldStillFor(step *declaration.Container, d *declaration.Declaration) []string {
if len(step.WhileStopped) == 0 {
return nil
}
byID := map[string]string{}
for _, r := range d.Resources {
if c, ok := r.(*declaration.Container); ok {
byID[c.Identity()] = c.Name
}
}
out := make([]string, 0, len(step.WhileStopped))
for _, id := range step.WhileStopped {
if name := byID[id]; name != "" {
out = append(out, name)
}
}
return out
}
// 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.
// held is the ids this node holds as found — what an adopted node keeps until its module is taken
// (novox/hq ADR 0100). A step of a module not yet taken is not armed: run on its cadence it would
// work on the predecessor's data, under the predecessor's service, which is the one thing an
// adopted node must not do. Nil on a converged node, where nothing is held.
func (s *Scheduler) Sync(d *declaration.Declaration, held map[string]bool) {
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
}
if module, untaken := d.Adoption.UntakenModuleOf(c.Identity()); untaken || held[c.Identity()] {
if module == "" {
module = "its module"
}
s.log(fmt.Sprintf("scheduled step %s: not armed while %s is held as found on this node",
c.Identity(), module))
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
// Nothing to read: a scheduled container may not declare restart-on — it runs to completion
// on its cadence rather than staying running to be restarted — so its identity cannot
// depend on another resource's content and there is nothing to pass.
spec := containerSpec(c, inputs{})
// The runtime stops containers by name; the declaration names them by id. Resolved here,
// against the declaration this job was armed from, so a fire never has to look anything up
// (novox/hq ADR 0189). The parser has already refused an id that is not a container here.
hold := heldStillFor(c, d)
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
existing.hold = hold
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, hold: hold,
}
}
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
}
// **The window opens here and closes in the defer, whatever happens** (novox/hq ADR 0189).
// Deferred before the first stop so a panic, a failing step or a step that runs long all end
// the same way: the service running. The one real risk of this field is a window that never
// closes, and the only defence against it is that closing is not conditional on anything.
if len(j.hold) > 0 {
defer s.letRun(ctx, cri, j)
s.holdStill(ctx, cri, j)
}
// 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())
}
}
}
// holdStill stops the containers this step runs instead of, in the order it named them.
//
// A stop that fails is said and not fatal. The step runs anyway: for the case this exists for —
// a collector walking storage nothing must be writing to — a writer that would not stop is worth
// knowing about, and refusing to run would mean the work never happens and the log says nothing
// new each night. What must not be skipped is the restart, and it is not: it is deferred.
func (s *Scheduler) holdStill(ctx context.Context, cri string, j *scheduledJob) {
for _, name := range j.hold {
if _, err := s.run(ctx, cri, "stop", name); err != nil {
s.log(fmt.Sprintf("scheduled step %s: could not stop %s for the run: %v", j.id, name, err))
continue
}
s.log(fmt.Sprintf("scheduled step %s: %s held still for the run", j.id, name))
}
}
// letRun starts them again, in the reverse of the order they were stopped, and says so loudly if
// one does not come back.
//
// **Reverse order**, because stopping walks a dependency the other way: a module that holds two
// containers still names the one that depends on the other first, and bringing them back the same
// way would start a dependant before what it depends on.
//
// Given its own context, because this runs in a defer and the one the run used may already be
// cancelled — a host shutting down mid-window would otherwise leave the service stopped, which is
// precisely the outcome this field must never have.
func (s *Scheduler) letRun(_ context.Context, cri string, j *scheduledJob) {
ctx, cancel := context.WithTimeout(context.Background(), closingWindow)
defer cancel()
for i := len(j.hold) - 1; i >= 0; i-- {
name := j.hold[i]
if _, err := s.run(ctx, cri, "start", name); err != nil {
// Said as loudly as this host says anything: a service the mesh stopped for a
// maintenance window and could not start again is down, and nothing else will notice
// until the next apply compares it.
s.log(fmt.Sprintf(
"scheduled step %s: %s was held still for the run and WILL NOT START AGAIN: %v",
j.id, name, err))
continue
}
s.log(fmt.Sprintf("scheduled step %s: %s running again", j.id, name))
}
}
// closingWindow is how long the host will spend putting back what it stopped. Generous: this is
// the half that must not be given up on.
const closingWindow = 5 * time.Minute