Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b602b223e1 | ||
|
|
e8c4824ae2 | ||
|
|
8c2e76f4c1 | ||
|
|
8280a82ef8 | ||
|
|
2d5e76434b | ||
|
|
d9ea387680 | ||
|
|
99ad145bec | ||
|
|
2ff3b50a84 | ||
|
|
a24670d77c | ||
|
|
5fc4052a2b |
+166
-2
@@ -247,6 +247,20 @@ func ApplyKeeping(
|
||||
// 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
|
||||
// **What a declared process replaces is handed over, not removed first** (novox/hq issue 213).
|
||||
// Every other orphan goes before anything is applied; one a process names under `replaces` is
|
||||
// kept until that process is applied and still running a moment later, so whatever it was —
|
||||
// the controller's container — answers until its replacement does, and goes on answering if
|
||||
// the replacement never comes up.
|
||||
replacedBy := map[string]string{}
|
||||
for _, r := range d.Resources {
|
||||
if p, ok := r.(*declaration.Process); ok {
|
||||
for _, id := range p.Replaces {
|
||||
replacedBy[id] = p.ID
|
||||
}
|
||||
}
|
||||
}
|
||||
handover := map[string][]store.Applied{}
|
||||
for _, orphan := range known.Orphans(declared, origin) {
|
||||
if d.Adoption == nil && strings.HasPrefix(orphan.ID, declaration.AdoptionPrefix) {
|
||||
protecting = append(protecting, orphan)
|
||||
@@ -264,6 +278,10 @@ func ApplyKeeping(
|
||||
log(fmt.Sprintf(" kept %s (%s): %s was left out of this declaration by the mesh, not removed", orphan.ID, orphan.Target, module))
|
||||
continue
|
||||
}
|
||||
if by, replaced := replacedBy[orphan.ID]; replaced {
|
||||
handover[by] = append(handover[by], orphan)
|
||||
continue
|
||||
}
|
||||
orphans = append(orphans, orphan)
|
||||
}
|
||||
ordered := d.Resources
|
||||
@@ -517,6 +535,11 @@ func ApplyKeeping(
|
||||
if c, ok := resource.(*declaration.Container); ok && c.RunOnce {
|
||||
gates = true
|
||||
}
|
||||
// And a run-once process, which is the same step hosted as a unit: a version whose
|
||||
// preparation did not complete must not be started (novox/hq ADR 0135, issue 213).
|
||||
if p, ok := resource.(*declaration.Process); ok && p.RunOnce {
|
||||
gates = true
|
||||
}
|
||||
if gates {
|
||||
failed.Gated = true
|
||||
// **A module's step gates that module, not the machine** (novox/hq ADR 0136).
|
||||
@@ -607,6 +630,28 @@ func ApplyKeeping(
|
||||
}
|
||||
log(line)
|
||||
}
|
||||
|
||||
// Its replacement applied: what it replaces goes now, once it is seen running (issue 213).
|
||||
if waiting, has := handover[resource.Identity()]; has {
|
||||
delete(handover, resource.Identity())
|
||||
if err := handOver(ctx, run, resource, waiting, removeOrphan, &report, log); err != nil {
|
||||
failures = append(failures, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A replacement that did not apply — failed, skipped behind its module's step, held — leaves
|
||||
// what it replaces running and recorded, said, for the next apply to hand over.
|
||||
unhanded := make([]string, 0, len(handover))
|
||||
for by := range handover {
|
||||
unhanded = append(unhanded, by)
|
||||
}
|
||||
sort.Strings(unhanded)
|
||||
for _, by := range unhanded {
|
||||
for _, orphan := range handover[by] {
|
||||
keptForReplacement(&report, log, orphan,
|
||||
fmt.Sprintf("kept: %s, which replaces it, did not apply", by))
|
||||
}
|
||||
}
|
||||
|
||||
if !orphansRemoved {
|
||||
@@ -1413,6 +1458,15 @@ func remove(ctx context.Context, sys system.System, a store.Applied, run Runner,
|
||||
// undo is never reported as done; not fatal for a former target, which was never dropped by anyone.
|
||||
var errNoRemoval = errors.New("no way to remove")
|
||||
|
||||
// exited is a command that ran and exited non-zero: its words, and the exit itself.
|
||||
type exited struct {
|
||||
words string
|
||||
exit *exec.ExitError
|
||||
}
|
||||
|
||||
func (e *exited) Error() string { return e.words }
|
||||
func (e *exited) Unwrap() error { return e.exit }
|
||||
|
||||
// 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...)
|
||||
@@ -1421,8 +1475,12 @@ func ExecRunner(ctx context.Context, name string, args ...string) (string, error
|
||||
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)))
|
||||
// The words as they always were, and the exit underneath them, so a caller asks the
|
||||
// code (system.ExitCode) rather than matching text that this line is free to reword.
|
||||
return string(out), &exited{
|
||||
words: fmt.Sprintf("%s exited %d: %s", name, exit.ExitCode(), strings.TrimSpace(string(exit.Stderr))),
|
||||
exit: exit,
|
||||
}
|
||||
}
|
||||
return string(out), fmt.Errorf("%s: %w", name, err)
|
||||
}
|
||||
@@ -1649,6 +1707,13 @@ func containerSpecReading(r *declaration.Container, declares, reads map[string]s
|
||||
// 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.
|
||||
// And which of its module's containers it holds still while it runs (novox/hq ADR 0189), for
|
||||
// the same reason: a declaration that changed the window while the machine reported no change
|
||||
// would be a machine quietly holding yesterday's containers. In declared order, which is the
|
||||
// order they are stopped in.
|
||||
for _, id := range r.WhileStopped {
|
||||
b.WriteString("while-stopped " + id + "\n")
|
||||
}
|
||||
if r.Schedule != "" {
|
||||
b.WriteString("schedule " + r.Schedule + "\n")
|
||||
}
|
||||
@@ -2453,3 +2518,102 @@ func moduleOf(identity string) (string, bool) {
|
||||
}
|
||||
return identity[:at], true
|
||||
}
|
||||
|
||||
// handOver removes what a process replaces, once the process is running and still is a moment later
|
||||
// (novox/hq issue 213, ADR 0184). A replacement that is not up keeps what it replaces in place
|
||||
// and recorded, and is this apply's failure: the old one answers until a later apply finds the new
|
||||
// one up.
|
||||
func handOver(ctx context.Context, run Runner, resource declaration.Resource,
|
||||
waiting []store.Applied, removeOrphan func(store.Applied) error, report *Report,
|
||||
log func(string)) *Error {
|
||||
p, ok := resource.(*declaration.Process)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
if why := stillUp(ctx, run, p.Name+".service"); why != "" {
|
||||
for _, orphan := range waiting {
|
||||
keptForReplacement(report, log, orphan,
|
||||
fmt.Sprintf("kept: %s, which replaces it, is not running (%s)", p.ID, why))
|
||||
}
|
||||
return &Error{Resource: p.ID, Err: fmt.Errorf(
|
||||
"%s was started and is not running a moment later (%s), so what it replaces was kept: %s",
|
||||
p.Name, why, appliedIDs(waiting)), Done: *report}
|
||||
}
|
||||
for _, orphan := range waiting {
|
||||
if err := removeOrphan(orphan); err != nil {
|
||||
var failed *Error
|
||||
if errors.As(err, &failed) {
|
||||
return failed
|
||||
}
|
||||
return &Error{Resource: orphan.ID, Err: err, Done: *report}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// handoverSettle is how long a replacement must stay up before what it replaces goes. Longer than a
|
||||
// service's second look (serviceSettle): this one decides whether the last thing answering is
|
||||
// removed, and a process that fails on its store or its bus does so after it opened them, not in the
|
||||
// first instant. A test sets it to nothing.
|
||||
var handoverSettle = 10 * time.Second
|
||||
|
||||
// stillUp says why a process's unit is not up and staying up, or nothing when it is: active and
|
||||
// running at two looks handoverSettle apart, the same main process both times, restarted by
|
||||
// nothing in between. Stricter than stayedRunning, which reads "activating" as running — the state
|
||||
// a crash-looping unit is in while it waits to be started again, which is exactly the replacement
|
||||
// that must not be handed anything.
|
||||
func stillUp(ctx context.Context, run Runner, unit string) string {
|
||||
look := func() (map[string]string, string) {
|
||||
out, err := run(ctx, "systemctl", "show", unit, "--property=ActiveState", "--property=SubState",
|
||||
"--property=MainPID", "--property=NRestarts")
|
||||
if err != nil {
|
||||
return nil, err.Error()
|
||||
}
|
||||
got := map[string]string{}
|
||||
for _, line := range strings.Split(out, "\n") {
|
||||
if key, value, found := strings.Cut(strings.TrimSpace(line), "="); found {
|
||||
got[key] = value
|
||||
}
|
||||
}
|
||||
if got["ActiveState"] != "active" || got["SubState"] != "running" {
|
||||
return got, fmt.Sprintf("%s/%s", got["ActiveState"], got["SubState"])
|
||||
}
|
||||
return got, ""
|
||||
}
|
||||
first, why := look()
|
||||
if why != "" {
|
||||
return why
|
||||
}
|
||||
timer := time.NewTimer(handoverSettle)
|
||||
defer timer.Stop()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err().Error()
|
||||
case <-timer.C:
|
||||
}
|
||||
second, why := look()
|
||||
if why != "" {
|
||||
return why
|
||||
}
|
||||
if first["MainPID"] != second["MainPID"] || first["NRestarts"] != second["NRestarts"] {
|
||||
return fmt.Sprintf("restarted while it was watched (pid %s → %s, restarts %s → %s)",
|
||||
first["MainPID"], second["MainPID"], first["NRestarts"], second["NRestarts"])
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// keptForReplacement reports an orphan kept because what replaces it is not yet in its place.
|
||||
func keptForReplacement(report *Report, log func(string), orphan store.Applied, detail string) {
|
||||
report.Outcomes = append(report.Outcomes, Outcome{
|
||||
ID: orphan.ID, Type: orphan.Type, Target: orphan.Target, Action: "kept", Detail: detail,
|
||||
})
|
||||
log(fmt.Sprintf(" kept %s (%s): %s", orphan.ID, orphan.Target, strings.TrimPrefix(detail, "kept: ")))
|
||||
}
|
||||
|
||||
func appliedIDs(applied []store.Applied) string {
|
||||
ids := make([]string, 0, len(applied))
|
||||
for _, a := range applied {
|
||||
ids = append(ids, a.ID)
|
||||
}
|
||||
return strings.Join(ids, ", ")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,318 @@
|
||||
package apply
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-host/internal/store"
|
||||
)
|
||||
|
||||
// novox/hq issue 213: the controller moves from a container to a process on the one machine that
|
||||
// runs it. Every orphan is removed before anything is applied, so without a handover the container
|
||||
// went first and nothing answered the mesh's verbs while the process was fetched, unpacked and
|
||||
// started — and for ever, if it did not start. A process that `replaces` the container is applied
|
||||
// first; the container goes only once the process is running a moment later.
|
||||
|
||||
// aMachine fakes the service manager and the container runtime: it records every command, answers
|
||||
// `systemctl show` with whether the process is running, and has a container until it is removed.
|
||||
type aMachine struct {
|
||||
commands []string
|
||||
running bool // what `systemctl show` says of the process once it was started
|
||||
started bool
|
||||
container bool
|
||||
timer bool // whether a timer, once started, is up
|
||||
crashing bool // up at the first look, waiting to restart at the next
|
||||
looks int
|
||||
}
|
||||
|
||||
func (m *aMachine) run(ctx context.Context, name string, args ...string) (string, error) {
|
||||
line := name + " " + strings.Join(args, " ")
|
||||
m.commands = append(m.commands, line)
|
||||
switch {
|
||||
case name == "systemctl" && len(args) > 0 && (args[0] == "restart" || args[0] == "start"):
|
||||
m.started = true
|
||||
case name == "systemctl" && len(args) > 0 && args[0] == "is-active" && strings.HasSuffix(args[len(args)-1], ".timer"):
|
||||
if m.timer {
|
||||
return "active", nil
|
||||
}
|
||||
return "inactive", errors.New("inactive")
|
||||
case name == "systemctl" && len(args) > 0 && args[0] == "is-active":
|
||||
if m.started && m.running {
|
||||
return "active", nil
|
||||
}
|
||||
return "inactive", errors.New("inactive")
|
||||
case name == "systemctl" && len(args) > 0 && args[0] == "show":
|
||||
m.looks++
|
||||
if m.crashing && m.started {
|
||||
// Up at the first look; waiting to be started again, a new process, at the second.
|
||||
if m.looks == 1 {
|
||||
return "ActiveState=active\nSubState=running\nMainPID=42\nNRestarts=0\n", nil
|
||||
}
|
||||
return "ActiveState=activating\nSubState=auto-restart\nMainPID=0\nNRestarts=1\n", nil
|
||||
}
|
||||
if m.started && m.running {
|
||||
return "ActiveState=active\nSubState=running\nMainPID=42\nNRestarts=0\n", nil
|
||||
}
|
||||
return "ActiveState=inactive\nSubState=dead\nMainPID=0\nNRestarts=0\n", nil
|
||||
case name == "docker" && len(args) > 1 && args[0] == "rm":
|
||||
m.container = false
|
||||
case name == "docker" && len(args) > 1 && args[0] == "container" && args[1] == "inspect":
|
||||
if !m.container {
|
||||
return "", errors.New("no such container")
|
||||
}
|
||||
return "true\t", nil
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
|
||||
func (m *aMachine) index(prefix string) int {
|
||||
for i, c := range m.commands {
|
||||
if strings.HasPrefix(c, prefix) {
|
||||
return i
|
||||
}
|
||||
}
|
||||
return -1
|
||||
}
|
||||
|
||||
// theController is a machine whose controller ran as a container, recorded, and a declaration that
|
||||
// runs it as a process from a bundle served here instead.
|
||||
func theController(t *testing.T, digest, source, extra string) (store.State, string) {
|
||||
t.Helper()
|
||||
known := store.State{Resources: []store.Applied{{
|
||||
Origin: store.OriginDeclared, ID: "mesh-controller.server", Type: "container", Target: "mesh-controller",
|
||||
}}}
|
||||
return known, `{"declaration":1,"resources":[
|
||||
{"id":"mesh-controller.controller","type":"process","name":"mesh-controller","source":"` + source +
|
||||
`","digest":"` + digest + `","run":["./mesh-controller","serve"]` + extra + `}]}`
|
||||
}
|
||||
|
||||
func onAMachine(t *testing.T) {
|
||||
t.Helper()
|
||||
serviceSettle, handoverSettle = 0, 0
|
||||
wasUnits, wasBundles := unitDir, daemonRoot
|
||||
unitDir, daemonRoot = t.TempDir(), t.TempDir()
|
||||
t.Cleanup(func() { unitDir, daemonRoot = wasUnits, wasBundles })
|
||||
}
|
||||
|
||||
func TestAContainerAProcessReplacesGoesOnlyOnceTheProcessRuns(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
known, raw := theController(t, digest, serving(t, body), `,"replaces":["mesh-controller.server"]`)
|
||||
m := &aMachine{running: true, container: true}
|
||||
|
||||
report, after, err := Apply(context.Background(), archHost(t), parse(t, raw), known,
|
||||
store.OriginDeclared, m.run, nil, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("the handover failed: %v", err)
|
||||
}
|
||||
started, removed := m.index("systemctl restart mesh-controller.service"), m.index("docker rm -f mesh-controller")
|
||||
if started < 0 || removed < 0 {
|
||||
t.Fatalf("the process was not started or the container not removed: %v", m.commands)
|
||||
}
|
||||
if removed < started {
|
||||
t.Fatalf("the container was removed before its replacement was started — a window with "+
|
||||
"nothing answering: %v", m.commands)
|
||||
}
|
||||
if looked := m.index("systemctl show mesh-controller.service"); looked < 0 || looked > removed {
|
||||
t.Errorf("the container was removed without looking whether the process runs: %v", m.commands)
|
||||
}
|
||||
if o := outcomeOf(report, "mesh-controller.server"); o.Action != "removed" {
|
||||
t.Errorf("the container's outcome is %+v, want removed", o)
|
||||
}
|
||||
if _, still := after.Find("mesh-controller.server"); still {
|
||||
t.Error("the host still records the container it removed")
|
||||
}
|
||||
if _, has := after.Find("mesh-controller.controller"); !has {
|
||||
t.Error("the process was not recorded")
|
||||
}
|
||||
}
|
||||
|
||||
// The case the handover exists for: the replacement does not stay up. The container keeps
|
||||
// answering, stays recorded so a later apply hands it over, and the apply says why it failed.
|
||||
func TestAContainerIsKeptWhenItsReplacementDoesNotRun(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
known, raw := theController(t, digest, serving(t, body), `,"replaces":["mesh-controller.server"]`)
|
||||
m := &aMachine{running: false, container: true}
|
||||
|
||||
report, after, err := Apply(context.Background(), archHost(t), parse(t, raw), known,
|
||||
store.OriginDeclared, m.run, nil, nil)
|
||||
if err == nil {
|
||||
t.Fatal("a replacement that is not running was reported as a clean apply")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "mesh-controller.server") {
|
||||
t.Errorf("the failure does not name what was kept: %v", err)
|
||||
}
|
||||
if m.index("docker rm") >= 0 {
|
||||
t.Fatalf("the container was removed though its replacement is not running: %v", m.commands)
|
||||
}
|
||||
if o := outcomeOf(report, "mesh-controller.server"); o.Action != "kept" {
|
||||
t.Errorf("the container's outcome is %+v, want kept", o)
|
||||
}
|
||||
if _, still := after.Find("mesh-controller.server"); !still {
|
||||
t.Fatal("the container was forgotten, so no later apply would ever remove it")
|
||||
}
|
||||
|
||||
// The next apply finds the process up and finishes the handover.
|
||||
m.running, m.commands = true, nil
|
||||
report, after, err = Apply(context.Background(), archHost(t), parse(t, raw), after,
|
||||
store.OriginDeclared, m.run, nil, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("the second apply failed: %v", err)
|
||||
}
|
||||
if o := outcomeOf(report, "mesh-controller.server"); o.Action != "removed" {
|
||||
t.Errorf("the second apply did not hand over: %+v (%v)", o, m.commands)
|
||||
}
|
||||
if _, still := after.Find("mesh-controller.server"); still {
|
||||
t.Error("the container is still recorded after the handover")
|
||||
}
|
||||
}
|
||||
|
||||
// A replacement that never applied — its bundle is not what was declared — touches nothing, and
|
||||
// what it replaces keeps running.
|
||||
func TestAContainerIsKeptWhenItsReplacementFailsToApply(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, _ := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
known, raw := theController(t, "sha256:"+strings.Repeat("b", 64), serving(t, body),
|
||||
`,"replaces":["mesh-controller.server"]`)
|
||||
m := &aMachine{running: true, container: true}
|
||||
|
||||
report, after, err := Apply(context.Background(), archHost(t), parse(t, raw), known,
|
||||
store.OriginDeclared, m.run, nil, nil)
|
||||
if err == nil {
|
||||
t.Fatal("a replacement whose bundle did not match was reported applied")
|
||||
}
|
||||
if m.index("docker rm") >= 0 {
|
||||
t.Fatalf("the container was removed though nothing replaced it: %v", m.commands)
|
||||
}
|
||||
if o := outcomeOf(report, "mesh-controller.server"); o.Action != "kept" {
|
||||
t.Errorf("the container's outcome is %+v, want kept", o)
|
||||
}
|
||||
if _, still := after.Find("mesh-controller.server"); !still {
|
||||
t.Fatal("the container was forgotten")
|
||||
}
|
||||
}
|
||||
|
||||
// Without `replaces` nothing changes: an orphan goes before anything is applied, as it always has.
|
||||
// Kept as a test because it is the window the field exists to close.
|
||||
func TestWithoutReplacesAnOrphanStillGoesFirst(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
known, raw := theController(t, digest, serving(t, body), ``)
|
||||
m := &aMachine{running: true, container: true}
|
||||
if _, _, err := Apply(context.Background(), archHost(t), parse(t, raw), known,
|
||||
store.OriginDeclared, m.run, nil, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if removed, started := m.index("docker rm -f mesh-controller"), m.index("systemctl restart mesh-controller.service"); removed < 0 || removed > started {
|
||||
t.Fatalf("an orphan nothing replaces was not removed first: %v", m.commands)
|
||||
}
|
||||
}
|
||||
|
||||
// novox/hq issue 213, beyond what the oneshot unit (process_step_test.go) already holds: a step
|
||||
// written `./name` runs its own bundle's binary — the controller's preparation is its own binary —
|
||||
// and a step is started, never enabled.
|
||||
func TestAStepRunsItsOwnBundlesBinaryAndIsNotEnabled(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
m := &aMachine{running: true}
|
||||
d := declare(t, `{"id":"mesh-controller.controller-prepare","type":"process","name":"mesh-controller-prepare",
|
||||
"source":"`+serving(t, body)+`","digest":"`+digest+`","run":["./mesh-controller","prepare"],"run-once":true}`)
|
||||
if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginDeclared, m.run, nil, nil); err != nil {
|
||||
t.Fatalf("the step failed: %v", err)
|
||||
}
|
||||
for _, c := range m.commands {
|
||||
if strings.HasPrefix(c, "./") || strings.HasPrefix(c, "systemctl enable") {
|
||||
t.Errorf("the step was run directly or enabled: %v", m.commands)
|
||||
}
|
||||
}
|
||||
unit, err := os.ReadFile(filepath.Join(unitDir, "mesh-controller-prepare.service"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
want := "ExecStart=" + filepath.Join(daemonRoot, "mesh-controller-prepare", "mesh-controller") + " prepare"
|
||||
if !strings.Contains(string(unit), want) {
|
||||
t.Errorf("the step does not run its own bundle's binary (%q):\n%s", want, unit)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAStepThatFailsGatesItsModule(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
src := serving(t, body)
|
||||
run := func(ctx context.Context, name string, args ...string) (string, error) {
|
||||
if name == "systemctl" && len(args) > 1 && args[0] == "start" {
|
||||
return "", errors.New("Job for mesh-controller-prepare.service failed")
|
||||
}
|
||||
if name == "systemctl" && len(args) > 0 && args[0] == "restart" {
|
||||
t.Errorf("the module's process was started after its step failed")
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
d := declare(t, `{"id":"mesh-controller.controller-prepare","type":"process","name":"mesh-controller-prepare",
|
||||
"source":"`+src+`","digest":"`+digest+`","run":["./mesh-controller","prepare"],"run-once":true},
|
||||
{"id":"mesh-controller.controller","type":"process","name":"mesh-controller",
|
||||
"source":"`+src+`","digest":"`+digest+`","run":["./mesh-controller","serve"]}`)
|
||||
report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginDeclared, run, nil, nil)
|
||||
if err == nil {
|
||||
t.Fatal("a failed step was reported as a clean apply")
|
||||
}
|
||||
if o := outcomeOf(report, "mesh-controller.controller"); o.Action != "skipped" {
|
||||
t.Errorf("the process after a failed step was %+v, want skipped", o)
|
||||
}
|
||||
}
|
||||
|
||||
// A completed step is done: applied again unchanged, it is not run again — its service is never up
|
||||
// between runs, and reading that as "a daemon that stopped" re-ran the controller's preparation on
|
||||
// every apply. Nor is a scheduled run started off its cadence; its timer is what is kept up.
|
||||
func TestACompletedStepIsNotRunAgainAndAScheduleIsItsTimer(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"job": "#!/bin/sh\n"})
|
||||
src := serving(t, body)
|
||||
for _, mode := range []string{`"run-once":true`, `"schedule":"0 3 * * *"`} {
|
||||
// The service of either is never up between runs; a scheduled one's timer is.
|
||||
m := &aMachine{running: false, timer: true}
|
||||
d := declare(t, `{"id":"m.job","type":"process","name":"m-job","source":"`+src+`","digest":"`+digest+
|
||||
`","run":["./job"],`+mode+`}`)
|
||||
_, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginDeclared, m.run, nil, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("%s: first apply: %v", mode, err)
|
||||
}
|
||||
m.commands = nil
|
||||
report, _, err := Apply(context.Background(), archHost(t), d, state, store.OriginDeclared, m.run, nil, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("%s: second apply: %v", mode, err)
|
||||
}
|
||||
if m.index("systemctl start m-job.service") >= 0 || m.index("systemctl restart m-job.service") >= 0 {
|
||||
t.Errorf("%s: an unchanged apply ran the job again: %v", mode, m.commands)
|
||||
}
|
||||
if o := outcomeOf(report, "m.job"); o.Action != "unchanged" {
|
||||
t.Errorf("%s: an unchanged apply reported %+v", mode, o)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A replacement crash-looping between its restarts is not up, though the service manager calls it
|
||||
// "activating" — the state the host's ordinary second look accepts as running. Removing the
|
||||
// container on that reading would leave nothing answering.
|
||||
func TestAContainerIsKeptWhenItsReplacementIsCrashLooping(t *testing.T) {
|
||||
onAMachine(t)
|
||||
body, digest := anArchive(t, map[string]string{"mesh-controller": "#!/bin/sh\n"})
|
||||
known, raw := theController(t, digest, serving(t, body), `,"replaces":["mesh-controller.server"]`)
|
||||
m := &aMachine{crashing: true, container: true}
|
||||
_, after, err := Apply(context.Background(), archHost(t), parse(t, raw), known,
|
||||
store.OriginDeclared, m.run, nil, nil)
|
||||
if err == nil || !strings.Contains(err.Error(), "auto-restart") {
|
||||
t.Fatalf("a crash-looping replacement was accepted: %v", err)
|
||||
}
|
||||
if m.index("docker rm") >= 0 {
|
||||
t.Fatalf("the container was removed for a replacement that keeps dying: %v", m.commands)
|
||||
}
|
||||
if _, still := after.Find("mesh-controller.server"); !still {
|
||||
t.Fatal("the container was forgotten")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,238 @@
|
||||
package apply
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-host/internal/declaration"
|
||||
)
|
||||
|
||||
// A scheduled step may hold its module's own containers still while it runs (novox/hq ADR 0189,
|
||||
// issue 108).
|
||||
//
|
||||
// What it exists for: the artifact store's collector walks the storage and requires every writer
|
||||
// stopped. A run-once step runs beside containers and a scheduled one is the same container again,
|
||||
// so the mesh had no way to say it — which is why the store it inherited has never collected
|
||||
// anything. The risk the field brings is one shape only: a window that opens and never closes.
|
||||
// Every test here is about that shape.
|
||||
|
||||
// windowRun records the order of stop / run / start, which is the whole of what is being asserted.
|
||||
type windowRun struct {
|
||||
mu sync.Mutex
|
||||
order []string
|
||||
failAt string // the arg[0] that should fail ("run" makes the step fail)
|
||||
wontGo string // a container name that refuses to start again
|
||||
}
|
||||
|
||||
func (w *windowRun) run(_ context.Context, _ string, args ...string) (string, error) {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
switch args[0] {
|
||||
case "info":
|
||||
return "27.0\n", nil
|
||||
case "stop", "start":
|
||||
w.order = append(w.order, args[0]+" "+args[1])
|
||||
if args[0] == "start" && args[1] == w.wontGo {
|
||||
return "", errors.New("the runtime refused")
|
||||
}
|
||||
case "run":
|
||||
w.order = append(w.order, "run")
|
||||
if w.failAt == "run" {
|
||||
return "", errors.New("the step exited non-zero")
|
||||
}
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
|
||||
func (w *windowRun) seen() []string {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
return append([]string{}, w.order...)
|
||||
}
|
||||
|
||||
// aStoreWithACollector is a module in the shape distribution has: a server that must not be
|
||||
// writing, and a nightly step that walks its storage with the server held still.
|
||||
func aStoreWithACollector(t *testing.T) *declaration.Declaration {
|
||||
t.Helper()
|
||||
return parseTrusted(t, `{"declaration":1,"resources":[
|
||||
{"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"},
|
||||
{"id":"collect","type":"container","name":"mesh-registry-collect","image":"`+pinned+`",
|
||||
"schedule":"30 3 * * *","while-stopped":["store"]}
|
||||
]}`)
|
||||
}
|
||||
|
||||
func fireOnce(t *testing.T, d *declaration.Declaration, w *windowRun) {
|
||||
t.Helper()
|
||||
clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)}
|
||||
s := NewScheduler(clock, w.run, func(string) {})
|
||||
s.Sync(d, nil)
|
||||
s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC))
|
||||
s.Wait()
|
||||
}
|
||||
|
||||
func TestAScheduledStepHoldsItsModulesContainerStillAndStartsItAgain(t *testing.T) {
|
||||
w := &windowRun{}
|
||||
fireOnce(t, aStoreWithACollector(t), w)
|
||||
|
||||
got := w.seen()
|
||||
want := []string{"stop mesh-registry", "run", "start mesh-registry"}
|
||||
var kept []string
|
||||
for _, line := range got {
|
||||
if strings.HasPrefix(line, "stop mesh-registry-collect") {
|
||||
// Clearing the step's own exited container by name; not part of the window.
|
||||
continue
|
||||
}
|
||||
kept = append(kept, line)
|
||||
}
|
||||
if len(kept) != len(want) {
|
||||
t.Fatalf("the window was not stop, run, start: %v", got)
|
||||
}
|
||||
for i := range want {
|
||||
if kept[i] != want[i] {
|
||||
t.Fatalf("the window was %v, want %v", kept, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The one that matters: a step that fails must leave the service running.
|
||||
func TestAFailedStepStillClosesTheWindow(t *testing.T) {
|
||||
w := &windowRun{failAt: "run"}
|
||||
fireOnce(t, aStoreWithACollector(t), w)
|
||||
|
||||
var started bool
|
||||
for _, line := range w.seen() {
|
||||
if line == "start mesh-registry" {
|
||||
started = true
|
||||
}
|
||||
}
|
||||
if !started {
|
||||
t.Fatalf("the step failed and the container it held still was never started again: %v", w.seen())
|
||||
}
|
||||
}
|
||||
|
||||
// A container that will not come back is said loudly: it is down, and nothing else notices until
|
||||
// the next apply compares it.
|
||||
func TestAContainerThatWillNotStartAgainIsSaidLoudly(t *testing.T) {
|
||||
w := &windowRun{wontGo: "mesh-registry"}
|
||||
var said []string
|
||||
clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)}
|
||||
s := NewScheduler(clock, w.run, func(line string) { said = append(said, line) })
|
||||
s.Sync(aStoreWithACollector(t), nil)
|
||||
s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC))
|
||||
s.Wait()
|
||||
|
||||
var loud bool
|
||||
for _, line := range said {
|
||||
if strings.Contains(line, "WILL NOT START AGAIN") && strings.Contains(line, "mesh-registry") {
|
||||
loud = true
|
||||
}
|
||||
}
|
||||
if !loud {
|
||||
t.Fatalf("a service left stopped by a maintenance window was not said loudly: %v", said)
|
||||
}
|
||||
}
|
||||
|
||||
// Several containers come back in the reverse of the order they were stopped: a module names the
|
||||
// dependant first, and starting it before what it depends on is not bringing it back.
|
||||
func TestTheWindowClosesInTheReverseOfTheOrderItOpened(t *testing.T) {
|
||||
d := parseTrusted(t, `{"declaration":1,"resources":[
|
||||
{"id":"web","type":"container","name":"web","image":"`+pinned+`"},
|
||||
{"id":"db","type":"container","name":"db","image":"`+pinned+`"},
|
||||
{"id":"collect","type":"container","name":"collect","image":"`+pinned+`",
|
||||
"schedule":"30 3 * * *","while-stopped":["web","db"]}
|
||||
]}`)
|
||||
w := &windowRun{}
|
||||
fireOnce(t, d, w)
|
||||
|
||||
var stops, starts []string
|
||||
for _, line := range w.seen() {
|
||||
switch {
|
||||
case line == "stop web" || line == "stop db":
|
||||
stops = append(stops, line)
|
||||
case strings.HasPrefix(line, "start "):
|
||||
starts = append(starts, line)
|
||||
}
|
||||
}
|
||||
if len(stops) != 2 || stops[0] != "stop web" || stops[1] != "stop db" {
|
||||
t.Fatalf("stopped in %v, want the order the step named them", stops)
|
||||
}
|
||||
if len(starts) != 2 || starts[0] != "start db" || starts[1] != "start web" {
|
||||
t.Fatalf("started in %v, want the reverse", starts)
|
||||
}
|
||||
}
|
||||
|
||||
// And the refusals, each for what it says rather than that it says something.
|
||||
func TestAMaintenanceWindowIsRefusedWhereItCannotMean(t *testing.T) {
|
||||
for _, c := range []struct{ name, body, says string }{
|
||||
{
|
||||
"a window with no schedule",
|
||||
`{"id":"collect","type":"container","name":"c","image":"` + pinned + `","while-stopped":["store"]}`,
|
||||
"needs a schedule",
|
||||
},
|
||||
{
|
||||
"a window naming itself",
|
||||
`{"id":"collect","type":"container","name":"c","image":"` + pinned + `","schedule":"30 3 * * *","while-stopped":["collect"]}`,
|
||||
"this step itself",
|
||||
},
|
||||
{
|
||||
"a window naming something that is not a container here",
|
||||
`{"id":"collect","type":"container","name":"c","image":"` + pinned + `","schedule":"30 3 * * *","while-stopped":["elsewhere"]}`,
|
||||
"no container by that id",
|
||||
},
|
||||
} {
|
||||
_, err := declaration.ParseTrusted([]byte(`{"declaration":1,"resources":[` + c.body + `]}`))
|
||||
if err == nil {
|
||||
t.Errorf("%s was accepted", c.name)
|
||||
continue
|
||||
}
|
||||
if !strings.Contains(err.Error(), c.says) {
|
||||
t.Errorf("%s: the refusal does not say %q: %v", c.name, c.says, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A changed window is a changed declaration, and the install says so.
|
||||
//
|
||||
// The cadence already works this way: "a changed schedule is a changed spec — the marker moves and
|
||||
// the install is reported updated and re-established" (containerSpec). Which containers are held
|
||||
// still for the run is the same kind of statement, and a declaration that changed it while the
|
||||
// machine reported no change would be a machine quietly running the old window.
|
||||
func TestAChangedWindowMovesTheSpec(t *testing.T) {
|
||||
one := parseTrusted(t, `{"declaration":1,"resources":[
|
||||
{"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"},
|
||||
{"id":"other","type":"container","name":"other","image":"`+pinned+`"},
|
||||
{"id":"collect","type":"container","name":"collect","image":"`+pinned+`",
|
||||
"schedule":"30 3 * * *","while-stopped":["store"]}
|
||||
]}`)
|
||||
two := parseTrusted(t, `{"declaration":1,"resources":[
|
||||
{"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"},
|
||||
{"id":"other","type":"container","name":"other","image":"`+pinned+`"},
|
||||
{"id":"collect","type":"container","name":"collect","image":"`+pinned+`",
|
||||
"schedule":"30 3 * * *","while-stopped":["store","other"]}
|
||||
]}`)
|
||||
stepOf := func(d *declaration.Declaration) *declaration.Container {
|
||||
for _, r := range d.Resources {
|
||||
if c, ok := r.(*declaration.Container); ok && c.ID == "collect" {
|
||||
return c
|
||||
}
|
||||
}
|
||||
t.Fatal("no step in the fixture")
|
||||
return nil
|
||||
}
|
||||
if containerSpec(stepOf(one), inputs{}) == containerSpec(stepOf(two), inputs{}) {
|
||||
t.Fatal("the window changed and the spec did not; the machine would report no change " +
|
||||
"and keep holding the containers it held yesterday")
|
||||
}
|
||||
// And a container with no window is untouched by the field existing at all.
|
||||
plain := parseTrusted(t, `{"declaration":1,"resources":[
|
||||
{"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"}
|
||||
]}`)
|
||||
spec := containerSpec(plain.Resources[0].(*declaration.Container), inputs{})
|
||||
if strings.Contains(spec, "while-stopped") || strings.Contains(spec, "held") {
|
||||
t.Errorf("an ordinary container's spec mentions a field it does not set:\n%s", spec)
|
||||
}
|
||||
}
|
||||
@@ -82,11 +82,21 @@ func applyProcess(ctx context.Context, r *declaration.Process, run Runner,
|
||||
// next apply then found no record, re-created the daemon, and the one after that found a
|
||||
// record again — the node's runtime restarted every other cycle (novox/hq 04-ISSUES/210).
|
||||
out.wrote = want
|
||||
if active, err := run(ctx, "systemctl", "is-active", "--quiet", r.Name+".service"); err == nil {
|
||||
// **A step that ran is done, and a schedule is its timer** (novox/hq issue 213). Neither's
|
||||
// service is meant to be up between runs, so asking whether it is — and starting it when it
|
||||
// was not — ran a completed step again on every apply, and a scheduled run off its cadence.
|
||||
if r.RunOnce {
|
||||
return out, nil
|
||||
}
|
||||
unit := r.Name + ".service"
|
||||
if r.Schedule != "" {
|
||||
unit = r.Name + ".timer"
|
||||
}
|
||||
if active, err := run(ctx, "systemctl", "is-active", "--quiet", unit); err == nil {
|
||||
_ = active
|
||||
return out, nil
|
||||
}
|
||||
if _, err := run(ctx, "systemctl", "start", r.Name+".service"); err != nil {
|
||||
if _, err := run(ctx, "systemctl", "start", unit); err != nil {
|
||||
return out, fmt.Errorf("%s is installed and would not start: %w", r.Name, err)
|
||||
}
|
||||
out.Action = "updated"
|
||||
@@ -115,9 +125,22 @@ func applyProcess(ctx context.Context, r *declaration.Process, run Runner,
|
||||
// so the machine is not asked to start something that needed a migration that did not happen.
|
||||
// Nothing is left behind to ask afterwards: the record that it ran is the digest, which is why
|
||||
// the identity above includes the command.
|
||||
//
|
||||
// **Run as the unit a daemon would be, once** (novox/hq design 38 WP4c). Run directly, the
|
||||
// step started in the host's own working directory, without its environment, its environment
|
||||
// files or its user — `node bootstrap/index.js` resolved from wherever the host ran and was told
|
||||
// none of the words it was declared with. A oneshot unit carries all four exactly as a daemon's
|
||||
// does, and starting one waits for it to finish and fails when it fails.
|
||||
if r.RunOnce {
|
||||
if _, err := run(ctx, r.Run[0], r.Run[1:]...); err != nil {
|
||||
return out, fmt.Errorf("the %s step did not complete: %w", r.Name, err)
|
||||
unit := filepath.Join(unitDir, r.Name+".service")
|
||||
if err := os.WriteFile(unit, []byte(unitFor(r)), 0o644); err != nil {
|
||||
return out, err
|
||||
}
|
||||
if _, err := run(ctx, "systemctl", "daemon-reload"); err != nil {
|
||||
return out, err
|
||||
}
|
||||
if _, err := run(ctx, "systemctl", "start", r.Name+".service"); err != nil {
|
||||
return out, fmt.Errorf("the %s step did not complete (journalctl -u %s.service says why): %w", r.Name, r.Name, err)
|
||||
}
|
||||
out.Action = "created"
|
||||
if previous.Wrote != "" {
|
||||
@@ -215,8 +238,8 @@ func unitFor(r *declaration.Process) string {
|
||||
fmt.Fprintf(&b, "User=%s\n", r.User)
|
||||
}
|
||||
fmt.Fprintf(&b, "ExecStart=%s\n", strings.Join(runFrom(r), " "))
|
||||
if r.Schedule != "" {
|
||||
// Started by its timer and expected to finish. Restarting it would have it run
|
||||
if r.Schedule != "" || r.RunOnce {
|
||||
// Started by its timer, or once by the host, and expected to finish. Restarting it would have it run
|
||||
// continuously between fires, which is the opposite of a schedule.
|
||||
b.WriteString("Type=oneshot\n")
|
||||
b.WriteString("\n")
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
package apply
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"os/user"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-host/internal/store"
|
||||
)
|
||||
|
||||
// A run-once process is a step run where, how and as whom it was declared (novox/hq design 38
|
||||
// WP4c): its bundle's directory, its environment and environment files, its user. Run directly,
|
||||
// it started in the host's own directory with none of them.
|
||||
func TestARunOnceProcessRunsAsItsOneshotUnit(t *testing.T) {
|
||||
units, bundles := t.TempDir(), t.TempDir()
|
||||
wasUnits, wasBundles := unitDir, daemonRoot
|
||||
unitDir, daemonRoot = units, bundles
|
||||
t.Cleanup(func() { unitDir, daemonRoot = wasUnits, wasBundles })
|
||||
|
||||
me, err := user.Current()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
body, digest := anArchive(t, map[string]string{"bootstrap/index.js": "console.log(1)\n"})
|
||||
var commands []string
|
||||
fail := false
|
||||
run := func(ctx context.Context, name string, args ...string) (string, error) {
|
||||
commands = append(commands, name+" "+strings.Join(args, " "))
|
||||
if fail && name == "systemctl" && len(args) > 0 && args[0] == "start" {
|
||||
return "", errors.New("exit status 1")
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
d := declare(t, `{"id":"mosquitto.bootstrap","type":"process","name":"mosquitto-bootstrap","source":"`+serving(t, body)+
|
||||
`","digest":"`+digest+`","run":["node","bootstrap/index.js"],"run-once":true,"user":"`+me.Username+`",`+
|
||||
`"env":{"MESH_ADMIN":"mesh-admin"},"env-file":["/var/lib/mesh/mosquitto/bootstrap.env"]}`)
|
||||
|
||||
if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginDeclared, run, nil, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
unit, err := os.ReadFile(filepath.Join(units, "mosquitto-bootstrap.service"))
|
||||
if err != nil {
|
||||
t.Fatalf("no unit was written for the step: %v", err)
|
||||
}
|
||||
for _, want := range []string{
|
||||
"WorkingDirectory=" + filepath.Join(bundles, "mosquitto-bootstrap"),
|
||||
"EnvironmentFile=/var/lib/mesh/mosquitto/bootstrap.env",
|
||||
`Environment="MESH_ADMIN=mesh-admin"`,
|
||||
"User=" + me.Username,
|
||||
"Type=oneshot",
|
||||
"ExecStart=node bootstrap/index.js",
|
||||
} {
|
||||
if !strings.Contains(string(unit), want) {
|
||||
t.Errorf("the step's unit lacks %q:\n%s", want, unit)
|
||||
}
|
||||
}
|
||||
for _, never := range []string{"Restart=always", "[Install]", "Type=simple"} {
|
||||
if strings.Contains(string(unit), never) {
|
||||
t.Errorf("a step's unit says %q:\n%s", never, unit)
|
||||
}
|
||||
}
|
||||
joined := strings.Join(commands, "; ")
|
||||
if !strings.Contains(joined, "systemctl start mosquitto-bootstrap.service") {
|
||||
t.Errorf("the step was not started as its unit: %s", joined)
|
||||
}
|
||||
if strings.Contains(joined, "node bootstrap/index.js") || strings.Contains(joined, "enable mosquitto-bootstrap") {
|
||||
t.Errorf("the step was run directly or enabled: %s", joined)
|
||||
}
|
||||
|
||||
// A step that fails fails the apply, and is not recorded as done.
|
||||
fail = true
|
||||
d2 := declare(t, `{"id":"mosquitto.bootstrap","type":"process","name":"mosquitto-bootstrap","source":"`+serving(t, body)+
|
||||
`","digest":"`+digest+`","run":["node","bootstrap/index.js"],"run-once":true,"env":{"MESH_ADMIN":"changed"}}`)
|
||||
if _, _, err := Apply(context.Background(), archHost(t), d2, store.State{}, store.OriginDeclared, run, nil, nil); err == nil {
|
||||
t.Error("a step that failed did not fail the apply")
|
||||
}
|
||||
}
|
||||
@@ -65,6 +65,29 @@ type scheduledJob struct {
|
||||
container *declaration.Container
|
||||
next time.Time // the next minute at which it is due
|
||||
running bool // a run is in flight — the next due run is skipped rather than stacked
|
||||
// hold is the runtime names of the containers held still for the duration of a run, in the
|
||||
// order the step named them (novox/hq ADR 0189).
|
||||
hold []string
|
||||
}
|
||||
|
||||
// heldStillFor is the runtime names of the containers a step holds still, resolved from ids.
|
||||
func heldStillFor(step *declaration.Container, d *declaration.Declaration) []string {
|
||||
if len(step.WhileStopped) == 0 {
|
||||
return nil
|
||||
}
|
||||
byID := map[string]string{}
|
||||
for _, r := range d.Resources {
|
||||
if c, ok := r.(*declaration.Container); ok {
|
||||
byID[c.Identity()] = c.Name
|
||||
}
|
||||
}
|
||||
out := make([]string, 0, len(step.WhileStopped))
|
||||
for _, id := range step.WhileStopped {
|
||||
if name := byID[id]; name != "" {
|
||||
out = append(out, name)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// NewScheduler builds a scheduler. A nil clock is the system clock; a nil log says nothing.
|
||||
@@ -123,15 +146,20 @@ func (s *Scheduler) Sync(d *declaration.Declaration, held map[string]bool) {
|
||||
// on its cadence rather than staying running to be restarted — so its identity cannot
|
||||
// depend on another resource's content and there is nothing to pass.
|
||||
spec := containerSpec(c, inputs{})
|
||||
// The runtime stops containers by name; the declaration names them by id. Resolved here,
|
||||
// against the declaration this job was armed from, so a fire never has to look anything up
|
||||
// (novox/hq ADR 0189). The parser has already refused an id that is not a container here.
|
||||
hold := heldStillFor(c, d)
|
||||
if existing := s.jobs[c.Identity()]; existing != nil && existing.spec == spec {
|
||||
// Unchanged: keep where it is in its cadence, refresh the declaration pointer only.
|
||||
existing.container = c
|
||||
existing.hold = hold
|
||||
continue
|
||||
}
|
||||
// New or changed: arm it for the next due minute after now.
|
||||
next, _ := cron.Next(s.clock.Now())
|
||||
s.jobs[c.Identity()] = &scheduledJob{
|
||||
id: c.Identity(), spec: spec, cron: cron, container: c, next: next,
|
||||
id: c.Identity(), spec: spec, cron: cron, container: c, next: next, hold: hold,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -202,6 +230,15 @@ func (s *Scheduler) fire(ctx context.Context, j *scheduledJob) {
|
||||
return
|
||||
}
|
||||
|
||||
// **The window opens here and closes in the defer, whatever happens** (novox/hq ADR 0189).
|
||||
// Deferred before the first stop so a panic, a failing step or a step that runs long all end
|
||||
// the same way: the service running. The one real risk of this field is a window that never
|
||||
// closes, and the only defence against it is that closing is not conditional on anything.
|
||||
if len(j.hold) > 0 {
|
||||
defer s.letRun(ctx, cri, j)
|
||||
s.holdStill(ctx, cri, j)
|
||||
}
|
||||
|
||||
// A container by this name left exited by the previous run would collide with --name. Removing
|
||||
// one that is not there is the state we want, so its error is ignored — the same as run-once.
|
||||
_, _ = s.run(ctx, cri, "rm", "-f", j.container.Name)
|
||||
@@ -260,3 +297,51 @@ func (s *Scheduler) Run(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// holdStill stops the containers this step runs instead of, in the order it named them.
|
||||
//
|
||||
// A stop that fails is said and not fatal. The step runs anyway: for the case this exists for —
|
||||
// a collector walking storage nothing must be writing to — a writer that would not stop is worth
|
||||
// knowing about, and refusing to run would mean the work never happens and the log says nothing
|
||||
// new each night. What must not be skipped is the restart, and it is not: it is deferred.
|
||||
func (s *Scheduler) holdStill(ctx context.Context, cri string, j *scheduledJob) {
|
||||
for _, name := range j.hold {
|
||||
if _, err := s.run(ctx, cri, "stop", name); err != nil {
|
||||
s.log(fmt.Sprintf("scheduled step %s: could not stop %s for the run: %v", j.id, name, err))
|
||||
continue
|
||||
}
|
||||
s.log(fmt.Sprintf("scheduled step %s: %s held still for the run", j.id, name))
|
||||
}
|
||||
}
|
||||
|
||||
// letRun starts them again, in the reverse of the order they were stopped, and says so loudly if
|
||||
// one does not come back.
|
||||
//
|
||||
// **Reverse order**, because stopping walks a dependency the other way: a module that holds two
|
||||
// containers still names the one that depends on the other first, and bringing them back the same
|
||||
// way would start a dependant before what it depends on.
|
||||
//
|
||||
// Given its own context, because this runs in a defer and the one the run used may already be
|
||||
// cancelled — a host shutting down mid-window would otherwise leave the service stopped, which is
|
||||
// precisely the outcome this field must never have.
|
||||
func (s *Scheduler) letRun(_ context.Context, cri string, j *scheduledJob) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), closingWindow)
|
||||
defer cancel()
|
||||
for i := len(j.hold) - 1; i >= 0; i-- {
|
||||
name := j.hold[i]
|
||||
if _, err := s.run(ctx, cri, "start", name); err != nil {
|
||||
// Said as loudly as this host says anything: a service the mesh stopped for a
|
||||
// maintenance window and could not start again is down, and nothing else will notice
|
||||
// until the next apply compares it.
|
||||
s.log(fmt.Sprintf(
|
||||
"scheduled step %s: %s was held still for the run and WILL NOT START AGAIN: %v",
|
||||
j.id, name, err))
|
||||
continue
|
||||
}
|
||||
s.log(fmt.Sprintf("scheduled step %s: %s running again", j.id, name))
|
||||
}
|
||||
}
|
||||
|
||||
// closingWindow is how long the host will spend putting back what it stopped. Generous: this is
|
||||
// the half that must not be given up on.
|
||||
const closingWindow = 5 * time.Minute
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
package apply
|
||||
|
||||
// ForTests points where the host writes units and unpacks daemons at directories a test owns, and
|
||||
// waits nothing between its looks at a unit, until the returned function puts them back. For tests
|
||||
// in other packages that apply a process — the bootstrap's, which checks that what genesis raises is
|
||||
// what the controller's process takes over (novox/hq issue 223). Nothing outside a test calls it.
|
||||
func ForTests(units, daemons string) (restore func()) {
|
||||
wasUnits, wasDaemons, wasSettle, wasHandover := unitDir, daemonRoot, serviceSettle, handoverSettle
|
||||
unitDir, daemonRoot, serviceSettle, handoverSettle = units, daemons, 0, 0
|
||||
return func() {
|
||||
unitDir, daemonRoot, serviceSettle, handoverSettle = wasUnits, wasDaemons, wasSettle, wasHandover
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,8 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -112,6 +114,43 @@ func BuildControlPlane(ctx context.Context, run Runner, builderTag string, sourc
|
||||
return Built{}, nil
|
||||
}
|
||||
|
||||
// **Which form the controller is in at that commit** (novox/hq issue 223). An image the module
|
||||
// builds is the form genesis always raised, and the builder builds it as before. A process the
|
||||
// module runs from a Go bundle cannot be built here — its toolchain is one the mesh makes later —
|
||||
// so genesis builds the controller's own Dockerfile and raises it as the container that process
|
||||
// replaces (genesis_form.go).
|
||||
workspace, err := os.MkdirTemp("", "mesh-genesis-*")
|
||||
if err != nil {
|
||||
return Built{}, err
|
||||
}
|
||||
defer os.RemoveAll(workspace)
|
||||
dir, commit, err := cloneAt(ctx, run, builderTag, source, workspace)
|
||||
if err != nil {
|
||||
return Built{}, fmt.Errorf("the control plane could not be fetched from %s at %s: %w",
|
||||
source.Repository, shortRef(source.Ref), err)
|
||||
}
|
||||
raw, err := os.ReadFile(filepath.Join(dir, "module.json"))
|
||||
if err != nil {
|
||||
return Built{}, fmt.Errorf("%s at %s has no module manifest: %w", source.Repository, shortRef(source.Ref), err)
|
||||
}
|
||||
form, err := processFormOf(raw)
|
||||
if err != nil {
|
||||
return Built{}, err
|
||||
}
|
||||
if form.Found {
|
||||
image, err := buildGenesisImage(ctx, run, dir)
|
||||
if err != nil {
|
||||
return Built{}, err
|
||||
}
|
||||
manifest, err := genesisForm(raw, form, image)
|
||||
if err != nil {
|
||||
return Built{}, err
|
||||
}
|
||||
say(fmt.Sprintf(" built %s from %s, as the container its process %s replaces (%s)",
|
||||
ControlPlaneModule, shortRef(commit), form.Process, form.Replaces))
|
||||
return Built{Module: ControlPlaneModule, Commit: commit, Image: image, Manifest: manifest}, nil
|
||||
}
|
||||
|
||||
out, err := run(ctx, "docker", args...)
|
||||
if err != nil {
|
||||
return Built{}, fmt.Errorf("the control plane could not be built from %s at %s: %w",
|
||||
|
||||
@@ -49,10 +49,10 @@ func TestARepositoryAndACommitIsEnough(t *testing.T) {
|
||||
func TestTheBuildHandsOverTheManifestTheMeshWillHold(t *testing.T) {
|
||||
manifest := `{"module":"mesh-controller","version":"1","resources":[` +
|
||||
`{"id":"server","type":"container","name":"mesh-controller","image":"` + builtImage + `"}]}`
|
||||
runtime := &asked{answer: func(string, []string) (string, error) {
|
||||
runtime := &asked{answer: aRepository(t, imageFormManifest, func(string, []string) (string, error) {
|
||||
return `{"module":"mesh-controller","commit":"a1b2c3d4","manifest":` + manifest +
|
||||
`,"made":[{"name":"server","kind":"image","reference":"` + builtImage + `"}]}` + "\n", nil
|
||||
}}
|
||||
})}
|
||||
built, err := BuildControlPlane(context.Background(), runtime.run, "mesh-builder:test",
|
||||
Source{Repository: "https://example.invalid/mesh-controller.git", Ref: "a1b2c3d4"}, false, func(string) {})
|
||||
if err != nil {
|
||||
@@ -69,10 +69,10 @@ func TestTheBuildHandsOverTheManifestTheMeshWillHold(t *testing.T) {
|
||||
// A result without a manifest is a build the installer cannot finish, and it is refused beside the
|
||||
// builder that said it rather than at step 9 with a message about a missing file.
|
||||
func TestABuildReportingNoManifestIsRefused(t *testing.T) {
|
||||
runtime := &asked{answer: func(string, []string) (string, error) {
|
||||
runtime := &asked{answer: aRepository(t, imageFormManifest, func(string, []string) (string, error) {
|
||||
return `{"module":"mesh-controller","commit":"a1b2c3d4",` +
|
||||
`"made":[{"name":"server","kind":"image","reference":"` + builtImage + `"}]}`, nil
|
||||
}}
|
||||
})}
|
||||
_, err := BuildControlPlane(context.Background(), runtime.run, "mesh-builder:test",
|
||||
Source{Repository: "https://example.invalid/mesh-controller.git", Ref: "a1b2c3d4"}, false, func(string) {})
|
||||
if err == nil || !strings.Contains(err.Error(), "no manifest") {
|
||||
@@ -82,11 +82,11 @@ func TestABuildReportingNoManifestIsRefused(t *testing.T) {
|
||||
|
||||
// And a manifest that does not name the image the build produced describes some other build.
|
||||
func TestABuildWhoseManifestNamesAnotherImageIsRefused(t *testing.T) {
|
||||
runtime := &asked{answer: func(string, []string) (string, error) {
|
||||
runtime := &asked{answer: aRepository(t, imageFormManifest, func(string, []string) (string, error) {
|
||||
return `{"module":"mesh-controller","commit":"a1b2c3d4","manifest":{"module":"mesh-controller",` +
|
||||
`"resources":[{"id":"server","type":"container","image":"sha256:` + strings.Repeat("9", 64) + `"}]},` +
|
||||
`"made":[{"name":"server","kind":"image","reference":"` + builtImage + `"}]}`, nil
|
||||
}}
|
||||
})}
|
||||
_, err := BuildControlPlane(context.Background(), runtime.run, "mesh-builder:test",
|
||||
Source{Repository: "https://example.invalid/mesh-controller.git", Ref: "a1b2c3d4"}, false, func(string) {})
|
||||
if err == nil || !strings.Contains(err.Error(), "does not name that image") {
|
||||
|
||||
@@ -0,0 +1,205 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// The controller as a process, raised at genesis as a container (novox/hq issue 213, issue 223).
|
||||
//
|
||||
// **The mesh runs the controller as a Go bundle the host starts as a process; genesis cannot.** A
|
||||
// process's bundle is fetched from the mesh's artifact store, which genesis raises long after the
|
||||
// controller, and a Go bundle is compiled in a toolchain the mesh builds later still. So genesis
|
||||
// pivots to the controller as it always has — an image it built from the controller's own
|
||||
// Dockerfile, run as a container the temporary controller composes — and the first time the mesh
|
||||
// builds the controller from its repository, the controller's own declaration is the process, which
|
||||
// names that container under `replaces`, and the host hands over: the process is started, seen up,
|
||||
// and only then is the container removed (mesh-host `replaces`, issue 213).
|
||||
//
|
||||
// **The container is genesis's shape, not the manifest's.** The manifest declares only the process.
|
||||
// From it genesis takes what the process is given — its environment, which is host paths and words —
|
||||
// and the id the process replaces; the container around it is written here: the image genesis built,
|
||||
// the host's network, and every host path the environment names mounted at the same path read-only.
|
||||
// It runs as the image's own unprivileged user (65534), so the secrets belong to that number until
|
||||
// the process's account takes them over — the image is FROM scratch and knows no account by name. Its
|
||||
// id is the one the process replaces, so the first declaration the mesh composes for this machine
|
||||
// hands this container over rather than leaving two controllers running.
|
||||
|
||||
// genesisUser is who the genesis container runs as — the image's own USER — and who its secrets
|
||||
// belong to until the process takes them over: the image has no passwd to look an account up in.
|
||||
const genesisUser = "65534:65534"
|
||||
|
||||
// ProcessForm is the controller's process, as its manifest declares it, and the container id that
|
||||
// process replaces. Found is false for a manifest in the image form — an older controller — which
|
||||
// genesis installs as it always did.
|
||||
type ProcessForm struct {
|
||||
Found bool
|
||||
Process string // the process resource's id
|
||||
Replaces string // the id of the container genesis raises in its place
|
||||
}
|
||||
|
||||
// processFormOf finds the controller's process in its manifest: a process resource running a bundle
|
||||
// the module builds, saying which one resource it replaces.
|
||||
func processFormOf(manifest []byte) (ProcessForm, error) {
|
||||
var m struct {
|
||||
Build *struct {
|
||||
Artifacts []struct {
|
||||
Name string `json:"name"`
|
||||
Kind string `json:"kind"`
|
||||
} `json:"artifacts"`
|
||||
} `json:"build"`
|
||||
Resources []map[string]any `json:"resources"`
|
||||
}
|
||||
if err := json.Unmarshal(manifest, &m); err != nil {
|
||||
return ProcessForm{}, fmt.Errorf("the %s module's manifest is not readable: %w", ControlPlaneModule, err)
|
||||
}
|
||||
kinds := map[string]string{}
|
||||
if m.Build != nil {
|
||||
for _, a := range m.Build.Artifacts {
|
||||
kinds[a.Name] = a.Kind
|
||||
}
|
||||
}
|
||||
for _, kind := range kinds {
|
||||
if kind == "image" {
|
||||
return ProcessForm{}, nil // the image form: the builder builds it, as before
|
||||
}
|
||||
}
|
||||
var found []ProcessForm
|
||||
for _, r := range m.Resources {
|
||||
if r["type"] != "process" || kinds[fmt.Sprint(r["artifact"])] != "bundle" {
|
||||
continue
|
||||
}
|
||||
if once, _ := r["run-once"].(bool); once {
|
||||
continue
|
||||
}
|
||||
replaces, _ := r["replaces"].([]any)
|
||||
if len(replaces) != 1 {
|
||||
return ProcessForm{}, fmt.Errorf(
|
||||
"the %s module runs as the process %v and says it replaces %v. Genesis raises the "+
|
||||
"controller as a container that process takes over, so the process names exactly "+
|
||||
"one resource it replaces — the id genesis gives the container",
|
||||
ControlPlaneModule, r["id"], r["replaces"])
|
||||
}
|
||||
found = append(found, ProcessForm{Found: true, Process: fmt.Sprint(r["id"]),
|
||||
Replaces: fmt.Sprint(replaces[0])})
|
||||
}
|
||||
if len(found) != 1 {
|
||||
return ProcessForm{}, fmt.Errorf(
|
||||
"the %s module builds no image and runs %d process(es) of its own; genesis raises one "+
|
||||
"controller, from the process its manifest declares", ControlPlaneModule, len(found))
|
||||
}
|
||||
return found[0], nil
|
||||
}
|
||||
|
||||
// genesisForm is the manifest genesis registers: the controller's own manifest, its process
|
||||
// replaced by the container genesis runs in its place, under the id the process replaces, running
|
||||
// the image genesis built. Resolved as a build would resolve it — no build section, the image named
|
||||
// — because that is what the temporary controller is handed.
|
||||
func genesisForm(manifest []byte, form ProcessForm, image string) ([]byte, error) {
|
||||
var m map[string]any
|
||||
if err := json.Unmarshal(manifest, &m); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resources, _ := m["resources"].([]any)
|
||||
var out []any
|
||||
for _, raw := range resources {
|
||||
r, _ := raw.(map[string]any)
|
||||
if r == nil || r["id"] != form.Process {
|
||||
out = append(out, raw)
|
||||
continue
|
||||
}
|
||||
env, _ := r["env"].(map[string]any)
|
||||
var volumes []any
|
||||
for _, key := range sortedAnyKeys(env) {
|
||||
value := fmt.Sprint(env[key])
|
||||
if strings.HasPrefix(value, "/") || strings.HasPrefix(value, "${dir:") {
|
||||
volumes = append(volumes, value+":"+value+":ro")
|
||||
}
|
||||
}
|
||||
container := map[string]any{
|
||||
"id": form.Replaces, "type": "container", "name": ControlPlaneModule,
|
||||
"image": image, "network": "host", "args": []any{"serve"},
|
||||
}
|
||||
if len(env) > 0 {
|
||||
container["env"] = env
|
||||
}
|
||||
if len(volumes) > 0 {
|
||||
container["volumes"] = volumes
|
||||
}
|
||||
out = append(out, container)
|
||||
}
|
||||
m["resources"] = out
|
||||
// What the container reads must be readable by who it runs as. The process's account owns them
|
||||
// once the process takes over, and the host gives them to it in the same apply.
|
||||
m["secrets-owner"] = genesisUser
|
||||
// **Nothing to prepare at genesis.** The temporary controller — the same commit — migrated the
|
||||
// stores when the foundation raised it, and a preparation step is derived from a resource running
|
||||
// an artifact the module built, which a pinned image is not: the controller refuses `prepares`
|
||||
// with nothing to run it in. The process prepares the stores itself when it takes over.
|
||||
delete(m, "prepares")
|
||||
delete(m, "build")
|
||||
var b bytes.Buffer
|
||||
enc := json.NewEncoder(&b)
|
||||
enc.SetEscapeHTML(false)
|
||||
if err := enc.Encode(m); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return bytes.TrimSpace(b.Bytes()), nil
|
||||
}
|
||||
|
||||
func sortedAnyKeys(m map[string]any) []string {
|
||||
keys := make([]string, 0, len(m))
|
||||
for k := range m {
|
||||
keys = append(keys, k)
|
||||
}
|
||||
sort.Strings(keys)
|
||||
return keys
|
||||
}
|
||||
|
||||
// cloneAt fetches the controller's repository at the commit genesis builds, into a directory of
|
||||
// this machine's, with the carried builder's git — the machine is not assumed to have one. Returns
|
||||
// the module's directory and the commit that was checked out.
|
||||
func cloneAt(ctx context.Context, run Runner, builderTag string, source Source, into string) (string, string, error) {
|
||||
git := func(args ...string) (string, error) {
|
||||
return run(ctx, "docker", append([]string{"run", "--rm", "-v", into + ":/ws",
|
||||
"--entrypoint", "git", builderTag}, args...)...)
|
||||
}
|
||||
if _, err := git("clone", "--quiet", source.Repository, "/ws/src"); err != nil {
|
||||
return "", "", fmt.Errorf("cloning %s: %w", source.Repository, err)
|
||||
}
|
||||
if _, err := git("-C", "/ws/src", "checkout", "--quiet", "--detach", source.Ref); err != nil {
|
||||
return "", "", fmt.Errorf("checking out %s: %w", shortRef(source.Ref), err)
|
||||
}
|
||||
commit, err := git("-C", "/ws/src", "rev-parse", "HEAD")
|
||||
if err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
return filepath.Join(into, "src", source.Path), strings.TrimSpace(commit), nil
|
||||
}
|
||||
|
||||
// buildGenesisImage builds the controller's image from its own Dockerfile (the one `make image`
|
||||
// uses), on this machine, and names it by the digest of its own configuration, as the builder did.
|
||||
func buildGenesisImage(ctx context.Context, run Runner, dir string) (string, error) {
|
||||
if _, err := os.Stat(filepath.Join(dir, "Dockerfile")); err != nil {
|
||||
return "", fmt.Errorf("the %s repository has no Dockerfile, so genesis has no image to raise "+
|
||||
"the controller from: %w", ControlPlaneModule, err)
|
||||
}
|
||||
out, err := run(ctx, "docker", "build", "--quiet", dir)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("building the %s image: %w", ControlPlaneModule, err)
|
||||
}
|
||||
image := strings.TrimSpace(out)
|
||||
if i := strings.LastIndex(image, "\n"); i >= 0 {
|
||||
image = strings.TrimSpace(image[i+1:])
|
||||
}
|
||||
if !strings.HasPrefix(image, "sha256:") {
|
||||
return "", fmt.Errorf("docker build said %q, which is not an image id", firstLine(out))
|
||||
}
|
||||
return image, nil
|
||||
}
|
||||
@@ -0,0 +1,182 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// imageFormManifest is a controller from before issue 213: it builds an image and runs a container.
|
||||
const imageFormManifest = `{"module":"mesh-controller","version":"1","resources":[
|
||||
{"id":"server","type":"container","name":"mesh-controller","network":"host","artifact":"server"}],
|
||||
"build":{"artifacts":[{"name":"server","kind":"image","from":"Dockerfile"}]}}`
|
||||
|
||||
// processFormManifest is the controller's manifest as novox/mesh-controller declares it after issue
|
||||
// 213: a Go bundle the host runs as a process, replacing the container it ran as.
|
||||
const processFormManifest = `{
|
||||
"module": "mesh-controller", "version": "1", "slug": "control", "prepares": true,
|
||||
"claims": [{"name": "mesh-controller", "scope": "mesh"}],
|
||||
"accesses": [{"path": "/var/lib/mesh-broker-tls", "mode": "read"}],
|
||||
"own-secrets": {"inventory": "${dir:mesh-state}/inventory", "bus": "${dir:mesh-state}/bus"},
|
||||
"secrets-owner": "mesh-controller",
|
||||
"tools": ["status"],
|
||||
"resources": [
|
||||
{"id": "account", "type": "user", "name": "mesh-controller", "shell": "/usr/bin/nologin", "home": "/var/lib/mesh-controller"},
|
||||
{"id": "mesh-state", "type": "directory", "mode": "0700", "place": "mesh", "owner": "mesh-controller"},
|
||||
{"id": "controller", "type": "process", "name": "mesh-controller", "artifact": "controller",
|
||||
"run": ["./mesh-controller", "serve"], "user": "mesh-controller",
|
||||
"env": {"MESH_BROKER_CERTIFICATE": "/var/lib/mesh-broker-tls/tls.crt",
|
||||
"MESH_STORE_INVENTORY_FILE": "${dir:mesh-state}/inventory",
|
||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_BUS_NATS_FILE": "${dir:mesh-state}/bus"},
|
||||
"replaces": ["server"]}
|
||||
],
|
||||
"build": {"artifacts": [{"name": "controller", "kind": "bundle", "language": "go", "system": "arch",
|
||||
"from": "cmd/mesh-controller", "binary": "mesh-controller"}]}
|
||||
}`
|
||||
|
||||
// aRepository answers the carried builder's git as a clone of a repository holding this manifest and
|
||||
// a Dockerfile, and hands every other command to then.
|
||||
func aRepository(t *testing.T, manifest string, then func(string, []string) (string, error)) func(string, []string) (string, error) {
|
||||
t.Helper()
|
||||
return func(name string, args []string) (string, error) {
|
||||
if name == "docker" && len(args) > 5 && args[0] == "run" && args[4] == "--entrypoint" && args[5] == "git" {
|
||||
host := strings.TrimSuffix(args[3], ":/ws")
|
||||
joined := strings.Join(args, " ")
|
||||
switch {
|
||||
case strings.Contains(joined, " clone "):
|
||||
src := filepath.Join(host, "src")
|
||||
if err := os.MkdirAll(src, 0o755); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(src, "module.json"), []byte(manifest), 0o644); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return "", os.WriteFile(filepath.Join(src, "Dockerfile"), []byte("FROM scratch\n"), 0o644)
|
||||
case strings.Contains(joined, "rev-parse"):
|
||||
return "a1b2c3d4e5f6\n", nil
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
return then(name, args)
|
||||
}
|
||||
}
|
||||
|
||||
// novox/hq issue 223: a controller declared as a process is raised at genesis as a container built
|
||||
// from its own Dockerfile, not by the builder — which has no Go toolchain at genesis and would refuse.
|
||||
func TestAProcessFormControllerIsBuiltFromItsDockerfileAndRaisedAsAContainer(t *testing.T) {
|
||||
runtime := &asked{answer: aRepository(t, processFormManifest, func(name string, args []string) (string, error) {
|
||||
if name == "docker" && args[0] == "build" {
|
||||
return builtImage + "\n", nil
|
||||
}
|
||||
t.Fatalf("genesis ran %s %v; a process-form controller is built from its Dockerfile alone", name, args)
|
||||
return "", nil
|
||||
})}
|
||||
built, err := BuildControlPlane(context.Background(), runtime.run, "mesh-builder:test",
|
||||
Source{Repository: "https://example.invalid/mesh-controller.git", Ref: "a1b2c3d4"}, false, func(string) {})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if built.Image != builtImage || built.Commit != "a1b2c3d4e5f6" {
|
||||
t.Errorf("built %q from %q", built.Image, built.Commit)
|
||||
}
|
||||
if runtime.ran("mesh-builder:test build") {
|
||||
t.Error("the builder was asked to build a controller it cannot build at genesis")
|
||||
}
|
||||
|
||||
var m map[string]any
|
||||
if err := json.Unmarshal(built.Manifest, &m); err != nil {
|
||||
t.Fatalf("the genesis manifest is not JSON: %v", err)
|
||||
}
|
||||
if _, has := m["prepares"]; has {
|
||||
t.Error("the genesis manifest prepares its state; the temporary controller already did, and a pinned image is nothing the controller can derive a step from")
|
||||
}
|
||||
if _, has := m["build"]; has {
|
||||
t.Error("the genesis manifest still says how it is built; it is handed over resolved")
|
||||
}
|
||||
var container map[string]any
|
||||
for _, raw := range m["resources"].([]any) {
|
||||
r := raw.(map[string]any)
|
||||
if r["type"] == "process" {
|
||||
t.Errorf("the genesis manifest still runs the process: %v", r)
|
||||
}
|
||||
if r["type"] == "container" {
|
||||
container = r
|
||||
}
|
||||
}
|
||||
if container == nil {
|
||||
t.Fatal("the genesis manifest runs no container")
|
||||
}
|
||||
for key, want := range map[string]any{"id": "server", "name": "mesh-controller", "image": builtImage,
|
||||
"network": "host"} {
|
||||
if container[key] != want {
|
||||
t.Errorf("the genesis container's %s is %v, not %v", key, container[key], want)
|
||||
}
|
||||
}
|
||||
if _, has := container["user"]; has {
|
||||
t.Error("the genesis container says a user; the host's container has no such field, the image's USER is who it runs as")
|
||||
}
|
||||
if m["secrets-owner"] != "65534:65534" {
|
||||
t.Errorf("the secrets belong to %v, which the container cannot read as", m["secrets-owner"])
|
||||
}
|
||||
volumes, _ := json.Marshal(container["volumes"])
|
||||
for _, want := range []string{"${dir:mesh-state}/inventory:${dir:mesh-state}/inventory:ro",
|
||||
"/var/lib/mesh-broker-tls/tls.crt:/var/lib/mesh-broker-tls/tls.crt:ro"} {
|
||||
if !strings.Contains(string(volumes), want) {
|
||||
t.Errorf("the container does not mount %s: %s", want, volumes)
|
||||
}
|
||||
}
|
||||
if strings.Contains(string(volumes), "seat:") {
|
||||
t.Errorf("a word that is not a path was mounted: %s", volumes)
|
||||
}
|
||||
|
||||
// And it is what step 9 installs: pinned, its container found, its stores delivered from its
|
||||
// environment — the same path an image-form controller takes.
|
||||
pinned, _, err := pinImage(built.Manifest, built.Image, "registry.internal:5000/mesh-controller@sha256:"+strings.Repeat("e", 64), ControlPlaneModule)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if id := controlPlaneResourceIn(pinned); id != "server" {
|
||||
t.Errorf("step 9 finds the controller's container as %q", id)
|
||||
}
|
||||
wanted, err := secretsByVariableIn(pinned)
|
||||
if err != nil {
|
||||
t.Fatalf("step 9 cannot deliver the stores into the genesis container: %v", err)
|
||||
}
|
||||
if wanted["MESH_STORE_INVENTORY"] != "inventory" {
|
||||
t.Errorf("step 9 delivers %v", wanted)
|
||||
}
|
||||
}
|
||||
|
||||
// An older controller — an image and a container — is still built by the builder, as before.
|
||||
func TestAnImageFormControllerIsStillBuiltByTheBuilder(t *testing.T) {
|
||||
manifest := `{"module":"mesh-controller","version":"1","resources":[` +
|
||||
`{"id":"server","type":"container","name":"mesh-controller","image":"` + builtImage + `"}]}`
|
||||
runtime := &asked{answer: aRepository(t, imageFormManifest, func(name string, args []string) (string, error) {
|
||||
if name == "docker" && args[0] == "build" {
|
||||
t.Fatal("an image-form controller was built from its Dockerfile rather than by the builder")
|
||||
}
|
||||
return `{"module":"mesh-controller","commit":"a1b2c3d4","manifest":` + manifest +
|
||||
`,"made":[{"name":"server","kind":"image","reference":"` + builtImage + `"}]}`, nil
|
||||
})}
|
||||
built, err := BuildControlPlane(context.Background(), runtime.run, "mesh-builder:test",
|
||||
Source{Repository: "https://example.invalid/mesh-controller.git", Ref: "a1b2c3d4"}, false, func(string) {})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(built.Manifest) != manifest || !runtime.ran("mesh-builder:test build") {
|
||||
t.Errorf("the image form did not go through the builder: %s", built.Manifest)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAProcessFormNamingNoOneReplacementIsRefused(t *testing.T) {
|
||||
for _, bad := range []string{`"replaces": []`, `"replaces": ["a", "b"]`} {
|
||||
raw := strings.Replace(processFormManifest, `"replaces": ["server"]`, bad, 1)
|
||||
if _, err := processFormOf([]byte(raw)); err == nil {
|
||||
t.Errorf("a process saying %s was accepted; genesis would not know what to name its container", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,121 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"archive/tar"
|
||||
"bytes"
|
||||
"compress/gzip"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-host/internal/apply"
|
||||
"github.com/novox/mesh-host/internal/declaration"
|
||||
"github.com/novox/mesh-host/internal/store"
|
||||
"github.com/novox/mesh-host/internal/system"
|
||||
)
|
||||
|
||||
// novox/hq issue 223: what genesis raises is what the controller's process takes over. The temporary
|
||||
// controller composes the genesis container, and the host records it as `<module>.<its id>`; the
|
||||
// first declaration the mesh composes from the controller's real manifest names that same id under
|
||||
// the process's `replaces` (the composer prefixes both alike — mesh-controller's
|
||||
// TestTheControllerIsAProcessAndNoContainer). So the first apply hands over: the process is started,
|
||||
// seen up, and only then is the genesis container removed — one controller before, one after, never
|
||||
// none and never two left.
|
||||
func TestTheFirstApplyHandsTheGenesisContainerOverToTheProcess(t *testing.T) {
|
||||
restore := apply.ForTests(t.TempDir(), t.TempDir())
|
||||
defer restore()
|
||||
|
||||
form, err := processFormOf([]byte(processFormManifest))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
genesis, err := genesisForm([]byte(processFormManifest), form, builtImage)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// What the host records for the genesis container, as the composer names it.
|
||||
recorded := ""
|
||||
var m struct {
|
||||
Resources []map[string]any `json:"resources"`
|
||||
}
|
||||
if err := json.Unmarshal(genesis, &m); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, r := range m.Resources {
|
||||
if r["type"] == "container" {
|
||||
recorded = ControlPlaneModule + "." + r["id"].(string)
|
||||
}
|
||||
}
|
||||
known := store.State{Resources: []store.Applied{{Origin: store.OriginDeclared, ID: recorded,
|
||||
Type: "container", Target: ControlPlaneModule}}}
|
||||
|
||||
// The controller's first composed declaration: its process, replacing what the manifest names.
|
||||
body, digest := aBundle(t)
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write(body) }))
|
||||
defer server.Close()
|
||||
replaces, _ := json.Marshal([]string{ControlPlaneModule + "." + form.Replaces})
|
||||
d, err := declaration.Parse([]byte(`{"declaration":1,"resources":[{"id":"` + ControlPlaneModule + "." + form.Process +
|
||||
`","type":"process","name":"mesh-controller","source":"` + server.URL + `/c.tgz","digest":"` + digest +
|
||||
`","run":["./mesh-controller","serve"],"replaces":` + string(replaces) + `}]}`))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var commands []string
|
||||
started, container := false, true
|
||||
run := func(_ context.Context, name string, args ...string) (string, error) {
|
||||
line := name + " " + strings.Join(args, " ")
|
||||
commands = append(commands, line)
|
||||
switch {
|
||||
case strings.HasPrefix(line, "systemctl restart mesh-controller.service"):
|
||||
started = true
|
||||
case strings.HasPrefix(line, "systemctl show mesh-controller.service") && started:
|
||||
return "ActiveState=active\nSubState=running\nMainPID=7\nNRestarts=0\n", nil
|
||||
case strings.HasPrefix(line, "docker rm -f mesh-controller"):
|
||||
if !started {
|
||||
t.Error("the genesis container was removed before the process was started")
|
||||
}
|
||||
container = false
|
||||
case strings.HasPrefix(line, "docker container inspect") && !container:
|
||||
return "", errors.New("no such container")
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
sys, err := system.For("arch")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, after, err := apply.Apply(context.Background(), sys, d, known, store.OriginDeclared, run, nil, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("the first apply did not hand over: %v\n%s", err, strings.Join(commands, "\n"))
|
||||
}
|
||||
if container {
|
||||
t.Fatalf("the genesis container is still running beside the process — two controllers:\n%s",
|
||||
strings.Join(commands, "\n"))
|
||||
}
|
||||
if _, still := after.Find(recorded); still {
|
||||
t.Error("the host still records the genesis container")
|
||||
}
|
||||
}
|
||||
|
||||
func aBundle(t *testing.T) ([]byte, string) {
|
||||
t.Helper()
|
||||
var raw bytes.Buffer
|
||||
zipped := gzip.NewWriter(&raw)
|
||||
w := tar.NewWriter(zipped)
|
||||
content := "#!/bin/sh\n"
|
||||
if err := w.WriteHeader(&tar.Header{Name: "mesh-controller", Mode: 0o755, Size: int64(len(content)), Typeflag: tar.TypeReg}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, _ = w.Write([]byte(content))
|
||||
_ = w.Close()
|
||||
_ = zipped.Close()
|
||||
sum := sha256.Sum256(raw.Bytes())
|
||||
return raw.Bytes(), "sha256:" + hex.EncodeToString(sum[:])
|
||||
}
|
||||
@@ -569,6 +569,18 @@ type Process struct {
|
||||
// that runs once does not run on a schedule, and something that is not running cannot be
|
||||
// restarted when a file changes.
|
||||
Schedule string `json:"schedule,omitempty"`
|
||||
|
||||
// Replaces names resources this declaration no longer declares that this process takes the
|
||||
// place of (novox/hq issue 213). Such a resource is not removed with the other orphans, before
|
||||
// anything is applied: it is removed only once this process is applied and still running a
|
||||
// moment later, and kept when it is not. So the thing being replaced answers until the thing
|
||||
// replacing it does — the controller moving from its container to a process is the case: the
|
||||
// container removed first left nothing answering the mesh's verbs for as long as fetching,
|
||||
// unpacking and starting the process took, and for ever if the process did not start.
|
||||
//
|
||||
// For a process that stays up; a step or a scheduled run is not running a moment later by
|
||||
// design, so there is nothing to hand over to.
|
||||
Replaces []string `json:"replaces,omitempty"`
|
||||
}
|
||||
|
||||
func (d *Process) Identity() string { return d.ID }
|
||||
@@ -654,6 +666,20 @@ func (d *Process) validate(where string, _ bool) []string {
|
||||
problems = append(problems, where+": "+err.Error())
|
||||
}
|
||||
}
|
||||
if len(d.Replaces) > 0 && (d.RunOnce || d.Schedule != "") {
|
||||
problems = append(problems, where+
|
||||
": only a process that stays up replaces something — a step or a scheduled run is not "+
|
||||
"running a moment later, so what it replaced would be removed with nothing in its place, "+
|
||||
"or never")
|
||||
}
|
||||
for _, id := range d.Replaces {
|
||||
switch {
|
||||
case strings.TrimSpace(id) == "":
|
||||
problems = append(problems, where+": replaces names an empty id")
|
||||
case id == d.ID:
|
||||
problems = append(problems, where+": a process cannot replace itself")
|
||||
}
|
||||
}
|
||||
return problems
|
||||
}
|
||||
|
||||
@@ -1011,6 +1037,24 @@ type Container struct {
|
||||
// rather than stacked. It is exclusive with RunOnce and with restart-on: a container runs once
|
||||
// and gates, runs on a cadence, or stays up — never two of these.
|
||||
Schedule string `json:"schedule,omitempty"`
|
||||
|
||||
// WhileStopped names resources of the same module — containers — that must be held still for
|
||||
// the duration of this step's run (novox/hq ADR 0189). The host stops each before the run and
|
||||
// starts each again after it, **whatever the step did**: a step that failed must leave the
|
||||
// service running, because the one real risk of this field is a window that never closes.
|
||||
//
|
||||
// **For the work a service cannot have done underneath it.** The artifact store's collector
|
||||
// walks the storage and requires every writer stopped; a run-once step runs beside containers
|
||||
// and a scheduled one is the same container again, so until this there was no way for a module
|
||||
// to say it. The predecessor said it with a shell script, which is how the mesh inherited a
|
||||
// store that has never collected anything.
|
||||
//
|
||||
// **Its own module's containers, and only on a schedule.** A module that could quiesce a
|
||||
// neighbour could stop the mesh. And at apply time the host already has a window — the
|
||||
// declaration is applied in order and a run-once step gates what follows — so a one-time
|
||||
// offline job says *before*, not *instead of*; a recurring window is the case order cannot
|
||||
// express, and the only one this serves.
|
||||
WhileStopped []string `json:"while-stopped,omitempty"`
|
||||
}
|
||||
|
||||
func (c *Container) Identity() string { return c.ID }
|
||||
@@ -1042,6 +1086,20 @@ func (c *Container) validate(where string, _ bool) []string {
|
||||
problems = append(problems, where+": "+err.Error())
|
||||
}
|
||||
}
|
||||
// A maintenance window belongs to a recurring step (novox/hq ADR 0189). Refused on anything
|
||||
// else here, where the field is; that it names containers of the same module, and not itself,
|
||||
// is judged against the whole declaration (see whileStoppedNames).
|
||||
if len(c.WhileStopped) > 0 && c.Schedule == "" {
|
||||
problems = append(problems, where+": while-stopped needs a schedule; at apply the host "+
|
||||
"already has a window — the declaration is applied in order and a run-once step gates "+
|
||||
"what follows — so a one-time offline job is declared before what it works on")
|
||||
}
|
||||
for _, id := range c.WhileStopped {
|
||||
if id == c.ID {
|
||||
problems = append(problems, where+": while-stopped names "+strconv.Quote(id)+
|
||||
", which is this step itself")
|
||||
}
|
||||
}
|
||||
// The runtime's flags take addresses, and a name here would be handed to it verbatim and
|
||||
// refused at create — after the old container was already removed. Refused on arrival instead.
|
||||
for _, d := range c.Dns {
|
||||
@@ -1456,7 +1514,9 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
|
||||
problems = append(problems, resource.validate(where, allowActions)...)
|
||||
d.Resources = append(d.Resources, resource)
|
||||
}
|
||||
problems = append(problems, checkWhileStopped(d.Resources)...)
|
||||
problems = append(problems, checkAdoption(env.Adoption, d.Resources, allowActions)...)
|
||||
problems = append(problems, checkReplaces(d.Resources)...)
|
||||
if env.Adoption == nil {
|
||||
for _, r := range d.Resources {
|
||||
if r.Kind() == TypeOpening {
|
||||
@@ -1484,6 +1544,41 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
|
||||
return d, nil
|
||||
}
|
||||
|
||||
// checkWhileStopped judges a maintenance window against the whole declaration (novox/hq ADR 0189).
|
||||
//
|
||||
// A step may hold still only a container that is **here** — in this same declaration, which is to
|
||||
// say on this machine and placed by the mesh. That is what makes it the module's own: a node's
|
||||
// declaration carries one module's resources beside another's, so the id must also be a container
|
||||
// and not a file or a directory, which there would be nothing to stop.
|
||||
//
|
||||
// Refused on arrival rather than discovered at the first fire. A window that names something the
|
||||
// host cannot stop is a window that opens at 03:00 and reports nothing until somebody reads a log.
|
||||
func checkWhileStopped(resources []Resource) []string {
|
||||
containers := map[string]bool{}
|
||||
for _, r := range resources {
|
||||
if r.Kind() == TypeContainer {
|
||||
containers[r.Identity()] = true
|
||||
}
|
||||
}
|
||||
var problems []string
|
||||
for _, r := range resources {
|
||||
c, ok := r.(*Container)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
for _, id := range c.WhileStopped {
|
||||
if containers[id] {
|
||||
continue
|
||||
}
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"resource %q: while-stopped names %q, and this declaration has no container by "+
|
||||
"that id. A step may hold still only a container placed on this machine "+
|
||||
"beside it", c.ID, id))
|
||||
}
|
||||
}
|
||||
return problems
|
||||
}
|
||||
|
||||
func strictDecode(raw []byte, into any) error {
|
||||
// DisallowUnknownFields is the whole point rather than strictness for its own sake: a
|
||||
// field the host does not know is a thing the control plane believes it asked for.
|
||||
@@ -1640,3 +1735,36 @@ func stripComments(raw []byte) []byte {
|
||||
}
|
||||
return []byte(strings.Join(kept, "\n"))
|
||||
}
|
||||
|
||||
// checkReplaces refuses a `replaces` naming something this same declaration still declares, or
|
||||
// named by two processes (novox/hq issue 213). What is replaced is what the declaration no longer
|
||||
// says — a resource still declared is applied, not handed over, and one handed to two replacements
|
||||
// would go when the first of them came up, whatever became of the second.
|
||||
func checkReplaces(resources []Resource) []string {
|
||||
declared := map[string]bool{}
|
||||
for _, r := range resources {
|
||||
declared[r.Identity()] = true
|
||||
}
|
||||
var problems []string
|
||||
by := map[string]string{}
|
||||
for _, r := range resources {
|
||||
p, ok := r.(*Process)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
for _, id := range p.Replaces {
|
||||
if declared[id] {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"resource %q: replaces %q, which this declaration still declares — a process "+
|
||||
"replaces what the declaration no longer says", p.ID, id))
|
||||
}
|
||||
if other, taken := by[id]; taken && other != p.ID {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"resource %q: replaces %q, which %q replaces too — one thing has one replacement",
|
||||
p.ID, id, other))
|
||||
}
|
||||
by[id] = p.ID
|
||||
}
|
||||
}
|
||||
return problems
|
||||
}
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
package declaration
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// novox/hq issue 213: a process says what it takes the place of, so the host can keep the old one
|
||||
// answering until the new one does. What it may name is narrow, and each refusal is said here.
|
||||
|
||||
const aReplacingProcess = `{"id":"m.controller","type":"process","name":"m","source":"https://store.invalid/m",
|
||||
"digest":"sha256:` + "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + `","run":["./m","serve"]%s}`
|
||||
|
||||
func replacing(extra string) string {
|
||||
return `{"declaration":1,"resources":[` + strings.Replace(aReplacingProcess, "%s", extra, 1) + `]}`
|
||||
}
|
||||
|
||||
func TestAProcessMayNameWhatItReplaces(t *testing.T) {
|
||||
d, err := Parse([]byte(replacing(`,"replaces":["m.server"]`)))
|
||||
if err != nil {
|
||||
t.Fatalf("a process replacing an undeclared container was refused: %v", err)
|
||||
}
|
||||
if p := d.Resources[0].(*Process); len(p.Replaces) != 1 || p.Replaces[0] != "m.server" {
|
||||
t.Fatalf("replaces was not read: %+v", p)
|
||||
}
|
||||
}
|
||||
|
||||
// What is replaced is what the declaration no longer says. A resource still declared is applied,
|
||||
// and handing it over as well would remove something the same declaration asks to keep.
|
||||
func TestAProcessMayNotReplaceSomethingStillDeclared(t *testing.T) {
|
||||
raw := `{"declaration":1,"resources":[
|
||||
{"id":"m.server","type":"directory","path":"/tmp/x"},
|
||||
` + strings.Replace(aReplacingProcess, "%s", `,"replaces":["m.server"]`, 1) + `]}`
|
||||
if _, err := Parse([]byte(raw)); err == nil || !strings.Contains(err.Error(), "still declares") {
|
||||
t.Fatalf("a process replacing a declared resource was accepted: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOnlyAProcessThatStaysUpReplacesAnything(t *testing.T) {
|
||||
for _, mode := range []string{`,"run-once":true`, `,"schedule":"0 3 * * *"`} {
|
||||
if _, err := Parse([]byte(replacing(mode + `,"replaces":["m.server"]`))); err == nil {
|
||||
t.Errorf("a process with %s was allowed to replace something", mode)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAProcessCannotReplaceItselfOrNothing(t *testing.T) {
|
||||
for _, bad := range []string{`,"replaces":["m.controller"]`, `,"replaces":[""]`} {
|
||||
if _, err := Parse([]byte(replacing(bad))); err == nil {
|
||||
t.Errorf("replaces %s was accepted", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOneThingHasOneReplacement(t *testing.T) {
|
||||
second := strings.NewReplacer(`"id":"m.controller"`, `"id":"m.other"`, `"name":"m"`, `"name":"other"`).
|
||||
Replace(strings.Replace(aReplacingProcess, "%s", `,"replaces":["m.server"]`, 1))
|
||||
raw := `{"declaration":1,"resources":[` +
|
||||
strings.Replace(aReplacingProcess, "%s", `,"replaces":["m.server"]`, 1) + "," + second + `]}`
|
||||
if _, err := Parse([]byte(raw)); err == nil {
|
||||
t.Fatal("two processes replacing one resource were accepted")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
package system
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os/exec"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// A user that does not exist yet is "absent", in the words the host's own runner uses
|
||||
// ("getent exited 2"), not only Go's ("exit status 2"). Matched as text, the runner's wording read as
|
||||
// a user database that did not answer, and the controller's account was never created.
|
||||
func TestAMissingUserIsAbsentInTheRunnersWords(t *testing.T) {
|
||||
for _, words := range []string{"getent exited 2: ", "exit status 2"} {
|
||||
run := func(context.Context, string, ...string) (string, error) { return "", errors.New(words) }
|
||||
_, found, err := LookUpUser(context.Background(), run, "nobody-here")
|
||||
if err != nil || found {
|
||||
t.Errorf("%q: found %v, err %v; want absent", words, found, err)
|
||||
}
|
||||
}
|
||||
run := func(context.Context, string, ...string) (string, error) { return "", errors.New("getent exited 1: ") }
|
||||
if _, _, err := LookUpUser(context.Background(), run, "x"); err == nil {
|
||||
t.Error("a database that failed read as an answer")
|
||||
}
|
||||
}
|
||||
|
||||
// The real command, through a real exit: getent's code for a key not found.
|
||||
func TestAMissingUserIsAbsentFromTheRealGetent(t *testing.T) {
|
||||
if _, err := exec.LookPath("getent"); err != nil {
|
||||
t.Skip("no getent here")
|
||||
}
|
||||
run := func(ctx context.Context, name string, args ...string) (string, error) {
|
||||
out, err := exec.CommandContext(ctx, name, args...).Output()
|
||||
return string(out), err
|
||||
}
|
||||
_, found, err := LookUpUser(context.Background(), run, "mesh-no-such-user-0b1f")
|
||||
if err != nil || found {
|
||||
t.Errorf("found %v, err %v; want absent", found, err)
|
||||
}
|
||||
}
|
||||
@@ -19,6 +19,9 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os/exec"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/novox/mesh-host/internal/declaration"
|
||||
@@ -95,8 +98,10 @@ func LookUpUser(ctx context.Context, run Runner, name string) (Login, bool, erro
|
||||
out, err := run(ctx, "getent", "passwd", name)
|
||||
if err != nil {
|
||||
// getent's own convention: 2 means the key was not found, which is the only failure that
|
||||
// means "no such user".
|
||||
if strings.Contains(err.Error(), "exit status 2") {
|
||||
// means "no such user". Asked of the exit code: matched as text, it looked for Go's wording
|
||||
// ("exit status 2") while the host's runner says "getent exited 2", so a user that did not
|
||||
// exist yet read as a database that did not answer, and no account was ever created.
|
||||
if code, ok := ExitCode(err); ok && code == 2 {
|
||||
return Login{}, false, nil
|
||||
}
|
||||
return Login{}, false, fmt.Errorf(
|
||||
@@ -221,3 +226,27 @@ func For(name string) (System, error) {
|
||||
func All() []System {
|
||||
return []System{arch{}, alpine{}, android{}}
|
||||
}
|
||||
|
||||
// ExitCode is the code a command exited with, when err says one: from the exit itself where the
|
||||
// runner kept it, else from the words either runner shape uses ("exit status N", "<cmd> exited N").
|
||||
func ExitCode(err error) (int, bool) {
|
||||
if err == nil {
|
||||
return 0, false
|
||||
}
|
||||
var exit *exec.ExitError
|
||||
if errors.As(err, &exit) {
|
||||
return exit.ExitCode(), true
|
||||
}
|
||||
if m := exitWords.FindStringSubmatch(err.Error()); len(m) == 3 {
|
||||
for _, g := range m[1:] {
|
||||
if g != "" {
|
||||
if n, convErr := strconv.Atoi(g); convErr == nil {
|
||||
return n, true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return 0, false
|
||||
}
|
||||
|
||||
var exitWords = regexp.MustCompile(`exit status (\d+)|exited (\d+)`)
|
||||
|
||||
Reference in New Issue
Block a user