package apply import ( "context" "errors" "os" "path/filepath" "slices" "strings" "sync" "testing" "time" "github.com/novox/mesh-host/internal/declaration" "github.com/novox/mesh-host/internal/store" ) // A scheduled step may hold its module's own containers still while it runs (novox/hq ADR 0189, // issue 108). // // What it exists for: the artifact store's collector walks the storage and requires every writer // stopped. A run-once step runs beside containers and a scheduled one is the same container again, // so the mesh had no way to say it — which is why the store it inherited has never collected // anything. The risk the field brings is one shape only: a window that opens and never closes. // Every test here is about that shape. // windowRun records the order of stop / run / start, which is the whole of what is being asserted. type windowRun struct { mu sync.Mutex order []string failAt string // the arg[0] that should fail ("run" makes the step fail) wontGo string // a container name that refuses to start again } func (w *windowRun) run(_ context.Context, _ string, args ...string) (string, error) { w.mu.Lock() defer w.mu.Unlock() switch args[0] { case "info": return "27.0\n", nil case "stop", "start": w.order = append(w.order, args[0]+" "+args[1]) if args[0] == "start" && args[1] == w.wontGo { return "", errors.New("the runtime refused") } case "run": w.order = append(w.order, "run") if w.failAt == "run" { return "", errors.New("the step exited non-zero") } } return "", nil } func (w *windowRun) seen() []string { w.mu.Lock() defer w.mu.Unlock() return append([]string{}, w.order...) } // aStoreWithACollector is a module in the shape distribution has: a server that must not be // writing, and a nightly step that walks its storage with the server held still. func aStoreWithACollector(t *testing.T) *declaration.Declaration { t.Helper() return parseTrusted(t, `{"declaration":1,"resources":[ {"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"}, {"id":"collect","type":"container","name":"mesh-registry-collect","image":"`+pinned+`", "schedule":"30 3 * * *","while-stopped":["store"]} ]}`) } func fireOnce(t *testing.T, d *declaration.Declaration, w *windowRun) { t.Helper() clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)} s := NewScheduler(clock, w.run, func(string) {}) s.Sync(d, nil) s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC)) s.Wait() } func TestAScheduledStepHoldsItsModulesContainerStillAndStartsItAgain(t *testing.T) { w := &windowRun{} fireOnce(t, aStoreWithACollector(t), w) got := w.seen() want := []string{"stop mesh-registry", "run", "start mesh-registry"} var kept []string for _, line := range got { if strings.HasPrefix(line, "stop mesh-registry-collect") { // Clearing the step's own exited container by name; not part of the window. continue } kept = append(kept, line) } if len(kept) != len(want) { t.Fatalf("the window was not stop, run, start: %v", got) } for i := range want { if kept[i] != want[i] { t.Fatalf("the window was %v, want %v", kept, want) } } } // The one that matters: a step that fails must leave the service running. func TestAFailedStepStillClosesTheWindow(t *testing.T) { w := &windowRun{failAt: "run"} fireOnce(t, aStoreWithACollector(t), w) var started bool for _, line := range w.seen() { if line == "start mesh-registry" { started = true } } if !started { t.Fatalf("the step failed and the container it held still was never started again: %v", w.seen()) } } // A container that will not come back is said loudly: it is down, and nothing else notices until // the next apply compares it. func TestAContainerThatWillNotStartAgainIsSaidLoudly(t *testing.T) { w := &windowRun{wontGo: "mesh-registry"} var said []string clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)} s := NewScheduler(clock, w.run, func(line string) { said = append(said, line) }) s.Sync(aStoreWithACollector(t), nil) s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC)) s.Wait() var loud bool for _, line := range said { if strings.Contains(line, "WILL NOT START AGAIN") && strings.Contains(line, "mesh-registry") { loud = true } } if !loud { t.Fatalf("a service left stopped by a maintenance window was not said loudly: %v", said) } } // Several containers come back in the reverse of the order they were stopped: a module names the // dependant first, and starting it before what it depends on is not bringing it back. func TestTheWindowClosesInTheReverseOfTheOrderItOpened(t *testing.T) { d := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"web","type":"container","name":"web","image":"`+pinned+`"}, {"id":"db","type":"container","name":"db","image":"`+pinned+`"}, {"id":"collect","type":"container","name":"collect","image":"`+pinned+`", "schedule":"30 3 * * *","while-stopped":["web","db"]} ]}`) w := &windowRun{} fireOnce(t, d, w) var stops, starts []string for _, line := range w.seen() { switch { case line == "stop web" || line == "stop db": stops = append(stops, line) case strings.HasPrefix(line, "start "): starts = append(starts, line) } } if len(stops) != 2 || stops[0] != "stop web" || stops[1] != "stop db" { t.Fatalf("stopped in %v, want the order the step named them", stops) } if len(starts) != 2 || starts[0] != "start db" || starts[1] != "start web" { t.Fatalf("started in %v, want the reverse", starts) } } // And the refusals, each for what it says rather than that it says something. func TestAMaintenanceWindowIsRefusedWhereItCannotMean(t *testing.T) { for _, c := range []struct{ name, body, says string }{ { "a window with no schedule", `{"id":"collect","type":"container","name":"c","image":"` + pinned + `","while-stopped":["store"]}`, "needs a schedule", }, { "a window naming itself", `{"id":"collect","type":"container","name":"c","image":"` + pinned + `","schedule":"30 3 * * *","while-stopped":["collect"]}`, "this step itself", }, { "a window naming something that is not a container here", `{"id":"collect","type":"container","name":"c","image":"` + pinned + `","schedule":"30 3 * * *","while-stopped":["elsewhere"]}`, "no container by that id", }, } { _, err := declaration.ParseTrusted([]byte(`{"declaration":1,"resources":[` + c.body + `]}`)) if err == nil { t.Errorf("%s was accepted", c.name) continue } if !strings.Contains(err.Error(), c.says) { t.Errorf("%s: the refusal does not say %q: %v", c.name, c.says, err) } } } // A changed window is a changed declaration, and the install says so. // // The cadence already works this way: "a changed schedule is a changed spec — the marker moves and // the install is reported updated and re-established" (containerSpec). Which containers are held // still for the run is the same kind of statement, and a declaration that changed it while the // machine reported no change would be a machine quietly running the old window. func TestAChangedWindowMovesTheSpec(t *testing.T) { one := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"}, {"id":"other","type":"container","name":"other","image":"`+pinned+`"}, {"id":"collect","type":"container","name":"collect","image":"`+pinned+`", "schedule":"30 3 * * *","while-stopped":["store"]} ]}`) two := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"}, {"id":"other","type":"container","name":"other","image":"`+pinned+`"}, {"id":"collect","type":"container","name":"collect","image":"`+pinned+`", "schedule":"30 3 * * *","while-stopped":["store","other"]} ]}`) stepOf := func(d *declaration.Declaration) *declaration.Container { for _, r := range d.Resources { if c, ok := r.(*declaration.Container); ok && c.ID == "collect" { return c } } t.Fatal("no step in the fixture") return nil } if containerSpec(stepOf(one), inputs{}) == containerSpec(stepOf(two), inputs{}) { t.Fatal("the window changed and the spec did not; the machine would report no change " + "and keep holding the containers it held yesterday") } // And a container with no window is untouched by the field existing at all. plain := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"} ]}`) spec := containerSpec(plain.Resources[0].(*declaration.Container), inputs{}) if strings.Contains(spec, "while-stopped") || strings.Contains(spec, "held") { t.Errorf("an ordinary container's spec mentions a field it does not set:\n%s", spec) } } // ---- An apply arriving during a window (novox/hq issue 224) ---- // // A window holds a container stopped on purpose; the apply's rule for a stopped container is to // replace it. Until the apply could see the window, a push landing at 03:30 recreated the registry // in the middle of its collection — the collector, already past its mark phase, then swept a blob // that an upload had just written. These drive the scheduler and the apply against one pretend // runtime, so what is asserted is what one does to the other. // aRegistryMachine is a runtime holding one long-running container, mesh-registry, and running the // step in the foreground until the test lets it finish. type aRegistryMachine struct { mu sync.Mutex exists bool running bool spec string calls []string wontGo bool // the registry refuses to start again after the window failRun bool // the step exits non-zero stepIn chan struct{} // the step has started: the window is open stepGo chan struct{} // let the step finish } func newRegistryMachine(spec string) *aRegistryMachine { return &aRegistryMachine{exists: true, running: true, spec: spec, stepIn: make(chan struct{}, 4), stepGo: make(chan struct{})} } func (m *aRegistryMachine) run(_ context.Context, _ string, args ...string) (string, error) { m.mu.Lock() switch args[0] { case "info", "image": m.mu.Unlock() return "27.0\n", nil case "stop": m.calls = append(m.calls, "stop "+args[1]) if args[1] == "mesh-registry" { m.running = false } case "start": m.calls = append(m.calls, "start "+args[1]) if args[1] == "mesh-registry" { if m.wontGo { m.mu.Unlock() return "", errors.New("the runtime refused") } m.running = true } case "container": // inspect name := args[len(args)-1] if name != "mesh-registry" || !m.exists { m.mu.Unlock() return "", errors.New("no such container") } out := "false\t" + m.spec if m.running { out = "true\t" + m.spec } m.mu.Unlock() return out, nil case "rm": name := args[len(args)-1] m.calls = append(m.calls, "rm "+name) if name == "mesh-registry" { m.exists, m.running = false, false } case "run": if !slices.Contains(args, "--detach") { // The step, in the foreground, for as long as the test says. m.calls = append(m.calls, "step") fail := m.failRun m.mu.Unlock() m.stepIn <- struct{}{} <-m.stepGo if fail { return "", errors.New("the step exited non-zero") } return "", nil } m.calls = append(m.calls, "create "+args[slices.Index(args, "--name")+1]) for i, a := range args { if a == "--label" && strings.HasPrefix(args[i+1], specLabel+"=") { m.spec = strings.TrimPrefix(args[i+1], specLabel+"=") } } m.exists, m.running = true, true } m.mu.Unlock() return "", nil } // touched is every call since mark that recreated or removed the registry. func (m *aRegistryMachine) touched(mark int) []string { m.mu.Lock() defer m.mu.Unlock() var out []string for _, c := range m.calls[mark:] { if c == "rm mesh-registry" || c == "create mesh-registry" { out = append(out, c) } } return out } func (m *aRegistryMachine) mark() int { m.mu.Lock() defer m.mu.Unlock() return len(m.calls) } // theStoreMoved is the store's module with the server's declaration changed — a push that landed // at 03:30 with a new registry image, say. func theStoreMoved(t *testing.T) *declaration.Declaration { t.Helper() return parseTrusted(t, `{"declaration":1,"resources":[ {"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`","env":{"NEW":"1"}}, {"id":"collect","type":"container","name":"mesh-registry-collect","image":"`+pinned+`", "schedule":"30 3 * * *","while-stopped":["store"]} ]}`) } // openAWindow fires the collector and returns once it is running, the registry held still. func openAWindow(t *testing.T, d *declaration.Declaration, m *aRegistryMachine, windows *Windows) *Scheduler { t.Helper() clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)} s := NewScheduler(clock, m.run, func(string) {}) s.RecordWindowsIn(windows) s.Sync(d, nil) s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC)) select { case <-m.stepIn: case <-time.After(5 * time.Second): t.Fatal("the step never started") } return s } func TestAnApplyDuringAWindowLeavesWhatItHoldsAlone(t *testing.T) { d := aStoreWithACollector(t) want := containerSpec(d.Resources[0].(*declaration.Container), inputs{}) m := newRegistryMachine(want) windows := WindowsIn(t.TempDir()) s := openAWindow(t, d, m, windows) mark := m.mark() // The push lands mid-window, and with it a registry whose declaration moved: the very shape // that, before this, removed the server and created it running under the collector. report, state, err := ApplyMindingWindows(context.Background(), archHost(t), theStoreMoved(t), store.State{}, store.OriginDeclared, m.run, nil, nil, nil, windows) if err != nil { t.Fatalf("the apply failed: %v", err) } if got := m.touched(mark); len(got) > 0 { t.Fatalf("an apply during the window touched the container it holds still: %v", got) } o := outcomeOf(report, "store") if o.Action != heldStill { t.Fatalf("a held container was reported %q, want %q: %+v", o.Action, heldStill, o) } for _, says := range []string{"collect", "maintenance window", "first apply after", "declaration moved"} { if !strings.Contains(o.Detail, says) { t.Errorf("the outcome does not say %q: %s", says, o.Detail) } } // Held still is not a change: nothing about the machine moved (the step's own install, first // time here, is). if (Report{Outcomes: []Outcome{o}}).Changed() { t.Errorf("a container left alone reads as a change: %+v", o) } // The report says a window is open, and which. if len(report.Windows) != 1 || report.Windows[0].Step != "collect" || len(report.Windows[0].Holds) != 1 || report.Windows[0].Holds[0] != "mesh-registry" { t.Errorf("the report does not say the collector's window is open: %+v", report.Windows) } // The window closes: the registry is running again and the record is gone. close(m.stepGo) s.Wait() if open := windows.Open(time.Now()); len(open) != 0 { t.Fatalf("the window closed and is still recorded open: %+v", open) } // And the next apply converges it by the ordinary rule — the changed declaration recreates it. mark = m.mark() report, _, err = ApplyMindingWindows(context.Background(), archHost(t), theStoreMoved(t), state, store.OriginDeclared, m.run, nil, nil, nil, windows) if err != nil { t.Fatalf("the apply after the window failed: %v", err) } if got := m.touched(mark); len(got) != 2 || got[0] != "rm mesh-registry" || got[1] != "create mesh-registry" { t.Fatalf("the apply after the window did not converge the container: %v", got) } if o := outcomeOf(report, "store"); o.Action != "updated" { t.Errorf("the container was reported %q after the window, want updated: %+v", o.Action, o) } if len(report.Windows) != 0 { t.Errorf("a closed window is still reported open: %+v", report.Windows) } } // A step that fails releases its hold like one that succeeds: the record is erased in the same // defer that starts the server again, and an apply afterwards does what it always did — here, to a // registry that would not come back, the one thing that brings it back. func TestAFailedStepReleasesItsHold(t *testing.T) { d := aStoreWithACollector(t) m := newRegistryMachine(containerSpec(d.Resources[0].(*declaration.Container), inputs{})) m.failRun, m.wontGo = true, true windows := WindowsIn(t.TempDir()) s := openAWindow(t, d, m, windows) if open := windows.Open(time.Now()); len(open) != 1 { t.Fatalf("the window is not recorded while the step runs: %+v", open) } close(m.stepGo) s.Wait() if open := windows.Open(time.Now()); len(open) != 0 { t.Fatalf("a failed step left its window recorded open: %+v", open) } mark := m.mark() report, _, err := ApplyMindingWindows(context.Background(), archHost(t), d, store.State{}, store.OriginDeclared, m.run, nil, nil, nil, windows) if err != nil { t.Fatalf("the apply after the failed step failed: %v", err) } if got := m.touched(mark); len(got) != 2 { t.Fatalf("the registry left stopped by a failed window was not brought back: %v", got) } if o := outcomeOf(report, "store"); o.Action == heldStill { t.Errorf("the hold outlived the step: %+v", o) } } // A window whose host died mid-step — crashed, or stood aside for a successor — holds nothing: // the process that would start the server again is gone, so the apply must be free to. And a window // nothing closed ends on its own. func TestAWindowNothingClosesDoesNotHoldForEver(t *testing.T) { dir := t.TempDir() alive := true windows := &Windows{dir: dir, alive: func(int) bool { return alive }} now := time.Now() if err := windows.open("collect", []string{"mesh-registry"}, now); err != nil { t.Fatal(err) } if _, held := windows.holding("mesh-registry", now); !held { t.Fatal("an open window with a live host does not hold its container") } if _, held := windows.holding("mesh-registry", now.Add(windowAtMost)); held { t.Error("a window past its end still holds") } if err := windows.open("collect", []string{"mesh-registry"}, now); err != nil { t.Fatal(err) } alive = false if _, held := windows.holding("mesh-registry", now); held { t.Error("a window whose host is gone still holds") } if entries, _ := os.ReadDir(dir); len(entries) != 0 { t.Errorf("a window read as closed was left on disk: %v", entries) } // The real liveness check: this process is alive, and a pid nothing can have is not. if !processAlive(os.Getpid()) || processAlive(0) { t.Error("processAlive does not tell a live process from none") } } // A window that cannot be recorded is not opened: an apply would not know, and would recreate the // server under the step. One night's collection is the cheaper loss. func TestAWindowThatCannotBeRecordedIsNotOpened(t *testing.T) { notADir := filepath.Join(t.TempDir(), "state") if err := os.WriteFile(notADir, nil, 0o600); err != nil { t.Fatal(err) } d := aStoreWithACollector(t) m := newRegistryMachine(containerSpec(d.Resources[0].(*declaration.Container), inputs{})) var said []string var mu sync.Mutex clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)} s := NewScheduler(clock, m.run, func(line string) { mu.Lock(); said = append(said, line); mu.Unlock() }) s.RecordWindowsIn(WindowsIn(notADir)) s.Sync(d, nil) s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC)) s.Wait() for _, c := range m.calls { if c == "stop mesh-registry" || c == "step" { t.Fatalf("a window that could not be recorded was opened anyway: %v", m.calls) } } mu.Lock() defer mu.Unlock() if len(said) == 0 || !strings.Contains(strings.Join(said, "\n"), "not run") { t.Errorf("the step not running was not said: %v", said) } }