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.
This commit is contained in:
+17
-5
@@ -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)
|
||||
|
||||
+75
-6
@@ -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"
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"`
|
||||
|
||||
Reference in New Issue
Block a user