482 lines
16 KiB
Go
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 "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, 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")
|
|
}
|
|
}
|