diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 74266f1..dd342f2 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -440,6 +440,15 @@ func enrol(ctx context.Context, opts options) error { } fmt.Printf("generated this node's identity: %s\n", mine.PublicBase64()) + // Its key on the private network, generated here and now for the same reason: the private + // half must never have been anywhere else. The mesh receives only the public half and uses it + // to compute a graph it cannot impersonate. + mine.Overlay, err = identity.GenerateOverlayKey() + if err != nil { + return err + } + fmt.Printf("generated this node's overlay key: %s\n", mine.Overlay.Public) + // What this machine can be asked to do, gathered before joining rather than after. The // control plane cannot decide what a node should run without it, so it travels with the // request instead of being asked for in a second round trip. @@ -450,7 +459,7 @@ func enrol(ctx context.Context, opts options) error { } reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret, - mine.Public, reported, opts.timeout) + mine.Public, mine.Overlay.Public, reported, opts.timeout) if err != nil { return err } @@ -488,6 +497,14 @@ func enrol(ctx context.Context, opts options) error { "identity could not be used after a restart. Nothing was saved", reply.Node) } + // Written before the identity, so a node that dies between the two has a key file with no + // identity — which enrols again cleanly — rather than an identity naming a key that is not + // there, which looks joined and cannot come up. + if err := os.WriteFile(identity.OverlayKeyPath(opts.state), + []byte(mine.Overlay.Private+"\n"), 0o600); err != nil { + return fmt.Errorf("cannot write this node's overlay key: %w", err) + } + fmt.Printf("\nenrolled as %s\n", reply.Node) fmt.Printf(" identity %s\n", identityPath) fmt.Printf(" queue %s\n", reply.Queue) diff --git a/examples/substrate-first-node.lock b/examples/substrate-first-node.lock index f0cd0d3..5e63f2d 100644 --- a/examples/substrate-first-node.lock +++ b/examples/substrate-first-node.lock @@ -15,6 +15,11 @@ // belong to the registry the lab raises, which is what a real node pulls from anyway — what is // required is a reference that is exact and cannot move (novox/hq ADR 0006). // +// The store waits up to three minutes rather than one. A machine that has just pulled the +// image and is running initdb for the first time can take longer than sixty seconds, and it +// failed that way three times in the lab -- a flaky bootstrap that a second run always fixed, +// which is the worst kind because it teaches people to run things twice. +// // The store's data is a NAMED VOLUME, not a directory on the machine. A directory the host // creates is owned by root, and the database runs as somebody else inside the container — so it // could not write, and the container crash-looped. A named volume lets the image set up its own @@ -51,7 +56,7 @@ "id": "store-ready", "type": "action", "in": "mesh-store", - "command": ["sh", "-c", "for i in $(seq 1 60); do pg_isready -U postgres >/dev/null 2>&1 && exit 0; sleep 1; done; exit 1"], + "command": ["sh", "-c", "for i in $(seq 1 180); do pg_isready -U postgres >/dev/null 2>&1 && exit 0; sleep 1; done; exit 1"], "verify": ["pg_isready", "-U", "postgres"] }, { @@ -74,7 +79,7 @@ "command": ["docker", "run", "--rm", "--network", "container:mesh-store", "-e", "MESH_STORE_INVENTORY=postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable", "-e", "MESH_STORE_IDENTITY=postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable", - "192.0.2.250:5000/mesh-control@sha256:b40717d6435077d4513e2348351fd2fd72ad790e0ccc2bba6aec326a0bc325b8", + "192.0.2.250:5000/mesh-control@sha256:c0f56edfb629ecb6a69b991b47abdbbb31c0da7dd3ce4db887d94efaa0b21632", "migrate"], "verify": ["sh", "-c", "docker exec mesh-store psql -U postgres -d inventory -tAc \"select to_regclass('public.node')\" | grep -qx node && docker exec mesh-store psql -U postgres -d identity -tAc \"select to_regclass('public.signing_key')\" | grep -qx signing_key"] }, @@ -108,7 +113,7 @@ "id": "control-plane", "type": "container", "name": "mesh-control", - "image": "192.0.2.250:5000/mesh-control@sha256:b40717d6435077d4513e2348351fd2fd72ad790e0ccc2bba6aec326a0bc325b8", + "image": "192.0.2.250:5000/mesh-control@sha256:c0f56edfb629ecb6a69b991b47abdbbb31c0da7dd3ce4db887d94efaa0b21632", "network": "host", "args": ["serve"], "volumes": ["mesh-broker-tls:/broker-tls:ro"], diff --git a/internal/apply/apply.go b/internal/apply/apply.go index 8681670..45e8a9e 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -112,8 +112,13 @@ func Apply( log(fmt.Sprintf(" %s %s (%s)", action, orphan.ID, orphan.Target)) } + // What moved in this apply, so a service that must reflect a file can be told the file + // moved. Only within one apply: a change from an earlier one has already been reflected, and + // restarting for it every time would make a steady machine restart its services for ever. + changed := map[string]bool{} + for _, resource := range d.Resources { - outcome, err := applyOne(ctx, sys, resource, run) + outcome, err := applyOne(ctx, sys, resource, run, changed) if err != nil { return report, known, &Error{Resource: resource.Identity(), Err: err, Done: report} } @@ -126,20 +131,22 @@ func Apply( }) report.Outcomes = append(report.Outcomes, outcome) if outcome.Action != "unchanged" { + changed[resource.Identity()] = true log(fmt.Sprintf(" %s %s (%s)", outcome.Action, outcome.ID, outcome.Target)) } } return report, known, nil } -func applyOne(ctx context.Context, sys system.System, r declaration.Resource, run Runner) (Outcome, error) { +func applyOne(ctx context.Context, sys system.System, r declaration.Resource, run Runner, + changed map[string]bool) (Outcome, error) { switch res := r.(type) { case *declaration.Directory: return applyDirectory(res) case *declaration.File: return applyFile(res) case *declaration.Service: - return applyService(ctx, sys, res, run) + return applyService(ctx, sys, res, run, changed) case *declaration.Package: return applyPackage(ctx, sys, res, run) case *declaration.Container: @@ -318,7 +325,25 @@ func writeAtomically(path string, content []byte, mode os.FileMode) error { return os.Rename(tmp.Name(), path) } -func applyService(ctx context.Context, sys system.System, r *declaration.Service, run Runner) (Outcome, error) { +// reflects reports whether anything this service must mirror changed in this apply. +func reflects(r *declaration.Service, changed map[string]bool) bool { + return len(reflected(r, changed)) > 0 +} + +// reflected is which of them changed, so the outcome can say why the service was restarted. A +// restart with no reason given is indistinguishable from a service that keeps falling over. +func reflected(r *declaration.Service, changed map[string]bool) []string { + var which []string + for _, id := range r.RestartOn { + if changed[id] { + which = append(which, id) + } + } + return which +} + +func applyService(ctx context.Context, sys system.System, r *declaration.Service, run Runner, + changed map[string]bool) (Outcome, error) { out := begin(r) var changes []string @@ -365,6 +390,28 @@ func applyService(ctx context.Context, sys system.System, r *declaration.Service return out, fmt.Errorf("%s was asked to be %s and is %s", r.Unit, r.State, after) } changes = append(changes, before+" to "+after) + } else if r.State == "running" && reflects(r, changed) { + // The service is already in the state it was asked for, and something it must reflect + // changed in this same apply. A running service does not re-read its configuration, so + // leaving it alone here is how a machine ends up correct on disk and wrong in fact — + // with every check passing. + if err := sys.SetServiceState(ctx, run, r.Unit, "stopped"); err != nil { + return out, fmt.Errorf("restarting %s: stopping it: %w", r.Unit, err) + } + if err := sys.SetServiceState(ctx, run, r.Unit, "running"); err != nil { + return out, fmt.Errorf("restarting %s: starting it again: %w", r.Unit, err) + } + // Read back, for the same reason as above: a unit that starts and immediately dies + // satisfies a service manager and nothing else. + after, err := sys.ServiceState(ctx, run, r.Unit) + if err != nil { + return out, err + } + if after != "running" { + return out, fmt.Errorf( + "%s was restarted to pick up a change and is %s", r.Unit, after) + } + changes = append(changes, "restarted for "+strings.Join(reflected(r, changed), ", ")) } if len(changes) == 0 { diff --git a/internal/apply/apply_test.go b/internal/apply/apply_test.go index 0181cc3..0dfbb4b 100644 --- a/internal/apply/apply_test.go +++ b/internal/apply/apply_test.go @@ -3,6 +3,7 @@ package apply import ( "context" "errors" + "fmt" "os" "path/filepath" "strings" @@ -943,3 +944,118 @@ func archHost(t *testing.T) system.System { } return s } + +func TestAServiceIsRestartedWhenWhatItReflectsChanges(t *testing.T) { + // A running service does not re-read its configuration. Replace the file, find the service + // already running, do nothing — and the machine keeps behaving as it did while every check + // passes, because the file is right and the service is up. + // + // That is how a third node joining a mesh left the first two carrying a network that no + // longer existed. Found in the lab; this is the shape of the fix. + dir := t.TempDir() + path := filepath.Join(dir, "thing.conf") + + d := parse(t, fmt.Sprintf(`{"declaration":1,"resources":[ + {"id":"conf","type":"file","path":%q,"content":"first\n","mode":"0644"}, + {"id":"svc","type":"service","unit":"thing.service","state":"running","restart-on":["conf"]} + ]}`, path)) + + var commands []string + run := recordingServices(&commands) + + if _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, + store.OriginCarried, run, nil); err != nil { + t.Fatal(err) + } else { + // Second apply with the same content: nothing moved, so nothing restarts. A machine that + // restarted its services on every reconcile would never be steady. + commands = nil + if _, _, err := Apply(context.Background(), archHost(t), d, state, + store.OriginCarried, run, nil); err != nil { + t.Fatal(err) + } + for _, c := range commands { + if strings.Contains(c, "stop") { + t.Errorf("an unchanged declaration restarted the service: %s", c) + } + } + + // Now the file changes. The service is already running and must still be restarted. + changedDecl := parse(t, fmt.Sprintf(`{"declaration":1,"resources":[ + {"id":"conf","type":"file","path":%q,"content":"second\n","mode":"0644"}, + {"id":"svc","type":"service","unit":"thing.service","state":"running","restart-on":["conf"]} + ]}`, path)) + commands = nil + if _, _, err := Apply(context.Background(), archHost(t), changedDecl, state, + store.OriginCarried, run, nil); err != nil { + t.Fatal(err) + } + var stopped, started bool + for _, c := range commands { + if strings.Contains(c, "stop thing.service") { + stopped = true + } + if strings.Contains(c, "start thing.service") { + started = true + } + } + if !stopped || !started { + t.Errorf("the file changed and the service was not restarted; commands were %v", commands) + } + } +} + +func TestAServiceIsNotRestartedByAChangeItDoesNotName(t *testing.T) { + // The list is what it reflects, not everything in the declaration. A service restarted by any + // change anywhere would make every apply a fleet-wide bounce. + dir := t.TempDir() + conf := filepath.Join(dir, "thing.conf") + other := filepath.Join(dir, "unrelated") + + first := parse(t, fmt.Sprintf(`{"declaration":1,"resources":[ + {"id":"conf","type":"file","path":%q,"content":"same\n","mode":"0644"}, + {"id":"other","type":"file","path":%q,"content":"one\n","mode":"0644"}, + {"id":"svc","type":"service","unit":"thing.service","state":"running","restart-on":["conf"]} + ]}`, conf, other)) + + var commands []string + run := recordingServices(&commands) + _, state, err := Apply(context.Background(), archHost(t), first, store.State{}, + store.OriginCarried, run, nil) + if err != nil { + t.Fatal(err) + } + + second := parse(t, fmt.Sprintf(`{"declaration":1,"resources":[ + {"id":"conf","type":"file","path":%q,"content":"same\n","mode":"0644"}, + {"id":"other","type":"file","path":%q,"content":"two\n","mode":"0644"}, + {"id":"svc","type":"service","unit":"thing.service","state":"running","restart-on":["conf"]} + ]}`, conf, other)) + commands = nil + if _, _, err := Apply(context.Background(), archHost(t), second, state, + store.OriginCarried, run, nil); err != nil { + t.Fatal(err) + } + for _, c := range commands { + if strings.Contains(c, "stop") { + t.Errorf("a change to a file the service does not name restarted it: %s", c) + } + } +} + +// recordingServices answers the way a machine with a running unit would, and remembers what it +// was asked to do — which is what a restart has to be proved by, since "running" looks the same +// before and after one. +func recordingServices(commands *[]string) Runner { + return func(_ context.Context, name string, args ...string) (string, error) { + line := name + " " + strings.Join(args, " ") + *commands = append(*commands, line) + switch { + case strings.Contains(line, "is-enabled"): + return "enabled", nil + case strings.Contains(line, "show") && strings.Contains(line, "ActiveState"): + return "LoadState=loaded\nActiveState=active\nSubState=running", nil + } + return "", nil + } +} diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index 2a504db..2c79fd2 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -113,6 +113,20 @@ type Service struct { // Without this the host could start a unit and not make it survive a reboot, which is a // declaration that reports success and stops being true at the next power cut. Boot string `json:"boot,omitempty"` + + // RestartOn names resources whose change means this service must be restarted. + // + // Because a running service does not re-read its configuration. Replace the file, find the + // service already running, do nothing, and the machine keeps behaving the way it did before — + // while every check passes, because the file is right and the service is up. That is not + // hypothetical: it is how a third node joining a mesh left the first two carrying a network + // that no longer existed, and every part of it reported success. + // + // This is declared state rather than a command. The declaration says the running service must + // reflect these files; the host works out that it does not and acts. A *command* to restart + // would be an action, and the link may not carry one (novox/hq ADR 0005) — so this is not a + // way around that rule, it is the shape the rule leaves. + RestartOn []string `json:"restart-on,omitempty"` } func (s *Service) Identity() string { return s.ID } diff --git a/internal/link/enrol.go b/internal/link/enrol.go index ed18f2e..cb64a60 100644 --- a/internal/link/enrol.go +++ b/internal/link/enrol.go @@ -23,10 +23,19 @@ func QueueFor(node string) string { return "node." + node } // EnrolRequest is what this node says when joining. type EnrolRequest struct { - Node string `json:"node"` - Secret string `json:"secret"` - PublicKey []byte `json:"public_key"` - Profile map[string]any `json:"profile,omitempty"` + Node string `json:"node"` + Secret string `json:"secret"` + PublicKey []byte `json:"public_key"` + + // OverlayKey is the public half of this node's key on the private network — a different key + // from PublicKey above, generated at the same moment and for a different purpose. + // + // Sent with enrolment because the overlay is the first declaration a node receives, and the + // mesh cannot compose it without this. Asking for it afterwards would mean a node is enrolled + // and unreachable for a round trip, which is the state everything else here works to avoid. + OverlayKey string `json:"overlay_key,omitempty"` + + Profile map[string]any `json:"profile,omitempty"` } // EnrolReply is what the mesh says back. @@ -56,7 +65,7 @@ var ErrRefused = errors.New("the mesh refused this enrolment") // says once it is in, and the secret travels again because the control plane must not have to ask // the broker who connected. func Enrol(ctx context.Context, address, pin, node, secret string, public []byte, - profile map[string]any, timeout time.Duration) (EnrolReply, error) { + overlayKey string, profile map[string]any, timeout time.Duration) (EnrolReply, error) { config, err := PinnedConfig(pin) if err != nil { @@ -101,7 +110,8 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte return EnrolReply{}, err } - request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, Profile: profile} + request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, + OverlayKey: overlayKey, Profile: profile} body, err := json.Marshal(request) if err != nil { return EnrolReply{}, err