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