Unify trunk on main: initialization → main #3

Merged
jschoubben merged 58 commits from initialization into main 2026-09-05 01:13:33 +00:00
4 changed files with 96 additions and 7 deletions
Showing only changes of commit 4bff67ec69 - Show all commits
+46 -5
View File
@@ -535,14 +535,20 @@ func firstNonEmpty2(values ...[]byte) []byte {
// is down it is disconnected, which is an ordinary situation rather than a failure — the machine
// keeps running whatever it was last told, from its own store.
func runLink(ctx context.Context, opts options) error {
mine, err := identity.Load(identity.Path(opts.state))
if errors.Is(err, identity.ErrNoIdentity) {
return errors.New("this machine has not joined a mesh. Enrol it first: " +
"mesh-host enrol --token <token> --name <name>")
}
// A machine that has not enrolled waits here rather than failing. It is *hosted*: the host is
// running, it has no identity, and there is nobody to link to — an ordinary state, and the
// one every machine passes through (novox/hq 09-the-node-lifecycle).
//
// Exiting instead would be worse than untidy. The launcher counts a failed start, and three
// of them roll the binary back — so a freshly installed host, waiting to be enrolled exactly
// as intended, would undo its own installation.
mine, err := waitForEnrolment(ctx, opts)
if err != nil {
return err
}
if mine.Node == "" {
return nil // asked to stop while waiting
}
fmt.Printf("node %s, linking to %s\n", mine.Node, mine.Membership.Broker)
@@ -668,3 +674,38 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D
}
return report
}
// waitForEnrolment returns this node's identity, waiting for one if it has none.
//
// It applies the carried bundle first, if there is one, because that is what a first node does
// before there is a mesh at all — and a machine that has been installed and not yet enrolled
// should still be whatever its bundle says it is.
func waitForEnrolment(ctx context.Context, opts options) (identity.Identity, error) {
const look = 5 * time.Second
said := false
for {
mine, err := identity.Load(identity.Path(opts.state))
if err == nil {
return mine, nil
}
if !errors.Is(err, identity.ErrNoIdentity) {
// An identity that exists and cannot be read is a fault, not a wait. Treating it as
// "not enrolled yet" would leave a node sitting quietly for ever while the mesh
// believes it is a member.
return identity.Identity{}, err
}
if !said {
fmt.Println("this machine has not joined a mesh, and is waiting to be told which one.")
fmt.Println(" enrol it with: mesh-host enrol --token <token> --name <name>")
said = true
}
select {
case <-ctx.Done():
return identity.Identity{}, nil
case <-time.After(look):
}
}
}
+2 -2
View File
@@ -79,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:1900893c8d175f1f955569dbead167e58ec3a593fc2ad9170fd936e0a8a55c3a",
"192.0.2.250:5000/mesh-control@sha256:5aaaea5f3d0daa7a4cdde4e13dbddc61176d46487e6a89e6d992808e5e771593",
"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"]
},
@@ -113,7 +113,7 @@
"id": "control-plane",
"type": "container",
"name": "mesh-control",
"image": "192.0.2.250:5000/mesh-control@sha256:1900893c8d175f1f955569dbead167e58ec3a593fc2ad9170fd936e0a8a55c3a",
"image": "192.0.2.250:5000/mesh-control@sha256:5aaaea5f3d0daa7a4cdde4e13dbddc61176d46487e6a89e6d992808e5e771593",
"network": "host",
"args": ["serve"],
"volumes": ["mesh-broker-tls:/broker-tls:ro"],
+14
View File
@@ -8,8 +8,22 @@ package link
// queue, so it can say these things and nothing else.
const (
KeyReport = "report"
KeyAlive = "alive"
)
// Alive is a node saying nothing except that it is there.
//
// novox/hq 09-the-node-lifecycle: *how long it has been disconnected is a fact the mesh must
// hold, and nothing holds it today. Without it, a node running last month's assignments looks
// exactly like one that is current.*
//
// Separate from a report because the two happen at completely different rates: a node is alive
// constantly and applies something rarely, and reading one as the other would make a quiet node
// look like a stale one.
type Alive struct {
Node string `json:"node"`
}
// Signed is a declaration and the signature over it.
//
// novox/hq ADR 0004: the transport is verified once at connect, and **each declaration is
+34
View File
@@ -19,6 +19,13 @@ import (
// the first means somebody is trying, the second means something is broken.
var ErrForged = errors.New("this declaration was not signed by the mesh this node joined")
// AliveEvery is how often a node says it is there.
//
// Often enough that "no word for five minutes" means something, rarely enough that a hundred
// nodes are not a hundred messages a second. The mesh reads absence rather than presence, so what
// matters is the interval being known and steady.
const AliveEvery = 60 * time.Second
// Membership is what a node needs to reach its mesh again, held by the caller.
type Membership struct {
Node string
@@ -153,6 +160,13 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
// reads as still broken when it means the opposite.
say("in the mesh, consuming " + queue)
// A word every so often, so the mesh can tell a node that is quiet from one that is gone.
// Cheap on purpose: it carries a name and nothing else, because anything more would be a
// report, and reports are rare where this is constant.
beat := time.NewTicker(AliveEvery)
defer beat.Stop()
publishAlive(ctx, channel, m, say, timeout)
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
// Published mandatory, so the broker hands back anything it cannot route rather than
@@ -171,6 +185,8 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
select {
case <-ctx.Done():
return nil
case <-beat.C:
publishAlive(ctx, channel, m, say, timeout)
case reason := <-closed:
return fmt.Errorf("the link closed: %v", reason)
case delivery, ok := <-deliveries:
@@ -236,3 +252,21 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep
say(fmt.Sprintf("applied, and could not tell the mesh: %v", err))
}
}
// publishAlive says this node is here, and nothing else.
func publishAlive(ctx context.Context, channel *amqp.Channel, m Membership, say Announce,
timeout time.Duration) {
body, err := json.Marshal(Alive{Node: m.Node})
if err != nil {
return
}
publish, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
// Not mandatory, unlike a report. Losing one is nothing: the next is a minute away, and the
// mesh is reading a gap rather than counting arrivals. Insisting on delivery would turn a
// harmless miss into a logged failure every minute.
if err := channel.PublishWithContext(publish, Exchange, KeyAlive, false, false,
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
say("could not tell the mesh this node is here: " + err.Error())
}
}