Files
mesh-host/internal/apply/apply.go
T
jschoubben 27c4b765b2 Refuse a declaration for the other mode, or older than the mesh's last, and say what an apply would change first
An operator ran `mesh-host reconcile` on an adopted control-node with twelve
modules assigned. It applied the bundle the host carries — the genesis
declaration, foundation only, converged: recreated the store, failed on the
broker's held port, wrote the converged base filter and started its service,
and stopped at the first failing action. The filter closed the machine for
forty-five minutes. The host reported the node adopted in every report, the
declaration said converged, and nothing compared the two; nothing was printed
before acting (hq issue 104).

The host now records the node's mode — from every declaration the mesh sends,
and at genesis from what the operator said — and refuses, at the point of
application, a declaration that says the other mode, naming both and the act
that changes it. Only a declaration the link delivers, signed, changes the
mode: that is how `converge` and `adopt` arrive, so the flip still works and
nothing else can do it. Genesis marks the bundle consumed, with the digest of
what it applied, so `reconcile` holds a node the mesh has spoken to against
what the mesh last said and never the bundle, and refuses the carried bytes
when they are not what genesis applied. A file is refused when it is not what
the mesh last said: a declaration carries no sequence and no issued-at, so the
host cannot tell older from newer, and says so. Both commands print what they
would change — a hold, a removal, an action named as one — before touching
anything, and --dry-run is that list and nothing more.
2026-09-23 23:15:28 +02:00

1670 lines
67 KiB
Go

