Files
mesh-host/internal/apply/schedule_test.go
T
jschoubben 5dd439df50 inspect by kind, not the ambiguous bare form — a same-named network stops a container from ever being found
docker inspect <name> resolves across every object kind, not just
containers. A module regularly names a network the same as the
container that joins it (keycloak does this today, ordinarily) — so
when the container does not exist yet but the same-named network
already does, the bare form answers with the network's JSON instead
of reporting the container absent, and the template these callers use
(.State.Running) fails to execute against it entirely.

Live on novox tonight: minio's LB container, named the same as its
network ("minio"), could never be created — every apply crashed on
"the container runtime could not say whether minio is here", stuck
since first push, because the check itself never got a clean answer.

Fixed at every call site asking a container's state by name
(containerState, inspectFound, NamesFree, raiseGiteaServer,
containerRunning) by scoping to `docker container inspect`, matching
the type-scoped form this codebase already uses correctly for
networks, volumes and images elsewhere. Also scoped the one image
inspect that was still bare (publish.go), for the same reason.

mesh-host runs as a host-level service (nox-mesh-host.service), not a
Docker module — merging this does not redeploy it. The live novox
failure persists until the service itself is rebuilt and updated.
2026-09-24 19:50:16 +02:00

482 lines
16 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 "container":
// 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, nil)
// 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, "* * * * *"),
}}, nil)
// 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
}}, nil)
// 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, "* * * * *"),
}}, nil)
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],
}}, nil)
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)
}
}
// Defends novox/hq ADR 0103: a scheduled step of a module not yet taken is not armed — run on its
// cadence it would work on the predecessor's data.
func TestAScheduledStepOfAnUntakenModuleIsNotArmed(t *testing.T) {
clock := &fixedClock{now: time.Date(2026, 9, 7, 12, 0, 30, 0, time.UTC)}
rec := &recordingRun{}
var said []string
s := NewScheduler(clock, rec.run, func(line string) { said = append(said, line) })
d := &declaration.Declaration{Version: 1,
Adoption: &declaration.Adoption{Untaken: map[string][]string{"backups": {"sync"}}},
Resources: []declaration.Resource{
scheduledContainer(t, "* * * * *"),
}}
s.Sync(d, nil)
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC))
s.Wait()
if rec.fireCount() != 0 {
t.Errorf("a held module's scheduled step ran %d time(s)", rec.fireCount())
}
if len(said) == 0 || !strings.Contains(said[0], "held as found") {
t.Errorf("nothing said why the step was not armed: %v", said)
}
// Held by id — a step whose own container was found on the machine.
s2 := NewScheduler(clock, rec.run, nil)
s2.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
scheduledContainer(t, "* * * * *"),
}}, map[string]bool{"sync": true})
s2.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC))
s2.Wait()
if rec.fireCount() != 0 {
t.Errorf("a held scheduled step ran %d time(s)", rec.fireCount())
}
// Taken: the same step is armed and runs.
taken := &declaration.Declaration{Version: 1, Adoption: &declaration.Adoption{Taken: []string{"backups"}},
Resources: []declaration.Resource{scheduledContainer(t, "* * * * *")}}
s.Sync(taken, nil)
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 2, 5, 0, time.UTC))
s.Wait()
if rec.fireCount() == 0 {
t.Error("a taken module's scheduled step never ran")
}
}