Files
mesh-host/internal/apply/maintenance_window_test.go
T
jochen f08681f225 An apply leaves a container a maintenance window holds still (hq issue 224)
A while-stopped step stops its module's containers, and an apply arriving
mid-window read the stopped server as broken and recreated it running under
the step - for the store's collector, a registry taking an upload the sweep
then deletes. The scheduler now records each window under the node's state
directory before its first stop and erases it after its last start, so every
apply on the machine (daemon, reconcile by hand, installer) reports a held
container held-still and leaves it for the first apply after the window.
A window whose host died, or past six hours, holds nothing; the report says
which windows are open.
2026-10-05 18:27:04 +02:00

543 lines
19 KiB
Go

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)
}
}