Files
mesh-host/internal/apply/schedule_test.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

437 lines
14 KiB
Go

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. Installing it ensures the pinned image is present (a
// scheduled step is never run at apply, so nothing else pulls it), but it NEVER runs the
// container — no `docker run` fires it, which is the invariant that matters (novox/hq ADR 0053).
var fired bool // a `docker run` — the container was started
var pulled bool // an image reported present was pulled anyway
run := func(ctx context.Context, name string, args ...string) (string, error) {
switch args[0] {
case "info":
return "27.0\n", nil // a container runtime answers the apply's probe
case "image":
return "", nil // `image inspect`: the pinned image is already present
case "pull":
pulled = true
return "", nil
case "run":
fired = true
return "", errors.New("installing a schedule must not run the container")
}
return "", nil
}
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 fired {
t.Error("installing a schedule ran the container — it installs state, it does not fire it")
}
if pulled {
t.Error("installing a schedule pulled an image it had just found present")
}
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 fires 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)
}
if fired {
t.Error("re-installing a schedule ran the container")
}
}
func TestInstallingAScheduleEnsuresItsImageIsPresentWithoutStartingIt(t *testing.T) {
// A scheduled step is never run at apply, so `docker run` — which is what pulls a service's or a
// run-once step's image — never fetches it. Without an explicit pull the image is absent from the
// node until the first scheduled fire, which then pays the whole pull latency and, until it runs,
// leaves tooling that expects the image present after apply looking at a node without it. So the
// apply ensures the image present: when it is absent it is pulled, and still nothing is run.
imagePresent := false
var pulledImage string
var fired bool
run := func(ctx context.Context, name string, args ...string) (string, error) {
switch args[0] {
case "info":
return "27.0\n", nil
case "image": // `image inspect <image>`
if imagePresent {
return "", nil
}
return "", errors.New("no such image")
case "pull":
pulledImage = args[len(args)-1]
imagePresent = true // a real runtime leaves the image present after a pull
return "", nil
case "run":
fired = true
return "", nil
}
return "", nil
}
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 whose image was absent failed the apply: %v", err)
}
if pulledImage != pinned {
t.Errorf("installing a schedule did not pull its pinned image (pulled %q, want %q)", pulledImage, pinned)
}
if fired {
t.Error("installing a schedule ran the container — ensuring the image is present is not running it")
}
if report.Outcomes[0].Action != "created" {
t.Errorf("an installed schedule was not reported created: %+v", report.Outcomes[0])
}
if _, ok := state.Find("sync"); !ok {
t.Fatal("an installed schedule was not recorded as applied")
}
}
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 ... <name>`.
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)
}
}