Merge pull request 'An apply leaves a container a maintenance window holds still (hq issue 224)' (#23) from fix/an-apply-leaves-a-held-container-held into main
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
This commit was merged in pull request #23.
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