Unify trunk on main: initialization → main #3

Merged
jschoubben merged 58 commits from initialization into main 2026-09-05 01:13:33 +00:00
6 changed files with 223 additions and 14 deletions
Showing only changes of commit 1bc97ed50d - Show all commits
+18 -1
View File
@@ -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)
+8 -3
View File
@@ -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"],
+51 -4
View File
@@ -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 {
+116
View File
@@ -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
}
}
+14
View File
@@ -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 }
+16 -6
View File
@@ -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