// Package apply makes a machine match a declaration.
//
// Three properties, each following a recorded decision, and each of them the difference
// between this and a script that writes files:
//
// - A failed step fails the apply (novox/hq ADR 0010). Not "logs and continues": a partial
// apply that reports success is the mesh's most expensive shape.
// - Every applier READS BACK. Setting a value is not evidence the value took.
// - What was applied is recorded after it works, never before (ADR 0018). A failed apply
// leaves the machine in whatever state it reached, and nothing must claim otherwise.
package apply
import (
"context"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"errors"
"fmt"
"os"
"os/exec"
"path/filepath"
"sort"
"strconv"
"strings"
"time"
"github.com/novox/mesh-host/internal/declaration"
"github.com/novox/mesh-host/internal/store"
"github.com/novox/mesh-host/internal/system"
)
// Runner executes a command. The real one is used everywhere outside unit tests; behaviour
// against a real system is tested alongside rather than mocked (novox/hq ADR 0017).
type Runner = system.Runner
// Outcome is what happened to one resource.
type Outcome struct {
ID string `json:"id"`
Type string `json:"type"`
Target string `json:"target"`
// Action is created · updated · unchanged · corrected · removed.
//
// "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
// right either way; being unable to say which happened is not.
Action string `json:"action"`
Detail string `json:"detail,omitempty"`
// wrote is a digest of what this apply put there, kept so the next one can tell a machine
// that drifted from one the mesh changed its mind about. Not reported: it is bookkeeping.
wrote string
// into is what a file written into held before the mesh's keys (novox/hq ADR 0102).
into *store.Into
}
// Report is what an apply did, in the order it did it.
type Report struct {
Outcomes []Outcome `json:"outcomes"`
}
// Changed reports whether anything about the machine actually moved. An apply that changed
// nothing is the ordinary steady state, and saying so is not the same as saying it failed.
func (r Report) Changed() bool {
for _, o := range r.Outcomes {
// Holding is keeping the machine as it was found, which is not moving it.
if o.Action != "unchanged" && o.Action != "held" {
return true
}
}
return false
}
// Error is a failure part-way through, carrying what had already been done.
//
// The outcomes matter as much as the message: the machine is in whatever state the apply
// reached, and the only honest thing to hand back is the list of what did happen.
type Error struct {
Resource string
Err error
Done Report
// Others is how many more resources also failed. Named rather than folded into the message,
// because "one thing failed" and "eleven things failed" are different machines and the first
// line is what somebody reads.
Others int
// Gated is true when what failed was an action, so nothing after it was attempted. A person
// reading a report needs to know the difference between "these things failed" and "these
// things failed and the rest was never tried".
Gated bool
}
func (e *Error) Error() string {
also := ""
if e.Others == 1 {
also = ", and one other resource also failed"
} else if e.Others > 1 {
also = fmt.Sprintf(", and %d other resources also failed", e.Others)
}
rest := "everything was attempted, so what is not listed as failed was done."
if e.Gated {
// An action is a gate: it exists to make something true before the next thing needs it.
rest = "this is an action, so nothing after it was attempted — the machine is in " +
"whatever state that left it."
}
return fmt.Sprintf("applying %q: %v%s\n\n%d resource(s) were applied and remain; %s",
e.Resource, e.Err, also, len(e.Done.Outcomes), rest)
}
func (e *Error) Unwrap() error { return e.Err }
// Apply makes the machine match the declaration, and returns what it did.
//
// Removal happens FIRST, and the order is not arbitrary. A resource that leaves a declaration
// while another arrives at the same path is an ordinary rename: removing afterwards would
// delete the file that had just been written. Removing first risks losing the old state if the
// apply then fails — a recovery concern, where the other is a correctness one. The one exception
// is what protects an adopted node, the openings and the guard: that goes last, and only when
// everything else applied (novox/hq ADR 0103).
func Apply(
ctx context.Context,
sys system.System,
d *declaration.Declaration,
known store.State,
origin string,
run Runner,
log func(string),
unseal Unseal,
) (Report, store.State, error) {
return ApplyKeeping(ctx, sys, d, known, origin, run, log, unseal, nil)
}
// ApplyKeeping is Apply on a node that may be adopted: keep is where the original of a file found
// there is recorded before anything else happens to it (novox/hq ADR 0100). Nil is a caller that
// can never be handed an adopted declaration — the carried bundle, which may not say it.
func ApplyKeeping(
ctx context.Context,
sys system.System,
d *declaration.Declaration,
known store.State,
origin string,
run Runner,
log func(string),
unseal Unseal,
keep Keep,
) (Report, store.State, error) {
if log == nil {
log = func(string) {}
}
report := Report{}
declared := map[string]bool{}
for _, r := range d.Resources {
declared[r.Identity()] = true
}
// Which firewall is found here, before anything else, since an unsupported one refuses the
// whole declaration (novox/hq ADR 0100). Nothing for a converged node.
fw, err := foundFirewall(ctx, d, &known, run, log)
if err != nil {
return report, known, &Error{Resource: "the firewall found on this machine", Err: err, Done: report}
}
// What an adopted node's untaken modules find on the machine, looked at before anything in
// this apply — a removal included — could change it or its records (novox/hq ADR 0103).
before := lookBefore(ctx, sys, d, known, run)
removeOrphan := func(orphan store.Applied) error {
var action, detail string
var err error
if declaration.Type(orphan.Type) == declaration.TypeOpening {
action, detail, err = removeOpening(ctx, orphan, run, known.Firewall)
} else {
action, detail, err = remove(ctx, sys, orphan, run)
}
if err != nil {
return &Error{Resource: orphan.ID, Err: err, Done: report}
}
known.Forget(orphan.ID)
report.Outcomes = append(report.Outcomes, Outcome{
ID: orphan.ID, Type: orphan.Type, Target: orphan.Target,
Action: action, Detail: detail,
})
log(fmt.Sprintf(" %s %s (%s)", action, orphan.ID, orphan.Target))
return nil
}
// **What protects an adopted node goes last on the flip, and first on the way back** (novox/hq
// ADR 0103). The openings and the guard are what keep the mesh reachable through the found
// firewall and the store unreachable from outside.
//
// When a node is converged they leave the declaration, and removing them first would leave the
// store open from the moment the guard stops until the derived filter loads — and for ever, if
// the filter then fails. So on a converged declaration they are removed only once everything
// else applied and the found firewall is retired; if anything failed, they stay, recorded, for
// the next try.
//
// Returned to adopted, it is the mirror image: removing the derived filter first would leave
// the store open until the guard loads. So the guard's own resources are applied before any
// orphan is removed, and if a removal then fails the guard is already up. A stale opening on an
// adopted node is removed as any orphan is.
var protecting, orphans []store.Applied
for _, orphan := range known.Orphans(declared, origin) {
if d.Adoption == nil && strings.HasPrefix(orphan.ID, declaration.AdoptionPrefix) {
protecting = append(protecting, orphan)
continue
}
orphans = append(orphans, orphan)
}
ordered := d.Resources
guardFirst := 0
if d.Adoption != nil {
ordered = nil
for _, r := range d.Resources {
if strings.HasPrefix(r.Identity(), guardPrefix) {
ordered = append(ordered, r)
}
}
guardFirst = len(ordered)
for _, r := range d.Resources {
if !strings.HasPrefix(r.Identity(), guardPrefix) {
ordered = append(ordered, r)
}
}
}
orphansRemoved := false
removeOrphans := func() error {
orphansRemoved = true
for _, orphan := range orphans {
if err := removeOrphan(orphan); err != nil {
return err
}
}
return nil
}
// **A hold whose resource is no longer declared is let go, and nothing on disk is touched.**
// What was found stays as it was; only the host's note that it holds it for a module goes, so
// the node stops reporting a hold for a module no longer assigned. Should the module come back,
// what is there is present with no record and is found, and held, again — its first kept
// original is never overwritten (novox/hq ADR 0100). Only a declaration from the mesh says
// what is assigned: a carried bundle's silence is not an unassignment.
if origin == store.OriginDeclared {
for _, h := range append([]store.Held{}, known.Held...) {
if declared[h.ID] {
continue
}
known.Release(h.ID)
report.Outcomes = append(report.Outcomes, Outcome{ID: h.ID, Type: h.Kind, Target: h.Target,
Action: "forgotten", Detail: "no longer declared; left as found"})
log(fmt.Sprintf(" forgotten %s (%s): no longer declared; left as found", h.ID, h.Target))
}
}
// What moved in this apply, so a service that must reflect a file can be told the file
// moved. Only within one apply: a change from an earlier one has already been reflected, and
// restarting for it every time would make a steady machine restart its services for ever.
changed := map[string]bool{}
// **What each resource currently says, so a container can be identified by its inputs.**
//
// `changed` above is edge-triggered and only within one apply, which is right for "restart it
// because this just moved" and wrong for "is this container running the file that is there
// now". A file written in an earlier apply, or written before the container declared it as a
// dependency, leaves a container holding values nothing will ever re-read: it is up, the
// machine reports success, and what is inside is using a credential the mesh has replaced
// (novox/hq 04-ISSUES/045). Folding these into the container's spec makes the comparison a
// standing one instead.
declares := map[string]string{}
for _, resource := range d.Resources {
declares[resource.Identity()] = declaredDigest(resource)
}
// Everything is attempted, and every failure is reported.
//
// **It used to stop at the first one**, and that made one broken resource hold the whole
// machine hostage: a module declaring a package that does not exist meant every module
// ordered after it was never applied, for ever, and the mesh reported "failed" without
// saying that the rest had not been tried. A machine with one bad module and nine good ones
// ran none of the nine.
//
// The argument for stopping was that a resource may depend on an earlier one. It still may —
// and it will then fail its own check and be reported, which is more information than not
// attempting it. A service started against a file that was never written does not verify, and
// this host reads back after every write precisely so that is caught rather than assumed.
//
// What does not change: **a declaration that cannot be parsed is still refused whole.** That
// is a different thing — one is "this machine could not do it", the other is "this was never
// a declaration", and they are fixed in different places.
var failures []*Error
for i, resource := range ordered {
if !orphansRemoved && i == guardFirst {
// **Only a guard that is up may let the filter go.** Removing the derived filter's
// resources stops its unit, whose stop deletes the mesh's table; if a guard resource
// failed, doing that would leave the node with neither, and the store open until some
// later reconcile gets the guard up (novox/hq ADR 0103).
if len(failures) > 0 {
first := failures[0]
first.Done = report
first.Others = len(failures) - 1
log(" kept " + guardPrefix + "*: the guard is not up, so what it replaces was left in force")
return report, known, first
}
if err := removeOrphans(); err != nil {
return report, known, err
}
}
// **On an adopted node, what is found is kept until its module is taken** (novox/hq ADR
// 0100, ADR 0103). Before anything is applied: whatever of a module not yet taken is
// present with no record of this host making it — or would reach what is — is held as it
// is and reported. Once held it stays held until its module is taken, and it is never
// recorded as applied, so it is never removed as an orphan either.
if d.Adoption != nil {
isHeld, news, outcome, err := holdOnAdopted(ctx, sys, resource, d, &known, before, run, keep,
changed, time.Now().UTC())
if err != nil {
failures = append(failures, &Error{Resource: resource.Identity(), Err: err, Done: report})
log(fmt.Sprintf(" failed %s (%s): %v", resource.Identity(), resource.Target(), err))
continue
}
if isHeld {
report.Outcomes = append(report.Outcomes, outcome)
if news {
log(fmt.Sprintf(" held %s (%s): %s", outcome.ID, outcome.Target, outcome.Detail))
}
continue
}
}
was, _ := known.Find(resource.Identity())
var outcome Outcome
var err error
if o, isOpening := resource.(*declaration.Opening); isOpening {
outcome, err = applyOpening(ctx, o, run, fw)
} else {
// A file this host has no record of, under any id, is the machine's until the mesh
// writes over it — on any node, adopted or not: its original is kept first.
var keepFound Keep
if f, isFile := resource.(*declaration.File); isFile && was.ID == "" &&
!known.Recorded(string(declaration.TypeFile), f.Path) {
keepFound = keep
}
outcome, err = applyOne(ctx, sys, resource, run, changed, declares, was, unseal, keepFound)
}
if err != nil {
failed := &Error{Resource: resource.Identity(), Err: err, Done: report}
failures = append(failures, failed)
log(fmt.Sprintf(" failed %s (%s): %v", resource.Identity(), outcome.Target, err))
// **A failed action stops what follows. Nothing else does.**
//
// An action is the only shape whose purpose is to make something true *before* the
// next thing needs it — which is why it is the only one with a `verify`. The
// bootstrap is a row of them: the store answers, then its databases exist, then their
// schemas, then the broker. Carrying on past one that did not happen means starting
// things against a machine that is not ready, and on a small machine that is how a
// database still initialising gets its memory taken away and shuts down. Observed,
// in the lab, caused by an earlier version of this loop.
//
// Everything else is independent state. A package that will not install has nothing
// to do with a file on the other side of the declaration, and stopping there is what
// made one broken module hold a whole machine hostage
// (novox/hq 04-ISSUES/011).
// A run-once container is a step with the same purpose as an action: to make
// something true *before* the next thing needs it (novox/hq ADR 0052). A broker whose
// dynsec store was not seeded must not be started, so a run-once step that did not
// complete gates what follows exactly as a failed action does — that is the whole of
// how "before that container starts" is enforced, since the step is declared before it.
gates := resource.Kind() == declaration.TypeAction
if c, ok := resource.(*declaration.Container); ok && c.RunOnce {
gates = true
}
if gates {
failed.Done = report
failed.Others = len(failures) - 1
failed.Gated = true
return report, known, failed
}
continue
}
// Only now. The record follows the fact, never leads it.
known.Record(store.Applied{
Origin: origin,
ID: resource.Identity(), Type: string(resource.Kind()),
Target: outcome.Target, AppliedAt: time.Now().UTC(),
Wrote: outcome.wrote,
Into: outcome.into,
Holds: holds(resource),
})
// Its module has been taken, and what was held for it is now the mesh's.
if held, wasHeld := known.HeldAt(resource.Identity()); wasHeld {
known.Release(held.ID)
outcome.Detail = takenDetail(held)
}
report.Outcomes = append(report.Outcomes, outcome)
if outcome.Action != "unchanged" {
changed[resource.Identity()] = true
log(fmt.Sprintf(" %s %s (%s)", outcome.Action, outcome.ID, outcome.Target))
}
}
if !orphansRemoved {
if err := removeOrphans(); err != nil {
return report, known, err
}
}
// A converged node whose found firewall was in force retires it only now, once everything —
// the mesh's derived filter among it — applied cleanly (novox/hq ADR 0100).
if len(failures) == 0 {
if err := retireFirewall(ctx, d, origin, &known, run, log); err != nil {
return report, known, &Error{Resource: "the firewall found on this machine", Err: err, Done: report}
}
for _, orphan := range protecting {
if err := removeOrphan(orphan); err != nil {
return report, known, err
}
}
}
if len(failures) > 0 {
// The first, carrying everything that did happen. One error is what the caller reports
// and what a person reads first; the rest are in the report, which is what the mesh
// keeps.
first := failures[0]
first.Done = report
first.Others = len(failures) - 1
return report, known, first
}
return report, known, nil
}
// guardPrefix is the ids of the mesh's guard on an adopted node: its package, table, unit and
// service (novox/hq ADR 0100).
const guardPrefix = declaration.AdoptionPrefix + "guard"
// Unseal opens a value the mesh sealed to this node. Nil when the node has no sealing key, which
// makes every sealed file an error rather than a silently skipped one.
type Unseal func(sealed string) ([]byte, error)
func applyOne(ctx context.Context, sys system.System, r declaration.Resource, run Runner,
changed map[string]bool, declares map[string]string, previous store.Applied,
unseal Unseal, keepFound Keep) (Outcome, error) {
switch res := r.(type) {
case *declaration.Directory:
return applyDirectory(res)
case *declaration.File:
return applyFile(res, previous, unseal, keepFound)
case *declaration.Service:
return applyService(ctx, sys, res, run, changed)
case *declaration.Package:
return applyPackage(ctx, sys, res, run)
case *declaration.Container:
return applyContainer(ctx, res, run, changed, declares, previous)
case *declaration.User:
return applyUser(ctx, sys, res, run)
case *declaration.Archive:
return applyArchive(ctx, res, previous)
case *declaration.Process:
return applyProcess(ctx, res, run, changed, previous)
case *declaration.Action:
return applyAction(ctx, res, run)
case *declaration.Network:
return applyNetwork(ctx, res, run)
case *declaration.Access:
return applyAccess(res)
default:
// Unreachable: the declaration refused this already. Present because "unreachable"
// stops being true the moment someone adds a kind and forgets this switch.
return Outcome{}, fmt.Errorf("no applier for type %q", r.Kind())
}
}
// begin starts an outcome from any resource, so the three facts a report needs are read from
// the resource itself rather than restated by each applier.
func begin(r declaration.Resource) Outcome {
return Outcome{ID: r.Identity(), Type: string(r.Kind()), Target: r.Target()}
}
func modeOf(spec string, fallback os.FileMode) (os.FileMode, error) {
if spec == "" {
return fallback, nil
}
parsed, err := strconv.ParseUint(spec, 8, 32)
if err != nil {
return 0, fmt.Errorf("mode %q: %w", spec, err)
}
return os.FileMode(parsed), nil
}
func applyDirectory(r *declaration.Directory) (Outcome, error) {
out := begin(r)
mode, err := modeOf(r.Mode, 0o755)
if err != nil {
return out, err
}
before, err := os.Stat(r.Path)
existed := err == nil
if err != nil && !errors.Is(err, os.ErrNotExist) {
return out, err
}
if existed && !before.IsDir() {
return out, fmt.Errorf("%s exists and is not a directory", r.Path)
}
if !existed {
if err := os.MkdirAll(r.Path, mode); err != nil {
return out, err
}
}
// Set explicitly even when it existed: MkdirAll applies the mode only on creation, and a
// permission set at creation is not a permission maintained — a lesson this repository
// already paid for once, with world-readable environment files.
if err := os.Chmod(r.Path, mode); err != nil {
return out, err
}
// Read back.
after, err := os.Stat(r.Path)
if err != nil {
return out, fmt.Errorf("made %s and cannot stat it: %w", r.Path, err)
}
if !after.IsDir() {
return out, fmt.Errorf("%s is not a directory after applying", r.Path)
}
if after.Mode().Perm() != mode.Perm() {
return out, fmt.Errorf("%s is mode %o after setting %o", r.Path, after.Mode().Perm(), mode.Perm())
}
ownedAlready, err := ownedBy(r.Path, r.Owner)
if err != nil {
return out, err
}
if !ownedAlready {
if err := own(r.Path, r.Owner); err != nil {
return out, err
}
}
out.Action = "unchanged"
if !existed {
out.Action = "created"
} else if !ownedAlready {
out.Action = "updated"
} else if before.Mode().Perm() != mode.Perm() {
out.Action = "updated"
out.Detail = fmt.Sprintf("mode %o to %o", before.Mode().Perm(), mode.Perm())
}
return out, nil
}
// applyAccess confirms an operator-owned path is present, and owns nothing about it.
//
// **The mirror image of applyDirectory** (novox/hq ADR 0051). A directory the host makes, chmods,
// chowns and removes when empty. An access it does none of: the path is the operator's — a media
// library, a download spool that several modules share — and the host's only job is to be sure it
// is there before anything mounts it.
//
// **Absent is refused, not created.** A bind mount whose source does not exist is made for you by
// the container runtime, as root, with whatever mode it picks — which is exactly the silent
// wrong-ownership 04-ISSUES/026 records. So the host checks first and says plainly that the
// operator must provide the path, rather than conjuring a directory it does not own and cannot
// give the right owner. Nothing is written, so this never reports a change: the machine did not
// move, the host merely confirmed a fact about it.
func applyAccess(r *declaration.Access) (Outcome, error) {
out := begin(r)
info, err := os.Stat(r.Path)
if errors.Is(err, os.ErrNotExist) {
return out, fmt.Errorf(
"%s is not there, and the mesh does not own it — the operator must provide it. It is "+
"shared, pre-existing data (novox/hq ADR 0051): the host mounts it and creates "+
"nothing, so a missing one is said here rather than made as root by the container "+
"runtime", r.Path)
}
if err != nil {
return out, err
}
if !info.IsDir() {
return out, fmt.Errorf(
"%s is not a directory, and an access is a shared directory the operator provides",
r.Path)
}
mode := declaration.AccessRead
if r.Mode != "" {
mode = r.Mode
}
out.Action = "unchanged"
out.Detail = "operator-owned; present, " + mode + ", nothing managed"
return out, nil
}
// keepFound, when not nil, is where the original of a file this host has no record of is kept
// before it is written over (novox/hq ADR 0100): once, never overwritten, and named in the outcome.
func applyFile(r *declaration.File, previous store.Applied, unseal Unseal, keepFound Keep) (Outcome, error) {
if r.Into != "" {
return applyInto(r, previous)
}
out := begin(r)
// What actually goes on disk. For a sealed file the mesh never had this, and neither did
// whatever carried the declaration here.
content := r.Content
if r.Bytes != "" {
// Not text. Decoded here rather than written as base64, because what a declaration says
// is in a file has to be what ends up in it — a wallpaper stored as its own encoding is
// a wallpaper nothing can open.
decoded, err := base64.StdEncoding.DecodeString(r.Bytes)
if err != nil {
return out, fmt.Errorf("%s carries bytes that are not base64: %w", r.Path, err)
}
content = string(decoded)
}
// A secret written world-readable is a secret. The default differs from an ordinary file's
// for that reason alone; an explicit mode still wins, because a module may need its own user
// to read it and only the module knows which.
fallback := os.FileMode(0o644)
if r.Secret() {
fallback = 0o600
if unseal == nil {
// Refused rather than skipped. A machine that quietly does not apply the one resource
// carrying a credential is a machine that looks configured and cannot connect.
return out, fmt.Errorf(
"%s is sealed to this node and this node has no sealing key", r.Path)
}
opened, err := unseal(r.Sealed)
if err != nil {
return out, fmt.Errorf("cannot open %s: %w", r.Path, err)
}
content = string(opened)
}
// Sealed values into the holes the content left for them. **The host is the only thing that
// ever holds both** — the mesh discarded the value, and the module wrote the document without
// it (novox/hq ADR 0024).
if len(r.Secrets) > 0 {
if unseal == nil {
return out, fmt.Errorf(
"%s needs %d sealed value(s) and this node has no sealing key",
r.Path, len(r.Secrets))
}
for _, name := range r.SecretsUsed() {
opened, err := unseal(r.Secrets[name])
if err != nil {
return out, fmt.Errorf("cannot open the secret %q for %s: %w", name, r.Path, err)
}
content = strings.ReplaceAll(content, "${secret:"+name+"}", string(opened))
}
// A file carrying a credential is not world-readable, whatever else it also carries.
if r.Mode == "" {
fallback = 0o600
}
}
out.wrote = digestOf(content)
mode, err := modeOf(r.Mode, fallback)
if err != nil {
return out, err
}
existing, readErr := os.ReadFile(r.Path)
existed := readErr == nil
if readErr != nil && !errors.Is(readErr, os.ErrNotExist) {
return out, readErr
}
// A seed that is already there is left exactly as it is — whatever has grown in it since is
// not the mesh's to put back (novox/hq issue 035). What this host records is what it once
// wrote, so a later declaration that changes the seed is not mistaken for drift either.
if r.CreateOnce && existed {
out.wrote = previous.Wrote
if out.wrote == "" {
out.wrote = digestOf(string(existing))
}
out.Action = "kept"
out.Detail = "created once, and present; what is in it now is not the mesh's to change"
return out, nil
}
var beforeMode os.FileMode
if existed {
if info, err := os.Stat(r.Path); err == nil {
beforeMode = info.Mode().Perm()
}
}
contentSame := existed && string(existing) == content
// Whether the machine still holds what this host last put there. When it does not, and the
// declaration has not changed either, somebody edited it — and saying so is the whole
// difference between a change that vanishes mysteriously and one that is reported.
drifted := existed && previous.Wrote != "" && digestOf(string(existing)) != previous.Wrote
modeSame := existed && beforeMode == mode.Perm()
kept := ""
if !contentSame {
if existed && keepFound != nil {
// Before anything is written: a keep that fails stops the write, since the
// original could not be had back otherwise.
if kept, err = keepFound(r.Path, existing, beforeMode); err != nil {
return out, fmt.Errorf("keeping the original of %s before writing over it: %w", r.Path, err)
}
}
if err := os.MkdirAll(filepath.Dir(r.Path), 0o755); err != nil {
return out, err
}
if err := writeAtomically(r.Path, []byte(content), mode); err != nil {
return out, err
}
} else if !modeSame {
if err := os.Chmod(r.Path, mode); err != nil {
return out, err
}
}
// Read back — the file, not the call that wrote it.
written, err := os.ReadFile(r.Path)
if err != nil {
return out, fmt.Errorf("wrote %s and cannot read it back: %w", r.Path, err)
}
if string(written) != content {
return out, fmt.Errorf("%s does not contain what was declared after writing it", r.Path)
}
info, err := os.Stat(r.Path)
if err != nil {
return out, err
}
if info.Mode().Perm() != mode.Perm() {
return out, fmt.Errorf("%s is mode %o after setting %o", r.Path, info.Mode().Perm(), mode.Perm())
}
// And who it belongs to. Checked before setting, so a file already owned correctly is not
// reported as changed on every apply — which would make every reconcile look like work.
ownedAlready, err := ownedBy(r.Path, r.Owner)
if err != nil {
return out, err
}
if !ownedAlready {
if err := own(r.Path, r.Owner); err != nil {
return out, err
}
contentSame = false
}
switch {
case !existed:
out.Action = "created"
case drifted:
// Somebody changed this on the machine. The mesh puts it back either way — that is what
// holding a machine to what it was told means — but a change that vanishes with nothing
// said is how a person ends up editing the same file every five minutes, believing the
// machine is broken.
out.Action = "corrected"
out.Detail = "it had been changed on the machine since this host last wrote it"
case !contentSame && !modeSame:
out.Action = "updated"
out.Detail = "content and mode"
case !contentSame:
out.Action = "updated"
out.Detail = "content"
case !modeSame:
out.Action = "updated"
out.Detail = fmt.Sprintf("mode %o to %o", beforeMode, mode.Perm())
default:
out.Action = "unchanged"
}
if kept != "" {
if out.Detail != "" {
out.Detail += "; "
}
out.Detail += "the file found here, which the mesh had no record of, was kept at " + kept
}
return out, nil
}
// writeAtomically writes through a temporary file in the same directory.
//
// A reader of a managed file must never see half of one. The mesh's own configuration is read
// by daemons that reload on change, so a torn write is a service reading a truncated config.
func writeAtomically(path string, content []byte, mode os.FileMode) error {
tmp, err := os.CreateTemp(filepath.Dir(path), ".mesh-host-*")
if err != nil {
return err
}
defer os.Remove(tmp.Name())
if _, err := tmp.Write(content); err != nil {
tmp.Close()
return err
}
if err := tmp.Sync(); err != nil {
tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
if err := os.Chmod(tmp.Name(), mode); err != nil {
return err
}
return os.Rename(tmp.Name(), path)
}
// reflects reports whether anything this service must mirror changed in this apply.
func reflects(r *declaration.Service, changed map[string]bool) bool {
return len(reflected(r, changed)) > 0
}
// reflected is which of them changed, so the outcome can say why the service was restarted. A
// restart with no reason given is indistinguishable from a service that keeps falling over.
func reflected(r *declaration.Service, changed map[string]bool) []string {
return restartedBy(r.RestartOn, changed)
}
// serviceReloader is a service manager that can tell a running unit to read its configuration again.
type serviceReloader interface {
ReloadService(ctx context.Context, run system.Runner, unit string) error
}
// unitReloader is a service manager that caches unit files and must be told to read them again.
type unitReloader interface {
ReloadUnits(ctx context.Context, run system.Runner) error
}
func applyService(ctx context.Context, sys system.System, r *declaration.Service, run Runner,
changed map[string]bool) (Outcome, error) {
out := begin(r)
var changes []string
// A file the service reflects changed, and it may be the unit's own file or a drop-in: the
// service manager reads those again only when told to, and a restart without it runs the unit
// it had already loaded.
if reflects(r, changed) {
if u, ok := sys.(unitReloader); ok {
if err := u.ReloadUnits(ctx, run); err != nil {
return out, fmt.Errorf("reloading the service manager's units for %s: %w", r.Unit, err)
}
}
}
// Boot first. A unit asked to be running and enabled should survive this apply failing
// half way in the more useful direction: enabled-and-stopped comes back at the next boot,
// where running-and-disabled does not.
if r.Boot != "" {
bootBefore, err := sys.ServiceBoot(ctx, run, r.Unit)
if err != nil {
return out, err
}
if bootBefore != r.Boot {
if err := sys.SetServiceBoot(ctx, run, r.Unit, r.Boot); err != nil {
return out, fmt.Errorf("setting %s to %s at boot: %w", r.Unit, r.Boot, err)
}
bootAfter, err := sys.ServiceBoot(ctx, run, r.Unit)
if err != nil {
return out, err
}
if bootAfter != r.Boot {
return out, fmt.Errorf(
"%s was asked to be %s at boot and is %s", r.Unit, r.Boot, bootAfter)
}
changes = append(changes, "boot "+bootBefore+" to "+bootAfter)
}
}
before, err := sys.ServiceState(ctx, run, r.Unit)
if err != nil {
return out, err
}
if before != r.State {
if err := sys.SetServiceState(ctx, run, r.Unit, r.State); err != nil {
return out, fmt.Errorf("setting %s to %s: %w", r.Unit, r.State, err)
}
// Read back. A service manager accepting a command says the transaction was accepted,
// not that the unit is running — one that starts and immediately dies satisfies it.
after, err := sys.ServiceState(ctx, run, r.Unit)
if err != nil {
return out, err
}
if after != r.State {
return out, fmt.Errorf("%s was asked to be %s and is %s", r.Unit, r.State, after)
}
changes = append(changes, before+" to "+after)
} else if r.State == "running" && reflects(r, changed) {
// The service is already in the state it was asked for, and something it must reflect
// changed in this same apply. A running service does not re-read its configuration, so
// leaving it alone here is how a machine ends up correct on disk and wrong in fact —
// with every check passing.
if err := sys.SetServiceState(ctx, run, r.Unit, "stopped"); err != nil {
return out, fmt.Errorf("restarting %s: stopping it: %w", r.Unit, err)
}
if err := sys.SetServiceState(ctx, run, r.Unit, "running"); err != nil {
return out, fmt.Errorf("restarting %s: starting it again: %w", r.Unit, err)
}
// Read back, for the same reason as above: a unit that starts and immediately dies
// satisfies a service manager and nothing else.
after, err := sys.ServiceState(ctx, run, r.Unit)
if err != nil {
return out, err
}
if after != "running" {
return out, fmt.Errorf(
"%s was restarted to pick up a change and is %s", r.Unit, after)
}
changes = append(changes, "restarted for "+strings.Join(reflected(r, changed), ", "))
} else if r.State == "running" && len(restartedBy(r.ReloadOn, changed)) > 0 {
// Told to read its configuration again, not stopped: for a service whose restart would
// stop what it runs — every container, for the container runtime (novox/hq ADR 0102).
reloader, ok := sys.(serviceReloader)
if !ok {
return out, fmt.Errorf("%s must be reloaded for %s and this machine's service manager "+
"cannot reload a unit", r.Unit, strings.Join(restartedBy(r.ReloadOn, changed), ", "))
}
if err := reloader.ReloadService(ctx, run, r.Unit); err != nil {
return out, fmt.Errorf("reloading %s: %w", r.Unit, err)
}
after, err := sys.ServiceState(ctx, run, r.Unit)
if err != nil {
return out, err
}
if after != "running" {
return out, fmt.Errorf("%s was reloaded to pick up a change and is %s", r.Unit, after)
}
changes = append(changes, "reloaded for "+strings.Join(restartedBy(r.ReloadOn, changed), ", "))
}
if len(changes) == 0 {
out.Action = "unchanged"
out.Detail = before
return out, nil
}
out.Action = "updated"
out.Detail = strings.Join(changes, ", ")
return out, nil
}
// remove undoes one resource the host applied and the declaration no longer names, and reports
// what it actually did.
//
// Only ever called for something in the store, which is what bounds it: the host is
// authoritative over its own footprint and inert everywhere else (novox/hq ADR 0005).
//
// It returns the action rather than assuming "removed", because for half the vocabulary the
// honest word is "forgotten". A host that reported a package removed when it left the package
// installed would be describing an effect it declined to have.
func remove(ctx context.Context, sys system.System, a store.Applied, run Runner) (string, string, error) {
switch declaration.Type(a.Type) {
case declaration.TypeDirectory:
// **A directory with anything left in it is kept, and that is the rule that protects
// data.** Everything the mesh put inside is itself a declared resource, and orphans are
// removed in reverse declaration order — so by the time a directory comes to be removed,
// what the mesh wrote there is already gone. Anything still present is something nobody
// declared: a database's files, a mail spool, somebody's uploads.
//
// Before this, unassigning a module deleted its data directory and everything under it,
// and the report said "removed". Nothing anywhere said what had been in there.
//
// This is the host's own line, applied to the one shape where getting it wrong is not
// recoverable: it removes what it made and leaves what it merely configured. An empty
// directory is what it made. A full one is not.
entries, err := os.ReadDir(a.Target)
if errors.Is(err, os.ErrNotExist) {
return "forgotten", "no longer there", nil
}
if err != nil {
return "", "", err
}
if len(entries) > 0 {
return "kept", fmt.Sprintf(
"no longer declared, and %d item(s) inside that the mesh did not put there — "+
"remove it by hand once you know what it is", len(entries)), nil
}
if err := os.Remove(a.Target); err != nil {
return "", "", err
}
return "removed", "no longer declared, and empty", nil
case declaration.TypeFile:
if a.Into != nil {
return removeInto(a)
}
if err := os.RemoveAll(a.Target); err != nil {
return "", "", err
}
if _, err := os.Stat(a.Target); !errors.Is(err, os.ErrNotExist) {
return "", "", fmt.Errorf("%s is still there after removing it", a.Target)
}
return "removed", "no longer declared", nil
case declaration.TypeService:
// A unit that is no longer declared is stopped, not deleted. The host did not install
// it and does not own the unit file — only the state it put the unit into.
//
// A unit that no longer EXISTS is already in the state removal is trying to reach, and
// saying so matters: stopping it fails, and a failure here fails the whole apply. A
// host holding a record of an uninstalled unit would then be unable to apply anything,
// ever, with no way out but editing its state by hand. Removal is idempotent for the
// same reason `os.RemoveAll` is.
if _, err := sys.ServiceState(ctx, run, a.Target); err != nil {
if strings.Contains(err.Error(), "does not exist on this machine") {
return "forgotten", "the unit no longer exists", nil
}
return "", "", err
}
if err := sys.SetServiceState(ctx, run, a.Target, "stopped"); err != nil {
return "", "", fmt.Errorf("stopping %s: %w", a.Target, err)
}
return "removed", "stopped; the unit file is not the host's to delete", nil
case declaration.TypeContainer:
// The host CREATED this one, so the host removes it. That is the line: it removes what
// it made and leaves what it merely configured.
if _, err := run(ctx, "docker", "rm", "-f", a.Target); err != nil {
// Already gone is the state removal wants. Anything else is a real failure.
if _, alive := containerState(ctx, a.Target, run); alive == nil {
return "", "", fmt.Errorf("removing container %s: %w", a.Target, err)
}
}
if _, err := containerState(ctx, a.Target, run); err == nil {
return "", "", fmt.Errorf("container %s is still there after removing it", a.Target)
}
return "removed", "no longer declared", nil
case declaration.TypePackage:
// Deliberately not uninstalled, and this is a decision rather than an omission.
//
// The host cannot know what else on this machine needs the package. Uninstalling a
// container runtime because a declaration changed would stop every container on the
// node, and the machine may have had the package before the mesh ever saw it
// (novox/hq research 012: adopted, not installed). Undeclaring says "the mesh no
// longer requires this", which is not the same as "remove it".
return "forgotten", "left installed; the host does not uninstall what it cannot know is unused", nil
case declaration.TypeAction:
// An action has no footprint the host can undo — it ran, and whatever it did belongs
// to whatever it acted on.
return "forgotten", "an action leaves nothing the host owns", nil
case declaration.TypeAccess:
// The path is the operator's and the host never owned it (novox/hq ADR 0051). Undeclaring
// it says only that this module no longer reaches it — not that the media library should
// be touched. So the record is dropped and the path left exactly as it is; removing it
// would be the data loss ADR 0030 exists to prevent, on a directory the mesh never made.
return "forgotten", "an operator-owned path is never the host's to remove", nil
case declaration.TypeNetwork:
// **The reason this is a shape at all** (novox/hq ADR 0029). Orphans are removed in
// reverse declaration order, so a network written before the containers that join it is
// removed after they are gone — and a runtime refusing to remove one still in use is
// reported rather than swallowed, because that means something the mesh did not declare
// is holding it.
cri, err := containerRuntime(ctx, run)
if err != nil {
return "", "", fmt.Errorf("%w, so the network %q cannot be removed", err, a.Target)
}
if _, err := run(ctx, cri, "network", "inspect", a.Target); err != nil {
return "forgotten", "no longer there", nil
}
if _, err := run(ctx, cri, "network", "rm", a.Target); err != nil {
return "", "", fmt.Errorf("cannot remove the network %q: %w", a.Target, err)
}
return "removed", "no longer declared", nil
default:
return "", "", fmt.Errorf("no way to remove a %q", a.Type)
}
}
// ExecRunner runs a real command, with stdin closed and output captured.
func ExecRunner(ctx context.Context, name string, args ...string) (string, error) {
cmd := exec.CommandContext(ctx, name, args...)
cmd.Stdin = nil
out, err := cmd.Output()
if err != nil {
var exit *exec.ExitError
if errors.As(err, &exit) {
return string(out), fmt.Errorf("%s exited %d: %s",
name, exit.ExitCode(), strings.TrimSpace(string(exit.Stderr)))
}
return string(out), fmt.Errorf("%s: %w", name, err)
}
return string(out), nil
}
// applyPackage installs a package the machine does not have.
//
// It never upgrades and never removes. "Present" is the whole of what a package resource
// asserts, because version is the package manager's business and the mesh does not have a
// second opinion about it (novox/hq ADR 0005 — the host depends on nothing, and that includes
// not becoming a second package manager).
func applyPackage(ctx context.Context, sys system.System, r *declaration.Package, run Runner) (Outcome, error) {
out := begin(r)
installed, err := sys.PackageInstalled(ctx, run, r.Package)
if err != nil {
return out, err
}
if installed {
out.Action = "unchanged"
out.Detail = "already installed"
return out, nil
}
if err := sys.InstallPackage(ctx, run, r.Package); err != nil {
return out, fmt.Errorf("installing %s: %w", r.Package, err)
}
// Read back. A package manager exiting zero says the transaction was accepted.
installed, err = sys.PackageInstalled(ctx, run, r.Package)
if err != nil {
return out, err
}
if !installed {
return out, fmt.Errorf(
"%s was installed without error and the package database does not have it", r.Package)
}
out.Action = "created"
return out, nil
}
// Labels the host puts on every container it creates.
//
// specLabel carries a digest of the declaration that made the container. It is what lets a
// reconcile answer "is this container the one the current declaration describes" without
// comparing every field the runtime reports — which cannot be done reliably, because a runtime
// normalises, defaults and reorders what it is given, and the differences that produces are
// indistinguishable from real drift.
const (
specLabel = "mesh-host.spec"
idLabel = "mesh-host.id"
)
// containerSpec is the identity of a declared container: everything that, if changed, means
// the running container is no longer what was asked for.
func containerSpec(r *declaration.Container, declares map[string]string) string {
keys := make([]string, 0, len(r.Env))
for k := range r.Env {
keys = append(keys, k)
}
sort.Strings(keys)
var b strings.Builder
b.WriteString(r.Image + "\n" + r.Name + "\n")
for _, k := range keys {
b.WriteString("env " + k + "=" + r.Env[k] + "\n")
}
for _, p := range r.Ports {
b.WriteString("port " + p + "\n")
}
for _, v := range r.Volumes {
b.WriteString("volume " + v + "\n")
}
for _, a := range r.Args {
b.WriteString("arg " + a + "\n")
}
// The cadence is part of what was declared, so a changed schedule is a changed spec — the marker
// moves and the install is reported "updated" and re-established. Added only when present, so no
// ordinary container's or run-once step's digest moves for a field it does not set.
if r.Schedule != "" {
b.WriteString("schedule " + r.Schedule + "\n")
}
// **What this container reads is part of what it is.**
//
// A container takes its environment and its mounted files once, at start, and never looks
// again. Comparing only the fields above meant a container whose configuration had since been
// rewritten compared equal and was left alone — running values the machine no longer holds,
// while every check reported success (novox/hq 04-ISSUES/045). Naming what it depends on here
// makes that comparison standing rather than a tripwire that fires during one apply and never
// again. Sorted, so the digest does not move for a reordering nobody made.
depends := append([]string{}, r.RestartOn...)
sort.Strings(depends)
for _, id := range depends {
b.WriteString("reads " + id + "=" + declares[id] + "\n")
}
return fmt.Sprintf("%x", sha256.Sum256([]byte(b.String())))
}
// containerState reports whether a container is running and which spec made it.
// The error means the container does not exist.
func containerState(ctx context.Context, name string, run Runner) (state struct {
Running bool
Spec string
}, err error) {
out, err := run(ctx, "docker", "inspect", "--format",
"{{.State.Running}}\t{{index .Config.Labels \""+specLabel+"\"}}", name)
if err != nil {
return state, fmt.Errorf("no container named %s", name)
}
running, spec, _ := strings.Cut(strings.TrimSpace(out), "\t")
state.Running = running == "true"
state.Spec = strings.TrimSpace(spec)
return state, nil
}
// applyContainer makes the declared container the one that is running.
//
// There is no "update" for a container: a container's configuration is fixed when it is
// created, so any change is a replacement. Saying that plainly is better than a partial
// in-place update that leaves the running thing half-declared.
// applyNetwork creates a named network if the machine does not already have one.
//
// **Existence is the whole of the state.** A network the mesh declared and a network somebody
// made by hand are indistinguishable by name, and that is deliberate: the mesh owns the name, not
// the thing, so it will not tear down and rebuild one that is already there and working. What it
// records is that this resource is now present, which is what lets it be removed later.
//
// Nothing is reconciled beyond presence. A driver or a subnet changed underneath would not be
// noticed — and is not declarable either (novox/hq ADR 0029), so there is nothing to disagree
// with.
func applyNetwork(ctx context.Context, r *declaration.Network, run Runner) (Outcome, error) {
out := begin(r)
cri, err := containerRuntime(ctx, run)
if err != nil {
return out, fmt.Errorf("%w, so nothing can be said about the network %q", err, r.Name)
}
if _, err := run(ctx, cri, "network", "inspect", r.Name); err == nil {
out.Action = "unchanged"
out.Detail = "already there"
return out, nil
}
if _, err := run(ctx, cri, "network", "create", r.Name); err != nil {
return out, fmt.Errorf("cannot create the network %q: %w", r.Name, err)
}
// Read back rather than trusting the exit status (novox/hq ADR 0018). A runtime that reports
// success and made nothing leaves every container that joins it failing to start, with the
// cause one step away.
if _, err := run(ctx, cri, "network", "inspect", r.Name); err != nil {
return out, fmt.Errorf(
"the network %q was created and is not there afterwards: %w", r.Name, err)
}
out.Action = "created"
out.Detail = "a network for this module's own containers"
return out, nil
}
func applyContainer(ctx context.Context, r *declaration.Container, run Runner,
changed map[string]bool, declares map[string]string, previous store.Applied) (Outcome, error) {
out := begin(r)
want := containerSpec(r, declares)
cri, err := containerRuntime(ctx, run)
if err != nil {
return out, fmt.Errorf("%w, so nothing can be said about %q", err, r.Name)
}
// A scheduled step is state that is present, not a container to start (novox/hq ADR 0053).
// Installing it records the schedule and reports the node current at once — the deliberate
// inversion of run-once, which gates. The recurring run is fired by the host's Scheduler off the
// clock, re-established from this declaration each apply, and NEVER here — so installing does not
// RUN the container.
//
// It does, however, ensure the pinned image is present now. A service or a run-once container
// gets its image as a side effect of `docker run`; a scheduled step is never run at apply, so
// without this its image would be absent from the node until the first scheduled fire — which
// would pay the whole pull latency then, and leave tooling that expects the image present after
// apply looking at a node that does not have it. Ensuring the image is not running it, so the
// no-run invariant holds.
if r.Schedule != "" {
if err := ensureImage(ctx, cri, r.Image, run); err != nil {
return out, err
}
return applySchedule(r, want, previous)
}
if r.RunOnce {
return applyRunOnce(ctx, r, run, cri, want, previous)
}
before, err := containerState(ctx, r.Name, run)
existed := err == nil
// A container reads a mounted file once, at start. When one of its restart-on resources changed
// this pass — a settings-merged config the runtime read, say — the file on disk is new and the
// running process still holds the old value, and the spec (image, env, volumes) has not moved,
// so the plain "spec matches, leave it" below would keep the stale process for ever
// (novox/hq 04-ISSUES/009). Recreating is how a container gets restart-on, which a service
// already has.
reasons := restartedBy(r.RestartOn, changed)
switch {
case existed && before.Spec == want && before.Running && len(reasons) == 0:
out.Action = "unchanged"
return out, nil
case existed:
if _, err := run(ctx, cri, "rm", "-f", r.Name); err != nil {
return out, fmt.Errorf("replacing container %s: %w", r.Name, err)
}
}
args := []string{"run", "--detach", "--name", r.Name, "--restart", "unless-stopped"}
for _, file := range r.EnvFile {
args = append(args, "--env-file", file)
}
if r.Network != "" {
args = append(args, "--network", r.Network)
}
args = append(args,
"--label", specLabel+"="+want, "--label", idLabel+"="+r.ID)
for _, k := range sortedKeys(r.Env) {
args = append(args, "--env", k+"="+r.Env[k])
}
for _, p := range r.Ports {
args = append(args, "--publish", p)
}
for _, v := range r.Volumes {
args = append(args, "--volume", v)
}
for _, h := range r.Hosts {
// Written into the container's own hosts file by the runtime. Per container rather than
// by editing the machine's resolver configuration: that file belongs to something else on
// most machines, and a host that edited it would be fighting whatever owns it on every
// boot — the fault this host exists to avoid, in the place it would be hardest to see.
args = append(args, "--add-host", h)
}
args = append(args, r.Image)
args = append(args, r.Args...)
if _, err := run(ctx, cri, args...); err != nil {
return out, fmt.Errorf("starting container %s: %w", r.Name, err)
}
// Read back. `docker run --detach` returning an id says the container was created, not
// that it is still running — a container whose entrypoint exits immediately satisfies the
// command exactly as one that came up does.
after, err := containerState(ctx, r.Name, run)
if err != nil {
return out, fmt.Errorf("started container %s and it is not there: %w", r.Name, err)
}
if !after.Running {
return out, fmt.Errorf(
"container %s was started and is not running. It exited; ask the runtime for its "+
"logs", r.Name)
}
if after.Spec != want {
return out, fmt.Errorf("container %s is not the one that was declared after creating it", r.Name)
}
out.Action = "created"
if existed {
out.Action = "updated"
if len(reasons) > 0 {
out.Detail = "recreated to pick up " + strings.Join(reasons, ", ")
} else {
out.Detail = "replaced; a container's configuration is fixed when it is created"
}
}
return out, nil
}
// applyRunOnce runs a container to completion, once, and requires it to exit 0 (novox/hq ADR 0052).
//
// It is a step, not a service: the module's own code seeding a store, migrating a schema or
// gating on health, under the module's own account (ADR 0047), before the container that depends
// on it. Three things make it a step rather than an ordinary container:
//
// - It is run in the foreground, so the runtime returns the container's exit code. A non-zero
// exit is an error here, and — because a run-once step gates the apply the way a failed action
// does — that error stops everything the declaration places after it. That is how "before the
// broker starts" is enforced: the step is declared first, and the broker is never reached
// until it has completed.
// - Its record of having happened is the digest of its declaration, recorded by the caller only
// after it exits 0 (ADR 0018). The step leaves nothing running to inspect, so the persisted
// digest — not a live container — is the marker. A re-apply whose declaration digest already
// matches does nothing; a changed declaration re-runs it.
// - It writes a seed only if it is absent and never reconciles it, so what a running program
// grows in that seed afterward is never wiped (04-ISSUES/035). The host's marker keeps the
// step from re-running; the step's own code keeps it from clobbering on the pass it does run.
func applyRunOnce(ctx context.Context, r *declaration.Container, run Runner, cri, want string, previous store.Applied) (Outcome, error) {
out := begin(r)
// Already completed for this exact declaration. The marker is the store, because a run-once
// step leaves nothing running to ask.
if previous.Wrote != "" && previous.Wrote == want {
out.Action = "unchanged"
out.Detail = "run-once step already completed for this declaration"
out.wrote = want
return out, nil
}
// A container by this name from a previous, different declaration must not linger and be
// mistaken for this run. Removing a name that is not there is the state we want, so its error
// is ignored.
_, _ = run(ctx, cri, "rm", "-f", r.Name)
// Run in the foreground so the runtime waits for the container and hands back its exit code.
// No --detach and no --restart: a step that is restarted is not a step.
args := foregroundRunArgs(r, want)
if _, err := run(ctx, cri, args...); err != nil {
// A non-zero exit or a runtime that could not start it. Either way the step did not make
// the machine ready, so the apply must not go on to the container that needs it.
return out, fmt.Errorf("run-once step %s did not complete: %w", r.Name, err)
}
// It completed. Remove the exited container so a later apply is not confused by a stopped one;
// the record that it ran is the digest below, which the caller persists after the fact.
_, _ = run(ctx, cri, "rm", "-f", r.Name)
out.Action = "created"
out.Detail = "run-once step completed"
out.wrote = want
return out, nil
}
// foregroundRunArgs builds a `docker run` that runs a container to completion and hands back its
// exit code — no --detach, no --restart, because a step that is restarted is not a step. Shared by
// a run-once step (novox/hq ADR 0052) and by one fire of a scheduled step (ADR 0053), which are the
// same "run the container and let it exit" up to how often it happens.
func foregroundRunArgs(r *declaration.Container, want string) []string {
args := []string{"run", "--name", r.Name}
for _, file := range r.EnvFile {
args = append(args, "--env-file", file)
}
if r.Network != "" {
args = append(args, "--network", r.Network)
}
args = append(args, "--label", specLabel+"="+want, "--label", idLabel+"="+r.ID)
for _, k := range sortedKeys(r.Env) {
args = append(args, "--env", k+"="+r.Env[k])
}
for _, v := range r.Volumes {
args = append(args, "--volume", v)
}
for _, h := range r.Hosts {
args = append(args, "--add-host", h)
}
args = append(args, r.Image)
args = append(args, r.Args...)
return args
}
// ensureImage makes the pinned image present on the node without running anything.
//
// A service or a run-once container gets its image as a side effect of `docker run` — the first run
// pulls it. A scheduled step is installed but deliberately never run at apply (novox/hq ADR 0053), so
// nothing would pull its image until the first scheduled fire: the image is absent from the node
// right after a successful apply, the first run pays the whole pull latency, and tooling that expects
// the image present after apply finds it missing. This fetches the same bytes `docker run` would, and
// stops short of starting the container.
//
// Idempotent, and it reads back (novox/hq ADR 0018): an image already present is left as is, and a
// pull that reported success but left nothing there is a failure, not a convergence.
func ensureImage(ctx context.Context, cri, image string, run Runner) error {
if _, err := run(ctx, cri, "image", "inspect", image); err == nil {
return nil
}
// An image named by the digest of its own configuration is one this machine was supposed to
// already hold — built here, or handed over. There is no registry that answers for it, so
// pulling would fail somewhere that names a network problem instead of a missing image.
if strings.HasPrefix(image, "sha256:") {
return fmt.Errorf(
"image %s is not on this machine, and an image named by its own digest cannot be "+
"fetched: nothing serves it. Build it here, or load it, before applying this",
image)
}
if _, err := run(ctx, cri, "pull", image); err != nil {
return fmt.Errorf("pulling image %s: %w", image, err)
}
if _, err := run(ctx, cri, "image", "inspect", image); err != nil {
return fmt.Errorf("image %s is not present after pulling it: %w", image, err)
}
return nil
}
// applySchedule installs a scheduled step: it records the schedule as present and reports the node
// current, without running anything (novox/hq ADR 0053).
//
// This is the deliberate inversion of run-once. A run-once step gates the apply — it runs to
// completion here and a non-zero exit halts everything after it — because "seed the store before the
// broker starts" is a precondition of convergence. A scheduled step is the opposite: it runs *after*
// the machine is up, on its own clock, and a single failed run is an ordinary operational event. So
// installing it is pure state: the schedule is present, like a running service, and the apply is
// current at once. The recurring run is fired by the host's Scheduler off the clock (see
// schedule.go), re-established from the applied declaration each pass because the declaration is the
// source of truth (ADR 0018) — never from here, and never persisted beyond what the mesh already
// owns.
//
// The marker is the declaration's digest, exactly as for run-once, so a re-apply of the same
// declaration reports the schedule unchanged and a changed image, environment or cadence reports it
// re-installed. A failed run touches none of this: it happens entirely in the Scheduler, outside the
// apply and the store, which is why a routine job's failure can never flip the node's state.
func applySchedule(r *declaration.Container, want string, previous store.Applied) (Outcome, error) {
out := begin(r)
out.wrote = want
switch {
case previous.Wrote == want:
out.Action = "unchanged"
out.Detail = "scheduled step; already installed for this declaration"
case previous.Wrote != "":
out.Action = "updated"
out.Detail = "scheduled step re-installed; its image, environment or cadence changed"
default:
out.Action = "created"
out.Detail = "scheduled step installed; the host runs it on its cadence"
}
return out, nil
}
// restartedBy is which of the named resources changed this pass — the reason a container or service
// must be brought back rather than left as it is (novox/hq 04-ISSUES/009).
func restartedBy(restartOn []string, changed map[string]bool) []string {
var which []string
for _, id := range restartOn {
if changed[id] {
which = append(which, id)
}
}
return which
}
func sortedKeys(m map[string]string) []string {
keys := make([]string, 0, len(m))
for k := range m {
keys = append(keys, k)
}
sort.Strings(keys)
return keys
}
// applyAction runs something the bundle declared, and never learns what it means.
//
// Verify does double duty, and that is the design rather than a convenience: it is both the
// idempotency check and the read-back. Running it first is how the host knows whether there is
// anything to do — it does not know what a database is, so "is the database there" is a
// question only the declaration can ask. Running it again afterwards is how the host knows the
// command had the effect it claimed (novox/hq ADR 0005).
func applyAction(ctx context.Context, r *declaration.Action, run Runner) (Outcome, error) {
out := begin(r)
if _, err := runAction(ctx, r, r.Verify, run); err == nil {
out.Action = "unchanged"
out.Detail = "already true"
return out, nil
}
if _, err := runAction(ctx, r, r.Command, run); err != nil {
return out, fmt.Errorf("running the action: %w", err)
}
if _, err := runAction(ctx, r, r.Verify, run); err != nil {
return out, fmt.Errorf(
"the action ran without error and its own verify still fails: %w\n\n"+
"The command reported success and the thing it was for did not happen, which "+
"is exactly what verify exists to catch", err)
}
out.Action = "created"
out.Detail = "verify was false and is now true"
return out, nil
}
// runAction runs one of an action's command lines, on the machine or inside a container.
func runAction(ctx context.Context, r *declaration.Action, argv []string, run Runner) (string, error) {
if len(argv) == 0 {
return "", errors.New("no command")
}
if r.In != "" {
return run(ctx, "docker", append([]string{"exec", r.In}, argv...)...)
}
return run(ctx, argv[0], argv[1:]...)
}
// Container runtimes the host knows how to ask.
//
// Two, because two exist on machines the mesh runs on. The list is short on purpose: each entry
// is a claim that its probe and its CLI have been checked, not that a binary of that name might
// work (novox/hq ADR 0005).
//
// The probe differs and the rest does not, which is what makes this a lookup rather than an
// interface. `docker info --format {{.ServerVersion}}` fails on podman — the field does not
// exist in its report — while `run`, `inspect --format` and `rm -f` are identical, including
// docker's own Go template syntax for reading state and labels.
var containerRuntimes = []struct {
command string
// probe asks the runtime for its version in the form THAT runtime understands. It must
// prove the runtime is FUNCTIONING, never that a binary is on disk
// (novox/hq 04-ISSUES/007).
probe []string
}{
{command: "docker", probe: []string{"info", "--format", "{{.ServerVersion}}"}},
{command: "podman", probe: []string{"info", "--format", "{{.Version.Version}}"}},
}
// containerRuntime returns the runtime this machine actually has, or says there is none.
//
// Detected rather than declared, because a machine already carrying one keeps it: adoption
// takes over what is there rather than replacing it (novox/hq research 012). Which runtime a
// machine has is reported upward in the profile; what to install on a machine with none is the
// control plane's decision, not this one's.
func containerRuntime(ctx context.Context, run Runner) (string, error) {
var tried []string
for _, rt := range containerRuntimes {
if _, err := run(ctx, rt.command, rt.probe...); err == nil {
return rt.command, nil
}
tried = append(tried, rt.command)
}
return "", fmt.Errorf(
"no container runtime answers on this machine (tried %s)", strings.Join(tried, ", "))
}
// digestOf is how this host recognises what it wrote.
//
// A digest rather than the content: the store is read on every reconcile and sits beside the
// state on disk, and keeping every managed file twice would make it grow with the machine rather
// than with the number of resources.
func digestOf(content string) string {
sum := sha256.Sum256([]byte(content))
return hex.EncodeToString(sum[:])
}
// DigestOf names a declaration by its bytes, exactly as the mesh names what it sends: sha256 of
// the raw bytes, hex. The two sides never digest different things.
func DigestOf(raw []byte) string {
sum := sha256.Sum256(raw)
return hex.EncodeToString(sum[:])
}
// holds is the machine's own ports a resource occupies.
//
// **What the declaration binds, not what is open.** A machine's open ports are a moving target —
// something a person started, a connection the kernel handed out — and assigning around them would
// mean a port that was free when it was asked for and taken when it was used. What a resource
// declares is stable, and it is the half the mesh can be responsible for.
func holds(resource declaration.Resource) []int {
container, ok := resource.(*declaration.Container)
if !ok {
return nil
}
var out []int
for _, mapping := range container.Ports {
// "8080:80", or "127.0.0.1:8080:80" when an address was named. The machine's port is the
// one before the last colon; the last is inside the container and is not the machine's.
parts := strings.Split(mapping, ":")
if len(parts) < 2 {
continue
}
port, err := strconv.Atoi(strings.TrimSpace(parts[len(parts)-2]))
if err != nil {
continue
}
out = append(out, port)
}
return out
}
// declaredDigest is what a resource currently says it should be.
//
// **Content, not identity.** It exists so a container can be told apart by what it reads: a file
// whose text changed must produce a different digest, or the container mounting it compares equal
// to one started against the old text. Only the shapes something can read are digested; for
// everything else the identity is enough, because nothing mounts a package.
func declaredDigest(r declaration.Resource) string {
var material string
switch res := r.(type) {
case *declaration.File:
// The content as declared, before any sealing is opened — two machines are given different
// ciphertext for the same secret, and digesting that would make an unchanged file look
// changed on every apply and restart the container reading it for ever.
material = res.Content
case *declaration.Directory:
material = res.Path + "\n" + res.Mode
default:
return ""
}
return fmt.Sprintf("%x", sha256.Sum256([]byte(material)))
}