diff --git a/internal/apply/maintenance_window_test.go b/internal/apply/maintenance_window_test.go new file mode 100644 index 0000000..1649207 --- /dev/null +++ b/internal/apply/maintenance_window_test.go @@ -0,0 +1,196 @@ +package apply + +import ( + "context" + "errors" + "strings" + "sync" + "testing" + "time" + + "github.com/novox/mesh-host/internal/declaration" +) + +// 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) + } + } +} diff --git a/internal/apply/schedule.go b/internal/apply/schedule.go index 6c80de1..179d3c4 100644 --- a/internal/apply/schedule.go +++ b/internal/apply/schedule.go @@ -65,6 +65,29 @@ type scheduledJob struct { container *declaration.Container next time.Time // the next minute at which it is due running bool // a run is in flight — the next due run is skipped rather than stacked + // hold is the runtime names of the containers held still for the duration of a run, in the + // order the step named them (novox/hq ADR 0189). + hold []string +} + +// heldStillFor is the runtime names of the containers a step holds still, resolved from ids. +func heldStillFor(step *declaration.Container, d *declaration.Declaration) []string { + if len(step.WhileStopped) == 0 { + return nil + } + byID := map[string]string{} + for _, r := range d.Resources { + if c, ok := r.(*declaration.Container); ok { + byID[c.Identity()] = c.Name + } + } + out := make([]string, 0, len(step.WhileStopped)) + for _, id := range step.WhileStopped { + if name := byID[id]; name != "" { + out = append(out, name) + } + } + return out } // NewScheduler builds a scheduler. A nil clock is the system clock; a nil log says nothing. @@ -123,15 +146,20 @@ func (s *Scheduler) Sync(d *declaration.Declaration, held map[string]bool) { // on its cadence rather than staying running to be restarted — so its identity cannot // depend on another resource's content and there is nothing to pass. spec := containerSpec(c, inputs{}) + // The runtime stops containers by name; the declaration names them by id. Resolved here, + // against the declaration this job was armed from, so a fire never has to look anything up + // (novox/hq ADR 0189). The parser has already refused an id that is not a container here. + hold := heldStillFor(c, d) if existing := s.jobs[c.Identity()]; existing != nil && existing.spec == spec { // Unchanged: keep where it is in its cadence, refresh the declaration pointer only. existing.container = c + existing.hold = hold continue } // New or changed: arm it for the next due minute after now. next, _ := cron.Next(s.clock.Now()) s.jobs[c.Identity()] = &scheduledJob{ - id: c.Identity(), spec: spec, cron: cron, container: c, next: next, + id: c.Identity(), spec: spec, cron: cron, container: c, next: next, hold: hold, } } @@ -202,6 +230,15 @@ func (s *Scheduler) fire(ctx context.Context, j *scheduledJob) { return } + // **The window opens here and closes in the defer, whatever happens** (novox/hq ADR 0189). + // Deferred before the first stop so a panic, a failing step or a step that runs long all end + // the same way: the service running. The one real risk of this field is a window that never + // closes, and the only defence against it is that closing is not conditional on anything. + if len(j.hold) > 0 { + defer s.letRun(ctx, cri, j) + s.holdStill(ctx, cri, j) + } + // A container by this name left exited by the previous run would collide with --name. Removing // one that is not there is the state we want, so its error is ignored — the same as run-once. _, _ = s.run(ctx, cri, "rm", "-f", j.container.Name) @@ -260,3 +297,51 @@ func (s *Scheduler) Run(ctx context.Context) { } } } + +// holdStill stops the containers this step runs instead of, in the order it named them. +// +// A stop that fails is said and not fatal. The step runs anyway: for the case this exists for — +// a collector walking storage nothing must be writing to — a writer that would not stop is worth +// knowing about, and refusing to run would mean the work never happens and the log says nothing +// new each night. What must not be skipped is the restart, and it is not: it is deferred. +func (s *Scheduler) holdStill(ctx context.Context, cri string, j *scheduledJob) { + for _, name := range j.hold { + if _, err := s.run(ctx, cri, "stop", name); err != nil { + s.log(fmt.Sprintf("scheduled step %s: could not stop %s for the run: %v", j.id, name, err)) + continue + } + s.log(fmt.Sprintf("scheduled step %s: %s held still for the run", j.id, name)) + } +} + +// letRun starts them again, in the reverse of the order they were stopped, and says so loudly if +// one does not come back. +// +// **Reverse order**, because stopping walks a dependency the other way: a module that holds two +// containers still names the one that depends on the other first, and bringing them back the same +// way would start a dependant before what it depends on. +// +// Given its own context, because this runs in a defer and the one the run used may already be +// cancelled — a host shutting down mid-window would otherwise leave the service stopped, which is +// precisely the outcome this field must never have. +func (s *Scheduler) letRun(_ context.Context, cri string, j *scheduledJob) { + ctx, cancel := context.WithTimeout(context.Background(), closingWindow) + defer cancel() + for i := len(j.hold) - 1; i >= 0; i-- { + name := j.hold[i] + if _, err := s.run(ctx, cri, "start", name); err != nil { + // Said as loudly as this host says anything: a service the mesh stopped for a + // maintenance window and could not start again is down, and nothing else will notice + // until the next apply compares it. + s.log(fmt.Sprintf( + "scheduled step %s: %s was held still for the run and WILL NOT START AGAIN: %v", + j.id, name, err)) + continue + } + s.log(fmt.Sprintf("scheduled step %s: %s running again", j.id, name)) + } +} + +// closingWindow is how long the host will spend putting back what it stopped. Generous: this is +// the half that must not be given up on. +const closingWindow = 5 * time.Minute diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index c7419b4..020e105 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -997,6 +997,24 @@ type Container struct { // rather than stacked. It is exclusive with RunOnce and with restart-on: a container runs once // and gates, runs on a cadence, or stays up — never two of these. Schedule string `json:"schedule,omitempty"` + + // WhileStopped names resources of the same module — containers — that must be held still for + // the duration of this step's run (novox/hq ADR 0189). The host stops each before the run and + // starts each again after it, **whatever the step did**: a step that failed must leave the + // service running, because the one real risk of this field is a window that never closes. + // + // **For the work a service cannot have done underneath it.** 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 until this there was no way for a module + // to say it. The predecessor said it with a shell script, which is how the mesh inherited a + // store that has never collected anything. + // + // **Its own module's containers, and only on a schedule.** A module that could quiesce a + // neighbour could stop the mesh. And at apply time the host already has a window — the + // declaration is applied in order and a run-once step gates what follows — so a one-time + // offline job says *before*, not *instead of*; a recurring window is the case order cannot + // express, and the only one this serves. + WhileStopped []string `json:"while-stopped,omitempty"` } func (c *Container) Identity() string { return c.ID } @@ -1028,6 +1046,20 @@ func (c *Container) validate(where string, _ bool) []string { problems = append(problems, where+": "+err.Error()) } } + // A maintenance window belongs to a recurring step (novox/hq ADR 0189). Refused on anything + // else here, where the field is; that it names containers of the same module, and not itself, + // is judged against the whole declaration (see whileStoppedNames). + if len(c.WhileStopped) > 0 && c.Schedule == "" { + problems = append(problems, where+": while-stopped needs a schedule; at apply the host "+ + "already has a window — the declaration is applied in order and a run-once step gates "+ + "what follows — so a one-time offline job is declared before what it works on") + } + for _, id := range c.WhileStopped { + if id == c.ID { + problems = append(problems, where+": while-stopped names "+strconv.Quote(id)+ + ", which is this step itself") + } + } // The runtime's flags take addresses, and a name here would be handed to it verbatim and // refused at create — after the old container was already removed. Refused on arrival instead. for _, d := range c.Dns { @@ -1438,6 +1470,7 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) { problems = append(problems, resource.validate(where, allowActions)...) d.Resources = append(d.Resources, resource) } + problems = append(problems, checkWhileStopped(d.Resources)...) problems = append(problems, checkAdoption(env.Adoption, d.Resources, allowActions)...) if env.Adoption == nil { for _, r := range d.Resources { @@ -1466,6 +1499,41 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) { return d, nil } +// checkWhileStopped judges a maintenance window against the whole declaration (novox/hq ADR 0189). +// +// A step may hold still only a container that is **here** — in this same declaration, which is to +// say on this machine and placed by the mesh. That is what makes it the module's own: a node's +// declaration carries one module's resources beside another's, so the id must also be a container +// and not a file or a directory, which there would be nothing to stop. +// +// Refused on arrival rather than discovered at the first fire. A window that names something the +// host cannot stop is a window that opens at 03:00 and reports nothing until somebody reads a log. +func checkWhileStopped(resources []Resource) []string { + containers := map[string]bool{} + for _, r := range resources { + if r.Kind() == TypeContainer { + containers[r.Identity()] = true + } + } + var problems []string + for _, r := range resources { + c, ok := r.(*Container) + if !ok { + continue + } + for _, id := range c.WhileStopped { + if containers[id] { + continue + } + problems = append(problems, fmt.Sprintf( + "resource %q: while-stopped names %q, and this declaration has no container by "+ + "that id. A step may hold still only a container placed on this machine "+ + "beside it", c.ID, id)) + } + } + return problems +} + func strictDecode(raw []byte, into any) error { // DisallowUnknownFields is the whole point rather than strictness for its own sake: a // field the host does not know is a thing the control plane believes it asked for.