diff --git a/internal/apply/apply.go b/internal/apply/apply.go index 1ccb3ca..f20a8ec 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -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 { @@ -2453,3 +2498,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, ", ") +} diff --git a/internal/apply/handover_test.go b/internal/apply/handover_test.go new file mode 100644 index 0000000..a3904f0 --- /dev/null +++ b/internal/apply/handover_test.go @@ -0,0 +1,336 @@ +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: a step is run by its unit — as the process it prepares for is — so it has +// that process's user, directory and environment, and `./name` is its own bundle's binary. It used +// to be run directly by the host, as root, in the host's directory, with no environment. +func TestAStepIsRunByItsUnitWithItsUserAndEnvironment(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, + "env":{"MESH_STORE_INVENTORY_FILE":"/var/lib/mesh/x/inventory"}}`) + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginDeclared, m.run, nil, nil) + if err != nil { + t.Fatalf("the step failed: %v", err) + } + if m.index("systemctl start mesh-controller-prepare.service") < 0 { + t.Fatalf("the step was not started through its unit: %v", m.commands) + } + 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) + } + for _, want := range []string{ + "Type=oneshot", + "ExecStart=" + filepath.Join(daemonRoot, "mesh-controller-prepare", "mesh-controller") + " prepare", + `Environment="MESH_STORE_INVENTORY_FILE=/var/lib/mesh/x/inventory"`, + } { + if !strings.Contains(string(unit), want) { + t.Errorf("the step's unit does not say %q:\n%s", want, unit) + } + } + for _, not := range []string{"Restart=always", "WantedBy="} { + if strings.Contains(string(unit), not) { + t.Errorf("the step's unit says %q, which makes it a service:\n%s", not, unit) + } + } + if o := outcomeOf(report, "mesh-controller.controller-prepare"); o.Action != "created" { + t.Errorf("the step's outcome is %+v", o) + } +} + +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") + } +} diff --git a/internal/apply/process.go b/internal/apply/process.go index a548d44..ad7bb93 100644 --- a/internal/apply/process.go +++ b/internal/apply/process.go @@ -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" @@ -113,10 +123,24 @@ func applyProcess(ctx context.Context, r *declaration.Process, run Runner, // **A step is run to completion, not installed.** What follows it is gated on it finishing, // 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. + // The record that it ran is the digest, which is why the identity above includes the command. + // + // **Run by its unit, as the process it prepares for is run** (novox/hq issue 213). It used to be + // run directly by the host: as root, in the host's own directory, with none of its environment — + // so a step reading its store's connection from a file its environment names found no variable, + // `./name` was looked for beside the host rather than in the bundle, and a step that must not be + // root was. A oneshot unit carries the same user, directory and environment as a daemon's, and + // starting one waits for it to exit and fails when it did not exit cleanly. It is neither enabled + // nor left running; its output is in the journal under its own name. if r.RunOnce { - if _, err := run(ctx, r.Run[0], r.Run[1:]...); err != nil { + 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: %w", r.Name, err) } out.Action = "created" @@ -215,9 +239,10 @@ 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 - // continuously between fires, which is the opposite of a schedule. + if r.Schedule != "" || r.RunOnce { + // Started by its timer, or by the host as a step, and expected to finish. Restarting it + // would have it run continuously, which is the opposite of either; and a step is not + // installed, so it is not wanted by anything at boot. b.WriteString("Type=oneshot\n") b.WriteString("\n") return strings.Replace(b.String(), "Type=simple\n", "", 1) diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index 734ac1a..0093b5d 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -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 } @@ -1457,6 +1483,7 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) { d.Resources = append(d.Resources, resource) } 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 { @@ -1640,3 +1667,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 +} diff --git a/internal/declaration/replaces_test.go b/internal/declaration/replaces_test.go new file mode 100644 index 0000000..c470505 --- /dev/null +++ b/internal/declaration/replaces_test.go @@ -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") + } +}