Merge pull request 'A scheduled step may hold its module's own containers still (hq ADR 0189, issue 108)' (#77) from feat/the-store-keeps-what-the-records-name into main

This commit was merged in pull request #77.
This commit is contained in:
2026-10-04 01:39:07 +00:00
4 changed files with 399 additions and 1 deletions
+7
View File
@@ -1707,6 +1707,13 @@ func containerSpecReading(r *declaration.Container, declares, reads map[string]s
// The cadence is part of what was declared, so a changed schedule is a changed spec — the marker // The cadence is part of what was declared, so a changed schedule is a changed spec — the marker
// moves and the install is reported "updated" and re-established. Added only when present, so no // moves and the install is reported "updated" and re-established. Added only when present, so no
// ordinary container's or run-once step's digest moves for a field it does not set. // ordinary container's or run-once step's digest moves for a field it does not set.
// And which of its module's containers it holds still while it runs (novox/hq ADR 0189), for
// the same reason: a declaration that changed the window while the machine reported no change
// would be a machine quietly holding yesterday's containers. In declared order, which is the
// order they are stopped in.
for _, id := range r.WhileStopped {
b.WriteString("while-stopped " + id + "\n")
}
if r.Schedule != "" { if r.Schedule != "" {
b.WriteString("schedule " + r.Schedule + "\n") b.WriteString("schedule " + r.Schedule + "\n")
} }
+238
View File
@@ -0,0 +1,238 @@
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)
}
}
}
// 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)
}
}
+86 -1
View File
@@ -65,6 +65,29 @@ type scheduledJob struct {
container *declaration.Container container *declaration.Container
next time.Time // the next minute at which it is due 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 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. // 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 // 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. // depend on another resource's content and there is nothing to pass.
spec := containerSpec(c, inputs{}) 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 { if existing := s.jobs[c.Identity()]; existing != nil && existing.spec == spec {
// Unchanged: keep where it is in its cadence, refresh the declaration pointer only. // Unchanged: keep where it is in its cadence, refresh the declaration pointer only.
existing.container = c existing.container = c
existing.hold = hold
continue continue
} }
// New or changed: arm it for the next due minute after now. // New or changed: arm it for the next due minute after now.
next, _ := cron.Next(s.clock.Now()) next, _ := cron.Next(s.clock.Now())
s.jobs[c.Identity()] = &scheduledJob{ 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 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 // 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. // 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) _, _ = 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
+68
View File
@@ -1037,6 +1037,24 @@ type Container struct {
// rather than stacked. It is exclusive with RunOnce and with restart-on: a container runs once // 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. // and gates, runs on a cadence, or stays up — never two of these.
Schedule string `json:"schedule,omitempty"` 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 } func (c *Container) Identity() string { return c.ID }
@@ -1068,6 +1086,20 @@ func (c *Container) validate(where string, _ bool) []string {
problems = append(problems, where+": "+err.Error()) 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 // 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. // refused at create — after the old container was already removed. Refused on arrival instead.
for _, d := range c.Dns { for _, d := range c.Dns {
@@ -1482,6 +1514,7 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
problems = append(problems, resource.validate(where, allowActions)...) problems = append(problems, resource.validate(where, allowActions)...)
d.Resources = append(d.Resources, resource) d.Resources = append(d.Resources, resource)
} }
problems = append(problems, checkWhileStopped(d.Resources)...)
problems = append(problems, checkAdoption(env.Adoption, d.Resources, allowActions)...) problems = append(problems, checkAdoption(env.Adoption, d.Resources, allowActions)...)
problems = append(problems, checkReplaces(d.Resources)...) problems = append(problems, checkReplaces(d.Resources)...)
if env.Adoption == nil { if env.Adoption == nil {
@@ -1511,6 +1544,41 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
return d, nil 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 { func strictDecode(raw []byte, into any) error {
// DisallowUnknownFields is the whole point rather than strictness for its own sake: a // 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. // field the host does not know is a thing the control plane believes it asked for.