The bus on NATS: both transports behind seams, and the rollout switch #87

Merged
jschoubben merged 40 commits from feat/nats-genesis into main 2026-09-27 17:36:41 +00:00
8 changed files with 340 additions and 23 deletions
Showing only changes of commit eb72ec36ba - Show all commits
+9 -3
View File
@@ -86,13 +86,19 @@ proxy-image:
# The whole gate. Raises a database, runs everything against it, and takes it down again --
# including when the tests fail, which is why the teardown is not conditional.
#
# **One package at a time (-p 1), and it is not about speed.** The live tests reach one bus, and on
# it they assert, read and remove the mesh's own objects -- streams and consumers with fixed names,
# because those names are the mesh's and a test cannot choose others. Two packages doing that at once
# is one deleting a consumer the other is reading through, and the failure lands in whichever test
# was reading, as "no response from stream". That reads as a bug in the code under test.
check: fmt vet postgres
@go test ./... ; status=$$? ; $(MAKE) postgres-stop ; exit $$status
@go test -p 1 ./... ; status=$$? ; $(MAKE) postgres-stop ; exit $$status
# Without a database the live tests skip rather than fail, so this is the honest subset and not
# the gate.
# the gate. Serialised for the same reason check is: a bus may be configured even when a store is not.
test:
go test ./...
go test -p 1 ./...
vet:
go vet ./...
+89
View File
@@ -257,6 +257,21 @@ func moduleCommand(ctx context.Context, args []string) error {
return err
}
// Which bus this mesh is on. A module gets a credential for exactly one, and the two are
// made in entirely different ways: on the bus the mesh runs on today an account is a
// management call, and on the bus being built it is a row the next composition writes into
// the server's user list (novox/hq design 25 §4).
busAddress, onNATS, err := broker.OnNATS()
if err != nil {
return err
}
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
return err
}
if onNATS {
return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress)
}
management, err := broker.ManagementFromEnvironment()
if err != nil {
return err
@@ -567,3 +582,77 @@ func mayIssue(m catalogue.Manifest) error {
}
return nil
}
// issueOnTheNewBus gives an assigned module its credential on the bus being built.
//
// **Three things differ from a management call, and each is the point of the move.** The credential
// is minted into the mesh's records and becomes usable at the next composition, so there is no
// server to be reachable for this to work. The password travels beside the address rather than inside
// it, because the runtime's contract already separates them and a credential embedded in a URL is one
// that leaks into every log line that prints a connection. And the module's durable consumer is
// derived from what it declared rather than declared by name, so a module cannot ask for delivery of
// something it did not say it consumes.
func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue.Manifest,
node, busAddress string) error {
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username()
password, err := inv.MintBusPassword(ctx, inventory.BusUser{
Username: user, Kind: inventory.BusModule, Node: node, Module: m.Module,
})
if err != nil {
return err
}
// Where the module is told to find the bus, and what certificate it must present. The same pair
// a node is told, for the same reason: a mesh's bus presents its own certificate, in no public
// trust store, so an address alone fails at TLS.
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("cannot deliver a credential without knowing where the bus is: %w", err)
}
reachable, err := brokerReachableAt(ctx, inv, known, node)
if err != nil {
return err
}
held, err := json.Marshal(struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
Node string `json:"node"`
Module string `json:"module"`
User string `json:"user"`
Password string `json:"password"`
}{
URL: "nats://" + reachable, Fingerprint: known.Fingerprint,
Node: node, Module: m.Module, User: user, Password: password,
})
if err != nil {
return err
}
if err := inv.AcceptSecretForModule(ctx, node, m.Module, "broker", string(held)); err != nil {
return err
}
// And how it hears what it consumes. Derived from its declaration, and only when it declared
// something: a module that consumes nothing needs no consumer, and creating one would be a
// durable subscription nobody reads.
if consumer, needed := broker.ConsumerFor(broker.Principal{
Kind: broker.KindModule, Node: node, Module: m.Module,
Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools,
}); needed {
js, err := broker.Dial(busAddress)
if err != nil {
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
"how %s hears what it consumes: %w", m.Module, err)
}
defer js.Close()
if err := js.EnsureConsumer(consumer); err != nil {
return err
}
}
fmt.Printf("bus user %s minted for %s, scoped to what it emits and consumes\n", user, m.Module)
fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n", node, node)
fmt.Printf(" and it works once the bus has been told: the user list is composed into the " +
"machine holding mesh-broker\n")
return nil
}
+71
View File
@@ -585,12 +585,19 @@ func renderingFor(ctx context.Context, open *stores, node string,
if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err
}
// The bus's user list, for the machine that runs the bus. Composed per push rather than kept,
// because it is a function of the mesh's records and a kept copy could disagree with them.
busUsers, err := composeBusUsers(ctx, inv, plan.Modules)
if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err
}
return catalogue.Rendering{
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
Certificate: certificate, Authority: authority, Mesh: private, Names: names,
Machines: machines,
Suffix: overlay.Suffix(), Foundation: foundation, Kept: kept, Adopted: record.Adopted,
Given: given, Taken: taken, Seats: seats, ArtifactStore: artifactStore, Built: built,
BusUsers: busUsers,
}, record, nil
}
@@ -1135,3 +1142,67 @@ func portsOn(
}
return out, nil
}
// composeBusUsers is the bus's user list, for a push to the machine that runs the bus.
//
// Empty for every other machine, and for every machine while the mesh is on the bus it runs on
// today — where accounts are a management call and there is no file to write.
//
// **Composed on each push, never kept.** The list is a function of the mesh's records (who exists,
// what runs where, what each declares), and a stored copy would be a second account of who may reach
// the bus, able to disagree with the records while both looked internally consistent (ADR 0043).
//
// A user the mesh has never minted a password for is **left out and said**, not written as a user
// without one — the composer refuses that, because a user with no password is a user anybody is. That
// is an ordinary situation with an obvious remedy (`module issue`, or enrolling), so the push carries
// the rest rather than failing: a bus that is missing one module's user is a mesh where that module
// cannot connect, and a bus with no file at all is a mesh where nothing can.
func composeBusUsers(ctx context.Context, inv *inventory.Inventory,
onThisNode []catalogue.Manifest) (string, error) {
_, onNATS, err := broker.OnNATS()
if err != nil || !onNATS {
return "", err
}
// Only for the machine holding the bus. Asked of what this push resolves to rather than of the
// seat's holder mesh-wide: the file is a resource of that module, so the question is whether it
// is here.
holdsTheBus := false
for _, m := range onThisNode {
if m.BusUsers != "" && m.ClaimsSeat("mesh-broker") {
holdsTheBus = true
}
}
if !holdsTheBus {
return "", nil
}
records, err := inv.BusRecords(ctx)
if err != nil {
return "", err
}
users, err := broker.Users(records)
if err != nil {
return "", err
}
kept, err := inv.BusUsers(ctx)
if err != nil {
return "", err
}
hashes := make(map[string]string, len(kept))
for name, u := range kept {
hashes[name] = u.PasswordHash
}
filled, missing := broker.WithPasswords(users, hashes)
if len(missing) > 0 {
fmt.Printf("the bus's user list leaves out %d user(s) the mesh has minted no credential "+
"for: %s. Each is a user that cannot connect until one is issued\n",
len(missing), strings.Join(missing, ", "))
}
if len(filled) == 0 {
return "", fmt.Errorf(
"this machine runs the bus and not one user has a credential, so the composed list " +
"would refuse every connection in the mesh")
}
return broker.ComposeAccounts(filled)
}
+2 -1
View File
@@ -88,7 +88,8 @@ func serve(ctx context.Context) error {
return err
}
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known}
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
OnNATS: onNATS}
server, err := link.Connect(work, work)
if err != nil {
return err
+6 -9
View File
@@ -28,15 +28,12 @@ func aLiveBus(t *testing.T) *JetStream {
t.Fatal(err)
}
t.Cleanup(js.Close)
// Deleted before, so what this test asserts is what it finds — and after, so the next test does
// not inherit it. Deleting a stream takes its consumers with it, which is why this is enough.
clear := func() {
for _, s := range MeshStreams() {
_ = js.Context().DeleteStream(s.Name)
}
}
clear()
t.Cleanup(clear)
// **Nothing is deleted here, deliberately.** These objects are the mesh's own and every live
// test in every package shares one server: a test that deleted a stream to get a clean slate
// took it out from under whatever was running beside it, and the failure landed in the other
// test as "stream not found" — which reads as a bug in the code under test. Raise is idempotent
// by requirement, so asserting against whatever is already there is both safe and the realistic
// case.
return js
}
+107
View File
@@ -0,0 +1,107 @@
package link_test
import (
"crypto/ed25519"
"crypto/rand"
"testing"
"time"
"golang.org/x/crypto/bcrypt"
"github.com/novox/mesh-controller/internal/identity"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// What a node is given to come back with, on the bus being built.
//
// **The credential becomes usable at the next composition, not when it is made**, which is the one
// real difference from the bus the mesh runs on today: there a management call makes it live at once.
// So what has to be true here is that the mesh recorded it and told the node, and the rest is a push.
// aMeshReadyToEnrol is both stores with a signing key established, which a control plane does at
// start: one that cannot sign is one whose declarations every node correctly refuses.
func aMeshReadyToEnrol(t *testing.T) (*inventory.Inventory, *identity.Identity) {
t.Helper()
inv := inventory.ForTest(t)
ident := identity.ForTest(t)
if _, err := ident.Establish(t.Context()); err != nil {
t.Fatal(err)
}
return inv, ident
}
func aTokenFor(t *testing.T, inv *inventory.Inventory, node string) (string, ed25519.PublicKey) {
t.Helper()
ctx := t.Context()
if _, err := inv.AddNode(ctx, node); err != nil {
t.Fatal(err)
}
issued, err := inv.IssueToken(ctx, node, time.Hour)
if err != nil {
t.Fatal(err)
}
public, _, err := ed25519.GenerateKey(rand.Reader)
if err != nil {
t.Fatal(err)
}
return issued.Secret, public
}
// A node enrolling onto the bus being built is told a password of its own, and the mesh keeps only
// its hash — which is what the next composition writes into the bus's user list.
func TestANodeEnrollingOnTheNewBusIsMintedACredentialTheMeshOnlyHashes(t *testing.T) {
inv, ident := aMeshReadyToEnrol(t)
ctx := t.Context()
secret, public := aTokenFor(t, inv, "anchor")
reply, err := link.Enrolment{Inventory: inv, Identity: ident, OnNATS: true}.Enrol(ctx, link.EnrolRequest{
Node: "anchor", Secret: secret, PublicKey: public})
if err != nil {
t.Fatal(err)
}
if reply.Password == "" {
t.Fatal("the node was told no password, so it keeps a one-time secret as a credential")
}
if reply.Password == secret {
t.Fatal("the node was handed the token's own secret back: a credential that lives for years " +
"must not be the string that was pasted into a terminal")
}
// Recorded under the name the composed file will use, and as a hash: a credential recoverable
// from the mesh's store is one whose blast radius is the store's.
hash, known, err := inv.BusUserHash(ctx, "node.anchor")
if err != nil || !known {
t.Fatalf("the mesh kept no credential for the node it enrolled: %v %v", known, err)
}
if hash == reply.Password {
t.Fatal("the store holds the password itself")
}
if err := bcrypt.CompareHashAndPassword([]byte(hash), []byte(reply.Password)); err != nil {
t.Fatalf("what the mesh kept does not verify what it told the node: %v", err)
}
}
// On the bus the mesh runs on today, with no management configured, nothing is minted and the node is
// told so by being given no password — it keeps the token's secret, which it says out loud.
//
// **This is the check that the switch is a switch.** A node enrolling on one bus must not come away
// with a credential for the other: it would be half-moved, and nothing anywhere would say which half.
func TestANodeEnrollingOnTheOldBusIsMintedNoCredentialForTheNewOne(t *testing.T) {
inv, ident := aMeshReadyToEnrol(t)
ctx := t.Context()
secret, public := aTokenFor(t, inv, "anchor")
reply, err := link.Enrolment{Inventory: inv, Identity: ident}.Enrol(ctx, link.EnrolRequest{
Node: "anchor", Secret: secret, PublicKey: public})
if err != nil {
t.Fatal(err)
}
if reply.Password != "" {
t.Fatalf("a node on the old bus was given a password from nowhere: %q", reply.Password)
}
if _, known, err := inv.BusUserHash(ctx, "node.anchor"); err != nil || known {
t.Fatalf("a node enrolling on the old bus was given a credential for the new one: %v %v",
known, err)
}
}
+36 -1
View File
@@ -28,6 +28,15 @@ type Enrolment struct {
Identity *identity.Identity
Management *broker.Management
Broker broker.Broker
// OnNATS says the mesh's own traffic is on the bus being built, so a node's credential is
// minted into the mesh's records and composed into the bus's user list rather than pushed
// through a management call (novox/hq design 25 §4).
//
// **One bus, and a node gets a credential for exactly one** — refused at start if the
// controller is told about both (broker.MustBeOneBus), because a node holding a credential for
// each is one that could be half-moved, and nothing would say which half.
OnNATS bool
}
// Enrol records what the node presented and spends the token.
@@ -159,7 +168,33 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol
// the spend and not before: a replaced password on an attempt that failed would be held by
// nobody. If the broker will not take it now, the enrolment still stands — the node keeps
// the token's secret as its password, which it is told, and which is said here.
if e.Management != nil {
switch {
case e.OnNATS:
// **Minted into the mesh's records, not pushed to a server.** The bus's users are a file
// the controller composes, so a credential becomes usable at the next composition rather
// than at the moment it is made — and the plaintext is returned once, here, and then exists
// only on the machine it was sealed to.
//
// The node reconnects as itself and may be refused until that composition reaches the
// machine running the bus. That is what the host's reconnect backoff is for and it is
// survivable by design (ADR 0004: disconnection is an ordinary situation); waiting for the
// push here would hold an enrolment open for as long as a declaration takes to apply.
password, err := e.Inventory.MintBusPassword(ctx, inventory.BusUser{
Username: broker.Principal{Kind: broker.KindNode, Node: node.Name}.Username(),
Kind: inventory.BusNode,
Node: node.Name,
})
if err != nil {
// Not fatal to the enrolment: the node is recorded and the token is spent, and a node
// that keeps the token's secret is told so. Said loudly, because until this is minted
// the machine has no credential of its own.
log.Printf("%s is enrolled and the mesh could not mint its bus credential, so it keeps "+
"the token's secret as its password: %v", node.Name, err)
} else {
reply.Password = password
}
case e.Management != nil:
password, err := freshPassword()
if err != nil {
return EnrolReply{}, err
+20 -9
View File
@@ -36,19 +36,30 @@ func aBus(t *testing.T) *broker.JetStream {
}
t.Cleanup(js.Close)
// The mesh's own streams and consumers, asserted the way the controller asserts them — and
// torn down after, so one test's held message is never another's surprise.
for _, s := range broker.MeshStreams() {
_ = js.Context().DeleteStream(s.Name)
}
// **Streams purged, consumers removed.** Both halves, and each was learned by getting it wrong.
//
// The streams are emptied rather than deleted and recreated, because delete-then-add is not a
// reset: the server's teardown races the creation, and a test then inherits the previous one's
// messages — which reads as a redelivery bug in the code under test.
//
// The consumers are removed, because deleting a stream used to take them with it and purging
// does not. A durable *push* consumer that survives between tests keeps pushing to a delivery
// subject the previous test's subscription has gone from: the messages count as delivered, go
// nowhere, and the next test waits out its timeout for an announcement the server believes it
// already sent. The controller recreates what it needs on start, so leaving none is correct.
if err := broker.AssertMeshStreams(js); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
for _, s := range broker.MeshStreams() {
_ = js.Context().DeleteStream(s.Name)
clean := func() {
for _, c := range broker.MeshConsumers() {
_ = js.Context().DeleteConsumer(c.Stream, c.Name)
}
})
for _, s := range broker.MeshStreams() {
_ = js.Context().PurgeStream(s.Name)
}
}
clean()
t.Cleanup(clean)
return js
}