Files
mesh-host/internal/apply/apply.go
T
jschoubben 19e5dd83ea apply: a run-once container is a step the host runs to completion (ADR 0052)
A module can declare state but not a step that runs at first boot. This adds
`run-once: true` to the container shape: the host runs it in the foreground,
requires it to exit 0, and records that it did — as the digest of the
declaration, so a re-apply does not re-run it unless the declaration changed.

Because the declaration is applied in order and a failed run-once step gates the
apply the way a failed action does, whatever is declared after the step starts
only once it has completed. That is how "before the broker starts" is enforced,
with no dependency graph the host must resolve (ADR 0005): the step is declared
first, and the container that needs it is never reached until it is done.

No new host shape and no arbitrary host command — a run-once container is
strictly less powerful than an action. Validation refuses run-once with
restart-on (contradictory lifecycles). Six unit tests; go test ./... green.

Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
2026-09-05 23:55:45 +02:00

1248 lines
48 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
}
// 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 {
if o.Action != "unchanged" {
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.
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) {
if log == nil {
log = func(string) {}
}
report := Report{}
declared := map[string]bool{}
for _, r := range d.Resources {
declared[r.Identity()] = true
}
for _, orphan := range known.Orphans(declared, origin) {
action, detail, err := remove(ctx, sys, orphan, run)
if err != nil {
return report, known, &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))
}
// 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{}
// 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 _, resource := range d.Resources {
was, _ := known.Find(resource.Identity())
outcome, err := applyOne(ctx, sys, resource, run, changed, was, unseal)
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,
Holds: holds(resource),
})
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 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
}
// 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, previous store.Applied, unseal Unseal) (Outcome, error) {
switch res := r.(type) {
case *declaration.Directory:
return applyDirectory(res)
case *declaration.File:
return applyFile(res, previous, unseal)
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, previous)
case *declaration.User:
return applyUser(ctx, sys, res, run)
case *declaration.Archive:
return applyArchive(ctx, res, 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
}
func applyFile(r *declaration.File, previous store.Applied, unseal Unseal) (Outcome, error) {
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
}
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()
if !contentSame {
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"
}
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)
}
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
// 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), ", "))
}
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 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) 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")
}
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, previous store.Applied) (Outcome, error) {
out := begin(r)
want := containerSpec(r)
cri, err := containerRuntime(ctx, run)
if err != nil {
return out, fmt.Errorf("%w, so nothing can be said about %q", err, r.Name)
}
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 := []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...)
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
}
// 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[:])
}
// 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
}