Hand a replaced resource over to the process that replaces it (hq issue 213) #86
@@ -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
|
// 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.
|
// adopted node is removed as any orphan is.
|
||||||
var protecting, orphans []store.Applied
|
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) {
|
for _, orphan := range known.Orphans(declared, origin) {
|
||||||
if d.Adoption == nil && strings.HasPrefix(orphan.ID, declaration.AdoptionPrefix) {
|
if d.Adoption == nil && strings.HasPrefix(orphan.ID, declaration.AdoptionPrefix) {
|
||||||
protecting = append(protecting, orphan)
|
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))
|
log(fmt.Sprintf(" kept %s (%s): %s was left out of this declaration by the mesh, not removed", orphan.ID, orphan.Target, module))
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
if by, replaced := replacedBy[orphan.ID]; replaced {
|
||||||
|
handover[by] = append(handover[by], orphan)
|
||||||
|
continue
|
||||||
|
}
|
||||||
orphans = append(orphans, orphan)
|
orphans = append(orphans, orphan)
|
||||||
}
|
}
|
||||||
ordered := d.Resources
|
ordered := d.Resources
|
||||||
@@ -517,6 +535,11 @@ func ApplyKeeping(
|
|||||||
if c, ok := resource.(*declaration.Container); ok && c.RunOnce {
|
if c, ok := resource.(*declaration.Container); ok && c.RunOnce {
|
||||||
gates = true
|
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 {
|
if gates {
|
||||||
failed.Gated = true
|
failed.Gated = true
|
||||||
// **A module's step gates that module, not the machine** (novox/hq ADR 0136).
|
// **A module's step gates that module, not the machine** (novox/hq ADR 0136).
|
||||||
@@ -607,6 +630,28 @@ func ApplyKeeping(
|
|||||||
}
|
}
|
||||||
log(line)
|
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 {
|
if !orphansRemoved {
|
||||||
@@ -2453,3 +2498,102 @@ func moduleOf(identity string) (string, bool) {
|
|||||||
}
|
}
|
||||||
return identity[:at], true
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
// 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).
|
// record again — the node's runtime restarted every other cycle (novox/hq 04-ISSUES/210).
|
||||||
out.wrote = want
|
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
|
_ = active
|
||||||
return out, nil
|
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)
|
return out, fmt.Errorf("%s is installed and would not start: %w", r.Name, err)
|
||||||
}
|
}
|
||||||
out.Action = "updated"
|
out.Action = "updated"
|
||||||
|
|||||||
@@ -569,6 +569,18 @@ type Process struct {
|
|||||||
// that runs once does not run on a schedule, and something that is not running cannot be
|
// that runs once does not run on a schedule, and something that is not running cannot be
|
||||||
// restarted when a file changes.
|
// restarted when a file changes.
|
||||||
Schedule string `json:"schedule,omitempty"`
|
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 }
|
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())
|
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
|
return problems
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1457,6 +1483,7 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
|
|||||||
d.Resources = append(d.Resources, resource)
|
d.Resources = append(d.Resources, resource)
|
||||||
}
|
}
|
||||||
problems = append(problems, checkAdoption(env.Adoption, d.Resources, allowActions)...)
|
problems = append(problems, checkAdoption(env.Adoption, d.Resources, allowActions)...)
|
||||||
|
problems = append(problems, checkReplaces(d.Resources)...)
|
||||||
if env.Adoption == nil {
|
if env.Adoption == nil {
|
||||||
for _, r := range d.Resources {
|
for _, r := range d.Resources {
|
||||||
if r.Kind() == TypeOpening {
|
if r.Kind() == TypeOpening {
|
||||||
@@ -1640,3 +1667,36 @@ func stripComments(raw []byte) []byte {
|
|||||||
}
|
}
|
||||||
return []byte(strings.Join(kept, "\n"))
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user