Two faults that both reported success while being wrong, found while proving the firewall module actually delivers. A unit whose job is to apply something and exit — load a rule set, set a sysctl — is inactive the instant it succeeds. Reading that as stopped made it permanently unsatisfiable: the host started it, it worked, the host read back stopped and reported the machine as not doing what it was told, on every apply, for ever, with the rules correctly in place the whole time. That is what the firewall has been doing on every machine it was assigned to, and why the four-machine bed was red. And a container took its identity from its own fields, not from the files it reads. A file written in an earlier apply — or before the container declared it as a dependency — left a process holding a credential the mesh had already replaced, with everything reporting success (novox/hq 04-ISSUES/045). What a container reads is now part of what it is, so the comparison is a standing one rather than a tripwire that fires during one apply and never again.
251 lines
9.9 KiB
Go
251 lines
9.9 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
|
|
|
|
// 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, nil)
|
|
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())
|
|
}
|
|
}
|
|
}
|