diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 16e7266..3fa89e5 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -598,12 +598,15 @@ func runApply(ctx context.Context, opts options, d *declaration.Declaration, raw } fmt.Fprintf(out, "\napplying:\n") - report, updated, applyErr := apply.ApplyKeeping(ctx, sys, d, known, origin, + // Minding the windows the daemon's scheduler has open: this runs as its own process, and a + // `reconcile` by hand at 03:31 is exactly the apply that must not restart a server a collector + // is holding still (novox/hq issue 224). + report, updated, applyErr := apply.ApplyMindingWindows(ctx, sys, d, known, origin, apply.ExecRunner, func(line string) { if !opts.json { fmt.Fprintln(out, line) } - }, sealOpener(opts.state), apply.KeepIn(filepath.Dir(opts.state))) + }, sealOpener(opts.state), apply.KeepIn(filepath.Dir(opts.state)), apply.WindowsIn(filepath.Dir(opts.state))) // What this apply settles about the node, whichever way it went. The mode is what the // declaration said and the check above agreed with; the bundle, once applied, is consumed. @@ -1041,6 +1044,9 @@ func runLink(ctx context.Context, opts options) error { // across the reconcile loop; a host restart rebuilds it from the declaration the node kept, the // first time either path applies. Its own loop is the thing on the clock — no system timer. sched := apply.NewScheduler(apply.SystemClock(), apply.ExecRunner, say) + // A window it opens is written beside the node's state, where every apply on this machine — + // this process's, a `reconcile` run by hand, the installer's — reads it (novox/hq issue 224). + sched.RecordWindowsIn(apply.WindowsIn(filepath.Dir(opts.state))) go sched.Run(ctx) // **Standing aside for a successor happens between reconciles and nowhere else** (novox/hq ADR @@ -1300,7 +1306,7 @@ func worthSaying(report link.Report) bool { return false } return len(report.Held) > 0 || report.Firewall != "" || len(report.Outward) > 0 || - len(report.Filters) > 0 || report.FoundFirewall != nil + len(report.Filters) > 0 || report.FoundFirewall != nil || len(report.Windows) > 0 } // applyDeclared applies a declaration that has already been proved to come from the mesh. @@ -1399,8 +1405,9 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D // // `say` already reaches stdout, and the launcher's unit sends that to the journal, so this needs // no new mechanism — only for the argument to be passed. - outcome, updated, applyErr := apply.ApplyKeeping(ctx, built, declared, known, store.OriginDeclared, - apply.ExecRunner, announceOr(say), sealOpener(opts.state), apply.KeepIn(filepath.Dir(opts.state))) + outcome, updated, applyErr := apply.ApplyMindingWindows(ctx, built, declared, known, store.OriginDeclared, + apply.ExecRunner, announceOr(say), sealOpener(opts.state), apply.KeepIn(filepath.Dir(opts.state)), + apply.WindowsIn(filepath.Dir(opts.state))) // The mode the mesh said, recorded whichever way the apply went: the declaration is kept // either way, and the node is held to it from the next reconcile (novox/hq ADR 0100). @@ -1443,6 +1450,11 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D report.Held = append(report.Held, link.Held{ID: h.ID, Module: h.Module, Kind: h.Kind, Target: h.Target, Since: h.Since, Changed: h.Changed, Kept: h.Kept, Facts: factsAsReported(h.Facts)}) } + // And a maintenance window open as the apply ended (novox/hq issue 224): a server stopped because + // its collector is running reads as working, not broken. + for _, w := range outcome.Windows { + report.Windows = append(report.Windows, link.Window{Step: w.Step, Holds: w.Holds, Since: w.Opened, Until: w.Until}) + } // And what runs here that nobody asked for (novox/hq ADR 0163). if strays, err := apply.Strays(ctx, apply.ExecRunner, updated); err != nil { fmt.Fprintf(os.Stderr, "mesh-host: applied, and could not list what else runs here: %v\n", err) diff --git a/internal/apply/apply.go b/internal/apply/apply.go index ca0c3da..77a4dfd 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -41,7 +41,8 @@ type Outcome struct { ID string `json:"id"` Type string `json:"type"` Target string `json:"target"` - // Action is created · updated · unchanged · corrected · removed. + // Action is created · updated · unchanged · corrected · removed — and held-still, a container + // a maintenance window is holding stopped, left for the next apply after it (novox/hq issue 224). // // "corrected" is its own answer and not a kind of "updated": it means the machine had drifted // from what this host last wrote, so somebody changed it by hand. The mesh converging is @@ -86,6 +87,10 @@ type Report struct { // Tunnel is what this apply says about the tunnel the private network took over, when the // declaration names one (novox/hq ADR 0105). Tunnel *TakenTunnel `json:"tunnel,omitempty"` + // Windows is the maintenance windows open on this machine as the apply ended (novox/hq issue + // 224): a machine whose server is stopped at 03:31 because its collector is running is + // working, not broken, and the report is where that difference is said. + Windows []Window `json:"windows,omitempty"` } // Changed reports whether anything about the machine actually moved. An apply that changed @@ -93,8 +98,10 @@ type Report struct { func (r Report) Changed() bool { for _, o := range r.Outcomes { // Holding is keeping the machine as it was found, which is not moving it; waiting is a unit - // whose account's manager is not running, which nothing moved either (novox/hq ADR 0177). - if o.Action != "unchanged" && o.Action != "held" && o.Action != "waiting" { + // whose account's manager is not running, which nothing moved either (novox/hq ADR 0177); + // held still is a container a maintenance window stopped, which the apply did not touch + // (novox/hq issue 224). + if o.Action != "unchanged" && o.Action != "held" && o.Action != "waiting" && o.Action != heldStill { return true } } @@ -174,10 +181,36 @@ func ApplyKeeping( unseal Unseal, keep Keep, ) (Report, store.State, error) { + return ApplyMindingWindows(ctx, sys, d, known, origin, run, log, unseal, keep, nil) +} + +// ApplyMindingWindows is ApplyKeeping on a machine where a scheduled step may be holding containers +// still (novox/hq issue 224, ADR 0189): a container an open window holds is left exactly as it is, +// reported held-still, and converged by the first apply after the window closes. Nil windows is a +// caller with no state directory to read them from — a test — and minds none. +// +// Read per container, at the moment the apply decides about it, not once at the start: a window +// that opens while this apply is half way through must still be seen by the containers it reaches +// afterwards. +func ApplyMindingWindows( + ctx context.Context, + sys system.System, + d *declaration.Declaration, + known store.State, + origin string, + run Runner, + log func(string), + unseal Unseal, + keep Keep, + windows *Windows, +) (report Report, _ store.State, _ error) { + // Said whichever way the apply ended, failure included: a window being open is a fact about the + // machine, not about this apply. + defer func() { report.Windows = windows.Open(time.Now()) }() if log == nil { log = func(string) {} } - report := Report{} + report = Report{} declared := map[string]bool{} for _, r := range d.Resources { @@ -377,7 +410,7 @@ func ApplyKeeping( // already the new one (novox/hq 04-ISSUES/103). That needs the file applied before the // container, which is the declared order; a container declared ahead of its file sees the // change one apply late, and never misses it. - in := inputs{declares: map[string]string{}, known: &known} + in := inputs{declares: map[string]string{}, known: &known, windows: windows} for _, resource := range d.Resources { in.declares[resource.Identity()] = declaredDigest(resource) } @@ -665,7 +698,10 @@ func ApplyKeeping( } } report.Outcomes = append(report.Outcomes, outcome) - if outcome.Action != "unchanged" { + if outcome.Action == heldStill { + // Not changed: nothing about it moved, so nothing that restarts on it restarts. + log(fmt.Sprintf(" %s %s (%s): %s", outcome.Action, outcome.ID, outcome.Target, outcome.Detail)) + } else if outcome.Action != "unchanged" { changed[resource.Identity()] = true // With the detail, when there is one: "updated app" says a container was replaced; // which file made that happen is what somebody reading the log at the time needs @@ -1668,6 +1704,9 @@ type inputs struct { // earlier in the same pass is already its new self. Nil where nothing was written: a test, // or a scheduled fire, which reads no file at creation. known *store.State + // windows is where a maintenance window says which containers it holds still (novox/hq issue + // 224). Nil minds none. + windows *Windows } // fileDigest is what a file the container reads holds, by digest. @@ -1911,6 +1950,11 @@ func applyNetwork(ctx context.Context, r *declaration.Network, run Runner) (Outc return out, nil } +// heldStill is the action of a container left alone because a maintenance window holds it stopped +// (novox/hq issue 224). Its own word rather than "unchanged": a person reading a machine that looks +// half-stopped needs to see why, and "unchanged" beside a stopped server says the opposite. +const heldStill = "held-still" + func applyContainer(ctx context.Context, r *declaration.Container, run Runner, changed map[string]bool, in inputs, previous store.Applied) (Outcome, error) { out := begin(r) @@ -1992,6 +2036,31 @@ func applyContainer(ctx context.Context, r *declaration.Container, run Runner, legacy := len(reads) > 0 && before.Spec == containerSpecReading(r, in.declares, nil) && (wasReading == nil || sameReads(wasReading, reads)) + // **A container a maintenance window holds still is left as it is** (novox/hq issue 224, ADR + // 0189). Everywhere else stopped means broken, and the rule below replaces it; during a window + // stopped is what was asked for, and replacing it starts the server under the step that needed + // it still — for the store's collector, a registry taking an upload the sweep then deletes. + // + // Left whatever it is: stopped, running because a stop failed, or with a declaration that has + // moved since. A changed declaration is not a reason to reopen a window; it is a reason for the + // next apply, which finds the window closed and the container running and converges it by the + // ordinary rule. The record keeps what the container was created reading, not what this pass + // read — it was not recreated, so a file changed under it must still be seen as changed then. + // + // Asked after the container was inspected, never before: the window is recorded before its first + // stop, so a container this apply found stopped by a window is one whose window it can see. + if win, held := in.windows.holding(r.Name, time.Now()); held { + out.reads = wasReading + out.Action = heldStill + out.Detail = fmt.Sprintf("held still by the maintenance window of %s since %s; left as it is, "+ + "and converged by the first apply after the window closes (novox/hq issue 224)", + win.Step, win.Opened.Format(time.RFC3339)) + if existed && before.Spec != want && !legacy || len(reasons) > 0 || len(changedFiles) > 0 { + out.Detail += "; its declaration moved, so that apply recreates it" + } + return out, nil + } + switch { case existed && (before.Spec == want || legacy) && before.Running && len(reasons) == 0: out.Action = "unchanged" diff --git a/internal/apply/maintenance_window.go b/internal/apply/maintenance_window.go new file mode 100644 index 0000000..8275e7a --- /dev/null +++ b/internal/apply/maintenance_window.go @@ -0,0 +1,195 @@ +package apply + +// What a maintenance window is holding still, written where every apply on this machine can read it +// (novox/hq issue 224, ADR 0189). +// +// A scheduled step declared `while-stopped` stops its module's containers, runs, and starts them +// again. Until this, nothing told the **apply** that a window was open, and the apply's rule for a +// container it finds stopped is the right rule everywhere else: stopped is broken, so remove it and +// create it running. In the middle of a window that rule reopens the window — for the store's +// collector, a registry recreated mid-collection accepts an upload the sweep then deletes, and the +// build that made the image reported success. +// +// **The apply learns which containers a window holds; the window does not take the apply lock.** +// The issue weighed both. Holding the lock for a window makes the race impossible and makes a push +// that arrives mid-window wait for minutes, which is what issue 185's outages looked like from +// outside. Telling the apply blocks nothing: a held container is reported as held and left exactly +// as it is — even when its declaration changed — and the first apply after the window converges it +// by the ordinary rule. The cost, stated in the issue, is a second source for "is this container +// meant to be running"; it is kept honest by being short-lived and by expiring on its own. +// +// **A file under the node's state directory, not a variable in the scheduler.** The scheduler lives +// in the daemon, but the daemon is not the only thing that applies here: `mesh-host reconcile` and +// `mesh-host apply` run as their own processes beside it, and the installer applies the carried +// bundle — which is where the store's registry comes from — from a third. A map shared in memory +// would protect only the daemon's own applies, which is to say it would leave open exactly the case +// a person running `reconcile` by hand at 03:31 creates. The state directory is already where every +// one of them meets (the apply lock, the node's state, the kept originals), it is root's alone, and +// each window is one small file written atomically, so no reader ever sees half of one. +// +// **A window that is never closed must not hold for ever** — the same risk the field itself carries, +// one level up. Three things close it: +// +// - the scheduler removes the file in a defer that runs after the containers are started again, +// whatever the step did; +// - a window whose process is gone is closed: the file names the process that opened it, and a +// host that crashed or was replaced mid-window is no longer holding anything, so the apply is +// free to bring back what it left stopped — which is then the apply doing the job the dead +// process could not; +// - and every window carries an end, windowAtMost after it opened, past which it is read as +// closed whoever is still alive. A step that genuinely runs longer than that is pathological, and +// the apply then does what it did before this record existed, which is the lesser harm next to a +// service the mesh can never bring back. + +import ( + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "syscall" + "time" +) + +// WindowsName is the directory beside the node's state that holds the open windows, one file each. +const WindowsName = "windows" + +// windowAtMost is the longest a window is believed. The store's collector is minutes over a store of +// fifty-odd repositories; six hours is far past any collection this mesh will do and short enough +// that a window nothing closed is over by the next working day (novox/hq issue 224). +const windowAtMost = 6 * time.Hour + +// Window is one scheduled step holding containers still: which step, which containers by runtime +// name, since when, and the latest moment it is believed. +type Window struct { + Step string `json:"step"` + Holds []string `json:"holds"` + Opened time.Time `json:"opened"` + Until time.Time `json:"until"` + // PID is the process that opened it. Not reported: it is how a window whose host died is read + // as closed. + PID int `json:"pid"` +} + +// Windows is where this machine's open windows are recorded. A nil *Windows records nothing and +// holds nothing — a test, or a caller that has no state directory. +type Windows struct { + dir string + // alive says whether a process is still running. Replaced in tests; the real one asks the kernel. + alive func(pid int) bool +} + +// WindowsIn records windows under stateDir/windows — the directory the node's state lives in, which +// every process that applies on this machine already shares. +func WindowsIn(stateDir string) *Windows { + return &Windows{dir: filepath.Join(stateDir, WindowsName), alive: processAlive} +} + +// processAlive is whether pid names a running process. Signal 0 delivers nothing and only asks; a +// process this one may not signal is still a process (EPERM), which can only happen across users. +func processAlive(pid int) bool { + if pid <= 0 { + return false + } + err := syscall.Kill(pid, 0) + return err == nil || errors.Is(err, syscall.EPERM) +} + +// fileFor is where one step's window is written. A step id is the declaration's, which may carry a +// dot or a slash; the file name keeps what is safe and replaces the rest. +func (w *Windows) fileFor(step string) string { + safe := strings.Map(func(r rune) rune { + switch { + case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == '-', r == '_', r == '.': + return r + } + return '_' + }, step) + return filepath.Join(w.dir, "window-"+safe+".json") +} + +// open records that step holds these containers still from now, before the first of them is +// stopped — so no apply can find one stopped and not know why. +func (w *Windows) open(step string, holds []string, now time.Time) error { + if w == nil { + return nil + } + if err := os.MkdirAll(w.dir, 0o700); err != nil { + return err + } + raw, err := json.Marshal(Window{Step: step, Holds: holds, Opened: now.UTC(), + Until: now.UTC().Add(windowAtMost), PID: os.Getpid()}) + if err != nil { + return err + } + return writeAtomically(w.fileFor(step), raw, 0o600) +} + +// close removes the record of step's window. Called after its containers are started again; a file +// already gone is the state wanted. +func (w *Windows) close(step string) error { + if w == nil { + return nil + } + if err := os.Remove(w.fileFor(step)); err != nil && !errors.Is(err, os.ErrNotExist) { + return err + } + return nil +} + +// Open is the windows open at now, by step. One that has ended, or whose process is gone, is not +// open — and its file is removed on the way, since nothing else will (novox/hq issue 224). A file +// that cannot be read is not a window either: the alternative is a machine whose containers are +// never converged because of a byte nobody can see. +func (w *Windows) Open(now time.Time) []Window { + if w == nil { + return nil + } + entries, err := os.ReadDir(w.dir) + if err != nil { + return nil + } + var open []Window + for _, e := range entries { + if e.IsDir() || !strings.HasPrefix(e.Name(), "window-") || !strings.HasSuffix(e.Name(), ".json") { + continue + } + path := filepath.Join(w.dir, e.Name()) + raw, err := os.ReadFile(path) + if err != nil { + continue + } + var win Window + if err := json.Unmarshal(raw, &win); err != nil { + _ = os.Remove(path) + continue + } + if !now.Before(win.Until) || !w.alive(win.PID) { + _ = os.Remove(path) + continue + } + open = append(open, win) + } + sort.Slice(open, func(i, j int) bool { return open[i].Step < open[j].Step }) + return open +} + +// holding is the open window that holds the container named name, if one does. +func (w *Windows) holding(name string, now time.Time) (Window, bool) { + for _, win := range w.Open(now) { + for _, held := range win.Holds { + if held == name { + return win, true + } + } + } + return Window{}, false +} + +// String is the window as a person reads it in a log line. +func (win Window) String() string { + return fmt.Sprintf("%s holds %s still since %s", win.Step, strings.Join(win.Holds, ", "), + win.Opened.Format(time.RFC3339)) +} diff --git a/internal/apply/maintenance_window_test.go b/internal/apply/maintenance_window_test.go index 0acc358..239f967 100644 --- a/internal/apply/maintenance_window_test.go +++ b/internal/apply/maintenance_window_test.go @@ -3,12 +3,16 @@ 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, @@ -236,3 +240,303 @@ func TestAChangedWindowMovesTheSpec(t *testing.T) { 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) + } +} diff --git a/internal/apply/schedule.go b/internal/apply/schedule.go index 179d3c4..32ba578 100644 --- a/internal/apply/schedule.go +++ b/internal/apply/schedule.go @@ -51,6 +51,11 @@ type Scheduler struct { run Runner log func(string) + // windows is where a step that holds containers still says so, for every apply on this + // machine to read (novox/hq issue 224). Nil records nothing: a test, or a scheduler with no + // state directory. + windows *Windows + mu sync.Mutex cri string // the container runtime, detected once and cached jobs map[string]*scheduledJob @@ -101,6 +106,14 @@ func NewScheduler(clock Clock, run Runner, log func(string)) *Scheduler { return &Scheduler{clock: clock, run: run, log: log, jobs: map[string]*scheduledJob{}} } +// RecordWindowsIn makes every window this scheduler opens readable by the applies on this machine, +// in whichever process they run (novox/hq issue 224). Called once, before Run. +func (s *Scheduler) RecordWindowsIn(w *Windows) { + s.mu.Lock() + defer s.mu.Unlock() + s.windows = w +} + // Sync re-establishes the scheduled steps from a declaration: it adds ones newly declared, re-arms // any whose image, environment or cadence changed, and forgets those the declaration no longer // names. Rebuilt from the declaration each apply because the declaration is the source of truth @@ -234,7 +247,31 @@ func (s *Scheduler) fire(ctx context.Context, j *scheduledJob) { // 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. + // + // **And it is written down before the first stop, and erased after the last start** (novox/hq + // issue 224). An apply that finds a container stopped replaces it; one that finds it stopped + // *and held* leaves it. So the record must exist for every moment the container could be found + // stopped by the window: opened before holdStill, and — defers running last-registered first — + // closed after letRun. A record that cannot be written means no apply would know, and an apply + // arriving then would recreate the server under the step; the step is not run, and said, which + // costs one night's collection rather than a blob. if len(j.hold) > 0 { + s.mu.Lock() + windows := s.windows + s.mu.Unlock() + // The wall clock, not the scheduler's: the cadence is the scheduler's to decide, but how + // long a window is believed is read by applies in other processes, against theirs. + if err := windows.open(j.id, j.hold, time.Now()); err != nil { + s.log(fmt.Sprintf("scheduled step %s: cannot record the window it holds %v still for, "+ + "so an apply could undo it — not run: %v", j.id, j.hold, err)) + return + } + defer func() { + if err := windows.close(j.id); err != nil { + s.log(fmt.Sprintf("scheduled step %s: the window is closed and its record could not be "+ + "removed; it expires on its own: %v", j.id, err)) + } + }() defer s.letRun(ctx, cri, j) s.holdStill(ctx, cri, j) } diff --git a/internal/bootstrap/apply.go b/internal/bootstrap/apply.go index 54da41a..b88e559 100644 --- a/internal/bootstrap/apply.go +++ b/internal/bootstrap/apply.go @@ -66,9 +66,12 @@ func ApplyBundle(ctx context.Context, o Options, sys system.System, d *declarati // distribution's own /etc/nftables.conf among them — so it keeps the original of each file it // has no record of, beside the node's state, exactly as a declaration from the mesh does // (novox/hq ADR 0100). - report, updated, applyErr := apply.ApplyKeeping(ctx, sys, d, known, store.OriginCarried, run, + // Minding the maintenance windows a running host has open: the carried bundle is where the + // store's registry comes from, and an installer re-run at 03:31 must not restart it under its + // collector (novox/hq issue 224). + report, updated, applyErr := apply.ApplyMindingWindows(ctx, sys, d, known, store.OriginCarried, run, func(line string) { say(" " + strings.TrimPrefix(line, " ")) }, refuseSealed, - apply.KeepIn(filepath.Dir(o.State))) + apply.KeepIn(filepath.Dir(o.State)), apply.WindowsIn(filepath.Dir(o.State))) // The bundle is consumed, whichever way the apply went: what is on the machine came from these // bytes, and the carried ones must not be applied over it. The mode is the operator's word at diff --git a/internal/link/messages.go b/internal/link/messages.go index ed559c6..b593674 100644 --- a/internal/link/messages.go +++ b/internal/link/messages.go @@ -121,6 +121,12 @@ type Report struct { // (ADR 0168). Nil on a machine found with none, and on an adopted one, where Firewall says it. FoundFirewall *FoundFirewall `json:"found_firewall,omitempty"` + // Windows is the maintenance windows open on this machine when it reported (novox/hq issue 224, + // ADR 0189): a scheduled step holding its module's containers still. A machine whose store is + // stopped at 03:31 because its collector is running is working, and without this it reads as + // broken. + Windows []Window `json:"windows,omitempty"` + // Strays is what runs on the machine that the mesh neither wrote nor holds (novox/hq ADR // 0163): containers nobody declared and nobody holds, the ones a cutover leaves behind. Strays []Stray `json:"strays,omitempty"` @@ -225,6 +231,15 @@ type Held struct { Facts map[string]any `json:"facts,omitempty"` } +// A Window is one scheduled step holding containers still (novox/hq issue 224): the step's id, the +// containers by runtime name, since when, and the latest moment the host believes it. +type Window struct { + Step string `json:"step"` + Holds []string `json:"holds"` + Since time.Time `json:"since"` + Until time.Time `json:"until"` +} + // A Stray is a container the mesh neither wrote nor holds (ADR 0163). type Stray struct { Kind string `json:"kind"`