1547 lines
62 KiB
Go
1547 lines
62 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.
|
|
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}
|
|
}
|
|
|
|
for _, orphan := range known.Orphans(declared, origin) {
|
|
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 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{}
|
|
|
|
// **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 _, resource := range d.Resources {
|
|
// **On an adopted node, what is found is kept until its module is taken** (novox/hq ADR
|
|
// 0100). Before anything is applied: a file present with no record of this host writing
|
|
// it, or a container present under that name that no host made, is held as it is and
|
|
// reported. Once held it stays held — changed or gone — 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 && holdable(resource) {
|
|
if module, untaken := d.Adoption.UntakenModuleOf(resource.Identity()); untaken {
|
|
was, already := known.HeldAt(resource.Identity())
|
|
isFound := false
|
|
if !already {
|
|
var err error
|
|
if isFound, err = found(ctx, resource, run, known); 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 already || isFound {
|
|
outcome, held, err := hold(ctx, resource, module, was, already, run, keep, 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(), outcome.Target, err))
|
|
continue
|
|
}
|
|
known.RecordHeld(held)
|
|
report.Outcomes = append(report.Outcomes, outcome)
|
|
if !already || held.Changed != was.Changed {
|
|
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 {
|
|
outcome, err = applyOne(ctx, sys, resource, run, changed, declares, 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,
|
|
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))
|
|
}
|
|
}
|
|
|
|
// 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}
|
|
}
|
|
}
|
|
|
|
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, declares map[string]string, 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, 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
|
|
}
|
|
|
|
func applyFile(r *declaration.File, previous store.Applied, unseal Unseal) (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()
|
|
|
|
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)
|
|
}
|
|
|
|
// 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[:])
|
|
}
|
|
|
|
// 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)))
|
|
}
|