Compare commits

..
Author SHA1 Message Date
jschoubben 81e76fa485 The controller follows the subject it decodes
The decoder named the forge's merge subject as the fourth thing followed and the
list was three long: every message that fell through to that switch panicked the
control plane (2026-09-28). The entry was written and lost between two attempts
at the same edit. A test now walks the list; the composed grants and the genesis
template carry the subject.
2026-09-28 03:13:57 +02:00
mesh-admin 12ed35e87d Merge pull request 'A merge on the forge builds what it moved, bases first' (#110) from feat/a-merge-on-the-forge-builds-what-it-moved into main 2026-09-28 00:59:04 +00:00
jschoubben 525f10b858 A merge on the forge builds what it moved, bases first
The controller follows the forge's merges (novox/hq 04-ISSUES/131). For each
module recorded as built from that repository and branch it records the move to
the merge commit and builds it — bases first, because a module built before the
module it stands on is built against the old one and reports success, and a base
that fails stops what stands on it. Nothing is pushed here: what a finished build
does to the machines running the module stays the upgrade's decision.

Two more things the same ordering gives: `build --behind` builds bases first, and
`build --on <module>` rebuilds everything that stands on a module — the rebuild a
changed base needs, which "behind" does not see because their sources did not
move.
2026-09-28 02:59:02 +02:00
mesh-admin 0547316cf2 Merge pull request 'rollout hand: a machine's membership, minted afresh and handed to an operator once' (#109) from feat/rollout-hand into main 2026-09-28 00:40:25 +00:00
jschoubben 9fe9b5349c Merge pull request 'A store row keeps its seat's protocol, and a holder may take work from its queue' (#108) from fix/store-seats-keep-their-protocol into main 2026-09-28 02:21:24 +02:00
jschoubben d1e488efaf Merge pull request 'A store row keeps its seat's protocol, and a holder may take work from its queue' (#108) from fix/store-seats-keep-their-protocol into main 2026-09-28 00:09:28 +00:00
jschoubben c5dc7e732a A store row keeps its seat's protocol, and a holder may take work from its queue
The seat table has name, scope, delivers and decision, and the protocol ADR 0129
gave a seat lives only in the compiled defaults; loading the rows dropped it, so
no role's work queue was ever raised and the first build submitted over the new
bus met "no response from stream". Until the table gains the columns, a row with
no protocol keeps the compiled one of its name. And the holder of a seat is
granted what taking work from its queue needs — asking about the worker consumer
it binds, and acknowledging on it — which the first machine to try was refused.

The control plane's own seat placeholders no longer include the old bus's port,
which the switch removed with the variable.
2026-09-28 02:08:56 +02:00
jschoubben 7efcccd013 Merge pull request 'The build machine takes work on the bus its credential names, and the work queue has a taker' (#107) from feat/the-build-machine-takes-work-on-nats into main 2026-09-27 23:59:56 +00:00
jschoubben 964285f08c The build machine takes work on the bus its credential names, and the work queue has a taker
Two halves of one gap the first build over the new bus met. The machine decided
its bus from a variable its container never received, so the credential the mesh
sealed to it went unread; a credential for the new bus names the bus by scheme and
carries user, password and fingerprint beside the address, and that is enough to
dial it, pinned. And the roles' work queues were raised with no holders, so the
consumer a machine binds to take work was never created: the holders are read
from the catalogue and the handover record, as the resolver reads them.
2026-09-28 01:59:20 +02:00
jschoubben 5698dda11f Merge pull request 'A principal may hear what its consumer delivers' (#106) from fix/a-principal-may-hear-its-consumer into main 2026-09-27 23:47:29 +00:00
jschoubben 6005a8471f A principal may hear what its consumer delivers
A push consumer delivers on _DELIVER.<its name>, and a client bound to it
subscribes exactly that. No principal was granted it, and the server refused
every one the first time it bound a consumer: the control plane, each machine,
and a module would have been next. Each kind is granted its own consumers'
delivery subjects and no other's. The line announcing the raised bus printed the
URL with the credential in it; the address alone now.
2026-09-28 01:46:16 +02:00
jschoubben e2ee0dfe98 Merge pull request 'The bus account has JetStream, and the control plane's client has its own inbox' (#105) from fix/the-bus-account-has-jetstream into main 2026-09-27 23:40:38 +00:00
jschoubben 70341cfbc7 The bus account has JetStream, and the control plane's client has its own inbox
Two refusals the first live connections met. A user in the MESH account was told
"JetStream not enabled for account" the first time it bound a consumer: with
accounts defined, JetStream is enabled per account, not only globally — the
account's setting, which the mesh owns, not the server's block, which it does not.
And the control plane's client used a random inbox prefix where it is granted
exactly _INBOX.<its user>.>, so the server's first answer could not reach it. The
prefix now follows from the user in the URL, for every principal that dials so.
2026-09-28 01:40:10 +02:00
jschoubben 77643aa3f4 Merge pull request 'The control plane pins the bus's certificate, and keeps its password out of errors' (#104) from fix/the-controller-pins-the-bus-certificate into main 2026-09-27 23:35:56 +00:00
jschoubben 1fd6194ff8 The control plane pins the bus's certificate, and keeps its password out of errors
The bus presents the mesh's own certificate, which names nothing a public verifier
accepts; the client verified by name and failed against a bus that was answering
("certificate is not valid for any names", 2026-09-28). It now pins the leaf's
fingerprint from MESH_BROKER_CERTIFICATE, as every host does. And a connection
error named the whole URL, password included — the address alone now.
2026-09-28 01:35:20 +02:00
jschoubben c37018fdd2 Merge pull request 'The control plane serves and pushes on the bus it is told to' (#103) from feat/the-controller-serves-on-nats into main 2026-09-27 23:30:52 +00:00
jschoubben 3907ea0db0 The control plane serves and pushes on the bus it is told to
The seams were there and nothing chose a side: serve, push, ask and build all
opened the old bus's connection and declared over its channel, whatever
MESH_BUS_NATS said. So the switch moved every host and left the control plane
unable to follow — "this control plane has no MESH_BROKER_AMQP" with the new bus
named and standing (2026-09-28). That was task 4.3 of design 28, still open.

One place now decides: connectLink reads the switch, refuses both buses named at
once, raises the new bus's streams and this controller's consumers when it is
handed the inventory, and opens the link over whichever bus it is on. Every
caller that sent a declaration or asked a tool through the old channel goes
through the server's bus instead, which the new transport has and the channel is
not. OverNats is that outbound: a declaration is a JetStream publish into the
node's own subject, an event is announced on the subject its name derives to, a
tool is request and reply on the module's tool subject.
2026-09-28 01:29:44 +02:00
jschoubben 40f5e9a41c Merge pull request 'A machine may bind its consumer' (#102) from fix/a-node-may-bind-its-consumer into main 2026-09-27 23:18:20 +00:00
jschoubben 2c2eb51878 Only CONSUMER.INFO was missing from a machine's grants; the rest was already there 2026-09-28 01:17:40 +02:00
jschoubben aa2d0b51ea Golden: a machine's user may bind its consumer, ack, and hear its inbox 2026-09-28 01:17:15 +02:00
jschoubben 64d154d9d7 A machine may bind its consumer and hear the answer
Binding to a consumer asks the server about it and hears the answer on the
client's inbox; hearing a declaration acknowledges it. A machine's user was granted
none of that and was refused the first time one dialled a permissioned server:
"this node cannot read its declarations". Its inbox is its own prefix, which the
host now sets.
2026-09-28 01:16:46 +02:00
jschoubben ffa390f916 Merge pull request 'The mint leaves the control plane's old-bus secret alone' (#101) from fix/mint-leaves-the-control-planes-old-secret-alone into main 2026-09-27 23:08:35 +00:00
jschoubben 4d62e6caf1 The mint leaves the control plane's old-bus secret alone
The control plane is a module too, and its broker secret is the old bus's
credential it is still using while the mint runs. Writing the new bus's blob there
cut the mesh off from its own old bus mid-move. Its new-bus credential is the
controller principal's bus secret; the module principal is skipped.
2026-09-28 01:08:11 +02:00
jschoubben 386ae676ca Merge pull request 'The control plane mounts the bus secret it reads' (#100) from fix/the-controller-mounts-its-bus-secret into main 2026-09-27 23:03:42 +00:00
jschoubben f8a9c3d6bc The control plane mounts the bus secret it reads
MESH_BUS_NATS_FILE named /run/secrets/bus and nothing put a file there: the
manifest binds each secret explicitly, and the switch added the secret and the
variable but not the bind. Found live — the control plane came up on the new bus
and could not read its own credential.
2026-09-28 01:03:18 +02:00
jschoubben 83671fae5f Merge pull request 'The network map resolves each machine with the seat holders on record' (#99) from fix/the-network-map-knows-the-holders into main 2026-09-27 22:56:55 +00:00
jschoubben 9b715524a2 The network map resolves each machine with the seat holders on record
Without them, a machine running the next holder of a seat beside the current one
resolves as two holders, is refused, and drops out of the map — and with it the
address every other machine composes for what it offers. Found live: the control
node vanished from the private network the moment the new bus was assigned beside
the old one, and nothing on any machine could be composed.
2026-09-28 00:55:39 +02:00
jschoubben e06fc1ed16 Merge pull request 'The control plane speaks the new bus' (#98) from switch/the-controller-speaks-nats into main 2026-09-27 22:38:27 +00:00
jschoubben 13d7c5c5dd Merge pull request 'The mint names the bus by bare host, and can mint again' (#97) from fix/mint-bare-host-and-again into main 2026-09-27 22:37:36 +00:00
jschoubben 7b02feaebb The mint names the bus by bare host, and can mint again
BareAddress adds a scheme where none was, so the host it yielded carried one and
every URL built from it carried two. Caught before a push: the host is now taken
with no scheme and no port, and every URL adds its own. `rollout mint --again`
mints every credential afresh for exactly this case — a mint that was wrong before
anything received it.
2026-09-28 00:37:10 +02:00
jschoubben bd10e2c695 Merge pull request 'The mint finds the bus on the hub' (#96) from fix/mint-finds-the-bus-on-the-hub into main 2026-09-27 22:34:41 +00:00
jschoubben 4b209d944d The mint finds the bus on the hub
The map of who is where on the private network lists the machines placed around
the hub, not the hub — and the control node is the hub, and runs the new bus.
Found on the first live run: refused for having no address. The address every
machine already dials the current bus at is the same machine, so that host with
the new port is what they are told.
2026-09-28 00:34:17 +02:00
jschoubben 84024cbdb6 The control plane speaks the new bus
One environment variable moves it (design 25): MESH_BUS_NATS, read from the `bus`
secret `rollout mint` sealed to its machine, and the two that named the old bus go,
because being told about both is refused at start. Registered and built ahead of
the push that flips it, so the flip and the bus it flips to arrive in one
declaration — the same push that starts the new bus and hands this machine its
membership. Deliberately not merged until every credential is minted.
2026-09-28 00:31:57 +02:00
jschoubben ad2eed2f71 Merge pull request 'The move mints every credential and tells each machine its membership' (#95) from feat/the-move-mints-and-delivers into main 2026-09-27 22:30:58 +00:00
jschoubben e8aa7ed9e7 The move mints every credential and tells each machine its membership
`rollout mint` gives every principal the new bus will have a credential it does
not yet have and puts each where its owner reads it: a machine's as a membership
— bus address, fingerprint, password, transport — sealed into its declaration
(migration 0041, the `bus-membership` resource the host reads after applying); a
module's as its broker secret, through the same delivery `module issue` uses; the
control plane's own as its `bus` secret. Idempotent, and worked out from where the
bus's module is assigned rather than from this process's environment, because this
process is still on the old bus when it runs and must be.

This is the half of design 28 task 5.2 the first live attempt found missing: a
credential was minted only at enrolment, at `module issue` and for a person, so no
machine already enrolled could ever be moved. `rollout check` was right to refuse;
now there is something to run first.
2026-09-28 00:16:24 +02:00
jschoubben 337aaea123 Merge pull request 'The first handover records the standing holder without re-judging it' (#93) from fix/record-the-standing-holder into main 2026-09-27 21:35:08 +00:00
jschoubben 4c41628b20 The first handover records the standing holder without re-judging it
On a mesh that predates the record, every handover has to begin by writing down who
already holds the seat — otherwise the next holder cannot be assigned beside it,
because two eligible claimants with nothing on record are refused. Found on the
live mesh minutes after 0040 moved the bus seat's row: the standing holder no
longer satisfies what the seat delivers, on purpose, and so could not be recorded,
and so nothing could stand beside it.

Recording who already holds is not making a new holder. Derivation never read
what the seat delivers, so the standing holder holds regardless; when nothing is
on record and the named assignment claims the seat at its scope, only that claim
is checked. Every change of holder is still judged in full.
2026-09-27 23:34:40 +02:00
jschoubben ae7fb520d7 Merge pull request 'AMQP is not a provision: the bus seat delivers the bus, and the word is refused' (#92) from feat/amqp-is-not-a-provision into main 2026-09-27 21:30:09 +00:00
jschoubben f325073982 AMQP is not a provision: the bus seat delivers the bus, and the word is refused
Two halves of novox/hq ADR 0131. Migration 0040 moves the mesh-broker row from
`amqp` to `mesh-bus`, so the seat's holder answers for the mesh's bus and not for
the wire protocol the old broker spoke — which is what let only the retiring
broker hold the seat that names the bus. Safe under the current holder: the
control plane composes its own address through the seat by name and the overview
derives holders by name; only registration and provision resolution read the
column. What must not happen in between is re-registering the current holder.

And the parser refuses a manifest that provides or requires `amqp`, each refusal
saying what to do instead: a module reaches the mesh's bus through the sdk and
depends on the seat, not on a protocol. A whole-catalogue test asserts nothing
beside this checkout names it; the three modules that did are removed there.

Two tests that used the old broker as a fixture now use the module that replaces
it or a manifest this package owns.
2026-09-27 23:29:00 +02:00
jschoubben 33c4e4be34 Merge pull request 'A seat is handed over as one act, and the holder is on record' (#91) from feat/seat-handover into main 2026-09-27 21:22:50 +00:00
jschoubben 8d52a2cfb0 Merge pull request 'The move ends with the old broker going, not staying' (#90) from fix/rollout-retires-the-old-broker into main 2026-09-27 21:05:47 +00:00
jschoubben 1cfe6be9c4 The move ends with the old broker going, not staying
`rollout check` said the old broker stays running as an ordinary provider of
amqp, and this was not its retirement. That was ADR 0127, which ADR 0131 has
superseded: AMQP is not a provision, so once every machine reports on the new bus
nothing of the mesh speaks to the old broker and its module is unassigned. The
plan says so, as its last step.

The flag that made the "it stays" line conditional is gone with the line — there
is no case in which the broker is kept. The test that pinned the opposite now
pins this, and says which record changed under it. The stale citation of a
record numbered 0119 is corrected while here.
2026-09-27 23:04:59 +02:00
42 changed files with 1216 additions and 135 deletions
+29 -3
View File
@@ -135,6 +135,16 @@ func run() error {
// machine told about both would take work from one and answer on the other, and every log line would // machine told about both would take work from one and answer on the other, and every log line would
// say it was fine. // say it was fine.
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) { func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
// **The credential decides, before any variable does.** A machine moved to the new bus was
// handed a credential for it and nothing else changed in its environment; that credential
// names the bus by scheme, so it is enough to know which bus to take work from.
if credential.onTheNewBus() {
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
if err != nil {
return nil, err
}
return link.MachineOverNATS(js, on), nil
}
address, onNATS, err := broker.OnNATS() address, onNATS, err := broker.OnNATS()
if err != nil { if err != nil {
return nil, err return nil, err
@@ -460,10 +470,26 @@ func brokerFrom() (Credential, error) {
// **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels // **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels
// out of band — here, sealed with the credential — and the endpoint is verified once at connect. // out of band — here, sealed with the credential — and the endpoint is verified once at connect.
type Credential struct { type Credential struct {
URL string `json:"url"` URL string `json:"url"`
// Fingerprint is SHA-256 over the broker certificate's DER bytes, or empty to verify the
// ordinary way.
Fingerprint string `json:"fingerprint,omitempty"` Fingerprint string `json:"fingerprint,omitempty"`
// User and Password ride beside the address on the bus being built (design 25): a credential
// embedded in a URL leaks into every log line that prints a connection, so the mesh seals them
// as two fields and this machine joins them once, here, to dial.
User string `json:"user,omitempty"`
Password string `json:"password,omitempty"`
}
// onTheNewBus is whether a credential is for the bus being built: its address says so, and the
// mesh only ever seals such a credential with the user and password beside it.
func (c Credential) onTheNewBus() bool { return strings.HasPrefix(strings.TrimSpace(c.URL), "nats://") }
// natsURL is the address with this machine's credential in it, for the one dial that needs it.
func (c Credential) natsURL() string {
rest := strings.TrimPrefix(strings.TrimSpace(c.URL), "nats://")
if c.User == "" {
return "nats://" + rest
}
return "nats://" + c.User + ":" + c.Password + "@" + rest
} }
// dial opens the connection, pinning the broker's certificate when there is one to pin. // dial opens the connection, pinning the broker's certificate when there is one to pin.
+2 -2
View File
@@ -14,8 +14,8 @@ func TestBuilderDiagnosticsStayOffStdout(t *testing.T) {
allowed := map[string]bool{ allowed := map[string]bool{
"string(body)": true, // once.go: the result JSON, which IS stdout "string(body)": true, // once.go: the result JSON, which IS stdout
"version)": true, // --version "version)": true, // --version
`"stopping")`: true, // the loop.s shutdown line `"stopping")`: true, // the loop.s shutdown line
"usage)": true, // --help text, for a human "usage)": true, // --help text, for a human
} }
for _, file := range []string{"once.go", "main.go"} { for _, file := range []string{"once.go", "main.go"} {
src, err := os.ReadFile(file) src, err := os.ReadFile(file)
+1 -1
View File
@@ -37,7 +37,7 @@ func askCommand(ctx context.Context, args []string) error {
arguments = json.RawMessage(positionals[2]) arguments = json.RawMessage(positionals[2])
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
+62 -3
View File
@@ -31,6 +31,53 @@ import (
// control plane may send a machine is bounded by the declaration language. This is the shape the // control plane may send a machine is bounded by the declaration language. This is the shape the
// builder module will take when it is given work over the broker; today a person runs it, and the // builder module will take when it is given work over the broker; today a person runs it, and the
// mesh records the result the same way either way. // mesh records the result the same way either way.
// buildOn rebuilds every module the mesh holds that stands on the named module's artifacts — the
// rebuild a changed base needs, which nothing else asks for: their sources did not move, and
// "behind" does not see a base that did (novox/hq 04-ISSUES/131). Bases first among them too.
func buildOn(ctx context.Context, base string, wait time.Duration) error {
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
held, err := open.inventory.Catalogued(ctx)
if err != nil {
return err
}
var on []inventory.Entry
for _, e := range held {
if e.Manifest.Build == nil {
continue
}
for _, b := range e.Manifest.Build.On {
if b.Module == base {
on = append(on, e)
break
}
}
}
if len(on) == 0 {
fmt.Printf("nothing the mesh holds stands on %s\n", base)
return nil
}
on = orderByBases(on)
fmt.Printf("%d module(s) stand on %s:\n", len(on), base)
var failed []string
for _, e := range on {
fmt.Printf("--- %s\n", e.Manifest.Module)
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, wait); err != nil {
fmt.Printf(" %v\n", err)
failed = append(failed, e.Manifest.Module)
}
}
if len(failed) > 0 {
return fmt.Errorf("%d of %d could not be built: %s", len(failed), len(on), strings.Join(failed, ", "))
}
fmt.Printf("\n%d module(s) rebuilt on %s. `push --behind` sends them on\n", len(on), base)
return nil
}
func buildCommand(ctx context.Context, args []string) error { func buildCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("build", flag.ContinueOnError) set := flag.NewFlagSet("build", flag.ContinueOnError)
ref := set.String("ref", "", "the branch, tag or commit to build") ref := set.String("ref", "", "the branch, tag or commit to build")
@@ -46,6 +93,7 @@ func buildCommand(ctx context.Context, args []string) error {
// retype each repository is asking them to be the loop. Naming a repository and asking which // retype each repository is asking them to be the loop. Naming a repository and asking which
// ones need building are different requests, so they are not combined. // ones need building are different requests, so they are not combined.
behind := set.Bool("behind", false, "every module the mesh holds older than its source has") behind := set.Bool("behind", false, "every module the mesh holds older than its source has")
on := set.String("on", "", "rebuild every module that stands on this module's artifacts — the rebuild a changed base needs")
// A repository on the mesh's own forge, named by its path there (novox/hq ADR 0111). Without it // A repository on the mesh's own forge, named by its path there (novox/hq ADR 0111). Without it
// the repository is external, cloned exactly as given — see source.go. // the repository is external, cloned exactly as given — see source.go.
self := set.Bool("self", false, "the repository is a path on the forge holding the git seat") self := set.Bool("self", false, "the repository is a path on the forge holding the git seat")
@@ -53,6 +101,13 @@ func buildCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err return err
} }
if *on != "" {
if len(positionals) != 0 || *behind || *self {
return errors.New("build --on <module> names a base and nothing else")
}
return buildOn(ctx, *on, *wait)
}
if *behind { if *behind {
if len(positionals) != 0 || *self { if len(positionals) != 0 || *self {
return errors.New("build <repository> or build --behind, not both: one names a " + return errors.New("build <repository> or build --behind, not both: one names a " +
@@ -61,7 +116,7 @@ func buildCommand(ctx context.Context, args []string) error {
return buildBehind(ctx, *wait) return buildBehind(ctx, *wait)
} }
if len(positionals) != 1 { if len(positionals) != 1 {
return errors.New("build <repository> [--self] [--path P] [--ref R] [--wait D] [--dry-run]") return errors.New("build <repository> [--self] [--path P] [--ref R] [--wait D] [--dry-run] | build --behind | build --on <module>")
} }
source := buildSource{Repository: positionals[0]} source := buildSource{Repository: positionals[0]}
if *self { if *self {
@@ -329,6 +384,10 @@ func buildBehind(ctx context.Context, wait time.Duration) error {
} }
fmt.Println() fmt.Println()
// Bases first: a module built before the module it stands on is built against the old one
// and reports success (novox/hq 04-ISSUES/131).
stale = orderByBases(stale)
var failed []string var failed []string
for _, e := range stale { for _, e := range stale {
fmt.Printf("--- %s\n", e.Manifest.Module) fmt.Printf("--- %s\n", e.Manifest.Module)
@@ -368,7 +427,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
} }
defer ident.Close() defer ident.Close()
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -472,7 +531,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
return err return err
} }
defer ident.Close() defer ident.Close()
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
+1
View File
@@ -184,6 +184,7 @@ func usage() {
operator key show the operator key, and what it can recover operator key show the operator key, and what it can recover
build <repository> [--ref R] have a build machine build it, and record what came out build <repository> [--ref R] have a build machine build it, and record what came out
build --behind build every module the mesh holds older than its source build --behind build every module the mesh holds older than its source
build --on <module> rebuild every module that stands on this module's artifacts, bases first
builds [<module>] what has been built lately, and what came of it builds [<module>] what has been built lately, and what came of it
builder issue <name> a broker account for a build machine, scoped to build work, builder issue <name> a broker account for a build machine, scoped to build work,
delivered as the builder module's broker secret (module add it first) delivered as the builder module's broker secret (module add it first)
+22 -8
View File
@@ -614,6 +614,15 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
if err != nil { if err != nil {
return err return err
} }
return issueWith(ctx, inv, m, node, busAddress, known, reachable, user, password)
}
// issueWith is the delivery half: the minted password sealed to the machine as the module's broker
// secret, and the module's consumer created where the bus can be reached. Split from the minting
// so the move can issue every module against a bus whose address it worked out itself
// (`rollout mint`, design 28 task 5.2) rather than the one in this process's environment.
func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manifest,
node, busAddress string, known broker.Broker, reachable, user, password string) error {
held, err := json.Marshal(struct { held, err := json.Marshal(struct {
URL string `json:"url"` URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"` Fingerprint string `json:"fingerprint,omitempty"`
@@ -639,14 +648,19 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
Kind: broker.KindModule, Node: node, Module: m.Module, Kind: broker.KindModule, Node: node, Module: m.Module,
Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools, Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools,
}); needed { }); needed {
js, err := broker.Dial(busAddress) if busAddress == "" {
if err != nil { fmt.Printf(" %s consumes; its consumer is created when the bus is reachable (`push`, then "+
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+ "`rollout mint` again is harmless)\n", m.Module)
"how %s hears what it consumes: %w", m.Module, err) } else {
} js, err := broker.Dial(busAddress)
defer js.Close() if err != nil {
if err := js.EnsureConsumer(consumer); err != nil { return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
return err "how %s hears what it consumes: %w", m.Module, err)
}
defer js.Close()
if err := js.EnsureConsumer(consumer); err != nil {
return err
}
} }
} }
+10 -1
View File
@@ -537,6 +537,15 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
if err != nil { if err != nil {
return nil, err return nil, err
} }
// **With the seat holders on record**, or a machine running the next holder of a seat beside
// the current one resolves as two holders, is refused, and drops out of the map — taking the
// address every other machine composes for what it offers (novox/hq ADR 0131). Found live:
// the control node vanished from the private network the moment the new bus was assigned
// beside the old one.
holdings, err := inv.Holdings(ctx)
if err != nil {
return nil, err
}
var out []inventory.Overlay var out []inventory.Overlay
for _, p := range places { for _, p := range places {
if p.Address == "" { if p.Address == "" {
@@ -549,7 +558,7 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
caps, _ := inv.ProfileOf(ctx, p.Name) caps, _ := inv.ProfileOf(ctx, p.Name)
got, err := catalogue.Resolve(shelf, assigned, got, err := catalogue.Resolve(shelf, assigned,
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps}, catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps},
catalogue.World{Unchecked: true}) catalogue.World{Unchecked: true, Holdings: holdings})
if err != nil { if err != nil {
continue continue
} }
+63
View File
@@ -0,0 +1,63 @@
package main
import (
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
func entry(module string, on ...string) inventory.Entry {
b := &catalogue.Build{}
for _, o := range on {
b.On = append(b.On, catalogue.BuildsOn{Arg: "X", Module: o, Artifact: "runtime"})
}
return inventory.Entry{Manifest: catalogue.Manifest{Module: module, Build: b}}
}
// A module built before the module it stands on is built against the old one and reports success
// (novox/hq 04-ISSUES/131). So bases come first, however the set arrived.
func TestBasesAreBuiltBeforeWhatStandsOnThem(t *testing.T) {
in := []inventory.Entry{entry("app", "runtime"), entry("runtime", "base"), entry("other"), entry("base")}
got := orderByBases(in)
pos := map[string]int{}
for i, e := range got {
pos[e.Manifest.Module] = i
}
if !(pos["base"] < pos["runtime"] && pos["runtime"] < pos["app"]) {
t.Fatalf("bases not first: %v", pos)
}
if len(got) != 4 {
t.Fatalf("an entry was lost or doubled: %d", len(got))
}
// A base outside the set is not waited for: it is not being rebuilt.
got = orderByBases([]inventory.Entry{entry("app", "elsewhere")})
if len(got) != 1 {
t.Fatalf("a dependency outside the set changed the set: %v", got)
}
}
// A merge names a repository the way the forge does; a source is recorded the way a build was
// asked for. The two meet on owner/repo and branch, whichever form the record took.
func TestAMergeMatchesTheSourcesBuiltFromIt(t *testing.T) {
m := link.SourceMoved{Owner: "novox", Repo: "mesh-controller", Base: "main",
CloneURL: "http://forge.internal:20000/novox/mesh-controller.git"}
for _, s := range []inventory.Source{
{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "main"},
{Repository: "novox/mesh-controller", Seat: "git", Ref: ""},
{Repository: "https://elsewhere.example/novox/mesh-controller", Ref: "main"},
} {
if !sourceIs(s, m) {
t.Errorf("%+v was not matched by the merge", s)
}
}
for _, s := range []inventory.Source{
{Repository: "novox/mesh-host", Seat: "git"},
{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "release"},
} {
if sourceIs(s, m) {
t.Errorf("%+v was matched by a merge that is not its", s)
}
}
}
+6 -1
View File
@@ -636,8 +636,13 @@ func renderingFor(ctx context.Context, open *stores, node string,
if err != nil { if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err return catalogue.Rendering{}, inventory.Node{}, err
} }
memberships, err := inv.BusMemberships(ctx)
if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err
}
return catalogue.Rendering{ return catalogue.Rendering{
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports, BusMembership: memberships[node],
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
Certificate: certificate, Authority: authority, Mesh: private, Names: names, Certificate: certificate, Authority: authority, Mesh: private, Names: names,
Machines: machines, Machines: machines,
Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation, Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation,
+80 -16
View File
@@ -38,6 +38,33 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
// nothing in it was wrong, and no one edit was the one that should have been a new file. // nothing in it was wrong, and no one edit was the one that should have been a new file.
// serve is the control plane running: one connection to the broker, one queue, one consumer. // serve is the control plane running: one connection to the broker, one queue, one consumer.
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
// so nothing served here finds them missing.
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
busAddress, onNATS, err := broker.OnNATS()
if err != nil {
return nil, err
}
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
return nil, err
}
if !onNATS {
return link.Connect(enroller, listener)
}
if inv != nil {
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
return nil, err
}
}
js, err := broker.Dial(busAddress)
if err != nil {
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
broker.BareAddress(busAddress), err)
}
return link.ConnectNats(js, enroller, listener), nil
}
func serve(ctx context.Context) error { func serve(ctx context.Context) error {
open, err := openStores(ctx) open, err := openStores(ctx)
if err != nil { if err != nil {
@@ -90,7 +117,7 @@ func serve(ctx context.Context) error {
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known, work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
OnNATS: onNATS} OnNATS: onNATS}
server, err := link.Connect(work, work) server, err := connectLink(ctx, inv, work, work)
if err != nil { if err != nil {
return err return err
} }
@@ -100,11 +127,6 @@ func serve(ctx context.Context) error {
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was // somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
// replaced all have records and no objects — and a node whose consumer is missing hears nothing // replaced all have records and no objects — and a node whose consumer is missing hears nothing
// while everything else about it looks correct. // while everything else about it looks correct.
if onNATS {
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
return err
}
}
// And build results nobody was waiting for. A build triggered any other way than `build` // And build results nobody was waiting for. A build triggered any other way than `build`
// would otherwise be reported into the void, which is the same as not reporting it. // would otherwise be reported into the void, which is the same as not reporting it.
server.Records(builds{inv}) server.Records(builds{inv})
@@ -158,13 +180,13 @@ func declare(ctx context.Context, args []string) error {
return err return err
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
defer server.Close() defer server.Close()
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
return err return err
} }
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
@@ -271,7 +293,7 @@ func pushCommand(ctx context.Context, args []string) error {
return err return err
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -332,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err return err
} }
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
// After it is away, not before. A digest recorded for something that failed to send would // After it is away, not before. A digest recorded for something that failed to send would
@@ -415,7 +437,7 @@ func pushCommand(ctx context.Context, args []string) error {
return declarationWith(held, open, node, plan, settings, gens, Allocating) return declarationWith(held, open, node, plan, settings, gens, Allocating)
}, },
func(s readyNode, body []byte) error { func(s readyNode, body []byte) error {
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, if err := link.Declare(ctx, server.Bus(), ident, s.node, body,
15*time.Second); err != nil { 15*time.Second); err != nil {
return err return err
} }
@@ -628,7 +650,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
len(refusals), strings.Join(refusals, "\n\n")) len(refusals), strings.Join(refusals, "\n\n"))
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -639,7 +661,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
if err != nil { if err != nil {
return err return err
} }
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
record, err := inv.NodeByName(ctx, s.node) record, err := inv.NodeByName(ctx, s.node)
@@ -704,7 +726,7 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
js, err := broker.Dial(address) js, err := broker.Dial(address)
if err != nil { if err != nil {
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w", return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
address, err) broker.BareAddress(address), err)
} }
defer js.Close() defer js.Close()
@@ -739,10 +761,52 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
// The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder, // The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder,
// deliberately: work queues until somebody arrives to do it, so assigning a build machine a week // deliberately: work queues until somebody arrives to do it, so assigning a build machine a week
// after something started asking for builds flushes the backlog instead of having lost it. // after something started asking for builds flushes the backlog instead of having lost it.
if err := broker.RaiseSeats(js, inventory.MeshSeats(), nil); err != nil { // With the seats' holders, so each role's work queue gets the consumer its holder takes
// work from. Passed as nil until the first live raise, which left the build machine bound to a
// consumer nothing had created (2026-09-28).
holders, err := seatHolders(ctx, inv)
if err != nil {
return err
}
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
return err return err
} }
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n", fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
address, len(names)) broker.BareAddress(address), len(names))
return nil return nil
} }
// seatHolders is who holds each of the mesh's seats, by seat name: the record where a handover
// wrote one, and the assigned module claiming the seat otherwise — the same derivation the
// resolver makes, read from the catalogue rather than re-resolved.
func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]broker.Holder, error) {
out := map[string]broker.Holder{}
entries, err := inv.Catalogued(ctx)
if err != nil {
return nil, err
}
for _, e := range entries {
if len(e.On) == 0 {
continue
}
for _, c := range e.Manifest.Claims {
seat, known := catalogue.SeatNamed(c.Name)
if !known {
continue
}
if _, taken := out[seat.Name]; !taken {
out[seat.Name] = broker.Holder{Node: e.On[0], Module: e.Manifest.Module}
}
}
}
recorded, err := inv.Holdings(ctx)
if err != nil {
return nil, err
}
for _, h := range recorded {
if seat, known := catalogue.SeatNamed(h.Claim); known {
out[seat.Name] = broker.Holder{Node: h.Node, Module: h.Module}
}
}
return out, nil
}
+265 -7
View File
@@ -2,8 +2,10 @@ package main
import ( import (
"context" "context"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"os"
"strings" "strings"
"time" "time"
@@ -12,6 +14,7 @@ import (
"github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/secrets"
) )
// Moving the mesh's own traffic to the bus being built (novox/hq ADR 0116 step 5). // Moving the mesh's own traffic to the bus being built (novox/hq ADR 0116 step 5).
@@ -25,18 +28,27 @@ import (
// against a mesh that is serving. It answers from records: what is missing, and what would happen. // against a mesh that is serving. It answers from records: what is missing, and what would happen.
// `rollout` itself refuses unless the check is clean. // `rollout` itself refuses unless the check is clean.
// //
// **The old broker is not switched off by this.** It stays an ordinary provider of `amqp` for whatever // **The old broker goes with the move, and goes last** (novox/hq ADR 0131): AMQP is not a provision,
// else uses it — on this installation, a whole automation layer that has nothing to do with the mesh // so once every machine reports on the new bus its module is unassigned. Only the mesh's own traffic
// ([ADR 0119](../../02-DECISIONS/0119-amqp-is-a-provision-not-the-bus.md)). Only the mesh's own // is what moves, which is why this is survivable at all: what breaks if it goes wrong is the mesh's
// traffic moves, which is why this is survivable at all: what breaks if it goes wrong is the mesh's // ability to change things, not the services its modules are serving — measured on 2026-09-27, when
// ability to change things, not the services its modules are serving. // a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
const rolloutUsage = "rollout check | rollout --confirm" const rolloutUsage = "rollout check | rollout mint [--again] | rollout hand <node> | rollout --confirm"
func rolloutCommand(ctx context.Context, args []string) error { func rolloutCommand(ctx context.Context, args []string) error {
switch { switch {
case len(args) == 1 && args[0] == "check": case len(args) == 1 && args[0] == "check":
return rolloutCheck(ctx) return rolloutCheck(ctx)
case len(args) == 1 && args[0] == "mint":
return rolloutMint(ctx, false)
case len(args) == 2 && args[0] == "hand":
return rolloutHand(ctx, args[1])
case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
// Every credential minted afresh, whether or not one exists — for a mint that was wrong
// before anything was pushed. Afterwards nothing that received the old one still works,
// which is fine exactly when nothing received it.
return rolloutMint(ctx, true)
case len(args) == 1 && args[0] == "--confirm": case len(args) == 1 && args[0] == "--confirm":
return errors.New( return errors.New(
"the rollout itself is not built yet: `rollout check` answers whether it could run, and " + "the rollout itself is not built yet: `rollout check` answers whether it could run, and " +
@@ -105,7 +117,6 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines
ModuleCredentialled: map[string]bool{}, ModuleCredentialled: map[string]bool{},
// The old broker keeps its other clients on this installation, and saying so is how the plan // The old broker keeps its other clients on this installation, and saying so is how the plan
// stops reading as a retirement. // stops reading as a retirement.
OldBusHasOtherClients: true,
} }
address, _, err := broker.OnNATS() address, _, err := broker.OnNATS()
@@ -191,3 +202,250 @@ func wasSentTheUserList(ctx context.Context, inv *inventory.Inventory, node stri
// notReadyOf is the readiness reasoning, named here so a test can reach it without the command's // notReadyOf is the readiness reasoning, named here so a test can reach it without the command's
// printing. The reasoning itself is the broker package's, where it is pure. // printing. The reasoning itself is the broker package's, where it is pure.
func notReadyOf(state broker.Readiness) []string { return broker.NotReady(state) } func notReadyOf(state broker.Readiness) []string { return broker.NotReady(state) }
// rolloutMint gives every principal the new bus will have a credential it does not yet have, and
// puts each where its owner reads it (novox/hq design 28, task 5.2): a machine's as a membership
// sealed into its declaration, a module's as its broker secret, the control plane's own as its
// `bus` secret. Idempotent: what already has a hash is left alone, so running it again is harmless.
//
// **Before anything moves, and it is what makes moving possible.** A machine moved without a
// credential cannot come back, and afterwards there is no bus to tell it anything over — which is
// why `rollout check` refuses until this has run. The bus's address is worked out here, from where
// the module that provides it is assigned, rather than read from this process's environment: this
// process is still on the old bus when this runs, and must be.
func rolloutMint(ctx context.Context, again bool) error {
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("the bus's certificate is not known to this process, and every membership "+
"must carry its fingerprint: %w", err)
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
return err
}
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
var busNode, controllerNode string
for _, e := range entries {
switch {
case e.Manifest.ClaimsSeat("mesh-broker") && providesBus(e.Manifest) && len(e.On) > 0:
busNode = e.On[0]
case e.Manifest.Module == "mesh-controller" && len(e.On) > 0:
controllerNode = e.On[0]
}
}
if busNode == "" {
return errors.New("no assigned module provides mesh-bus and claims mesh-broker, so there is no " +
"bus to mint credentials for — register and assign it first")
}
onNetwork, err := whereEveryoneIs(ctx, inv, shelf)
if err != nil {
return err
}
busHost := onNetwork[busNode]
if busHost == "" {
// **The hub is not in that map.** The machine that took over the tunnel is where the current
// bus already answers, and every machine dials it at the address the mesh handed them — so
// when the new bus runs on the same machine, that address is the one to tell them, with the
// new port. Found live: the control node is the hub, and the map lists the machines placed
// around it.
// The host alone: no scheme (BareAddress adds one where none was, which is the wrong
// direction here — every URL built below adds its own) and no port.
_, _, host := broker.CredentialIn(known.Address)
if host == "" {
host = known.Address
}
if _, after, hasScheme := strings.Cut(host, "://"); hasScheme {
host = after
}
host = strings.TrimSpace(host)
if i := strings.LastIndex(host, ":"); i > 0 && !strings.Contains(host[i:], "]") {
host = host[:i]
}
if host == "" {
return fmt.Errorf("%s runs the new bus and has no address on the private network, and the "+
"current bus's address is unknown too, so no machine could be told where it is", busNode)
}
busHost = host
}
busAddress := busHost + ":4222"
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
}
_, missing := broker.WithPasswords(users, hashes)
wanted := map[string]bool{}
for _, m := range missing {
wanted[m] = true
}
var machines, modules, skipped int
for _, p := range users {
if !again && !wanted[p.Username()] {
continue
}
switch p.Kind {
case broker.KindController:
if controllerNode == "" {
return errors.New("the control plane is not assigned anywhere, so its credential has nowhere to go")
}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusController})
if err != nil {
return err
}
url := "nats://" + p.Username() + ":" + password + "@" + busAddress
if err := inv.AcceptSecretForModule(ctx, controllerNode, "mesh-controller", "bus", url); err != nil {
return fmt.Errorf("the control plane's credential is minted and could not be sealed to %s: %w", controllerNode, err)
}
fmt.Printf("control plane: credential minted, sealed to %s as its `bus` secret\n", controllerNode)
case broker.KindNode:
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: p.Node})
if err != nil {
return err
}
membership, _ := json.Marshal(map[string]string{
"broker": busAddress, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
})
key, err := inv.SealingKeyOf(ctx, p.Node)
if err != nil {
return fmt.Errorf("%s has no sealing key, so its membership cannot be sealed to it: %w", p.Node, err)
}
sealed, err := secrets.Seal(key, membership)
if err != nil {
return err
}
if err := inv.PutBusMembership(ctx, p.Node, sealed); err != nil {
return err
}
machines++
case broker.KindModule:
if p.Module == "mesh-controller" {
// The control plane is a module too, and its `broker` secret is the old bus's
// credential it is still using while this runs. Writing the new bus's blob there
// cut the mesh off from its own old bus mid-move (2026-09-28). Its new-bus credential
// is the controller principal's `bus` secret above; nothing else is needed here.
skipped++
continue
}
m, inShelf := shelf[p.Module]
if !inShelf {
skipped++
continue
}
if _, reads := m.OwnSecrets["broker"]; !reads {
fmt.Printf(" %s on %s speaks on the bus but declares no `broker` secret to receive a credential in; skipped\n", p.Module, p.Node)
skipped++
continue
}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusModule, Node: p.Node, Module: p.Module})
if err != nil {
return err
}
if err := issueWith(ctx, inv, m, p.Node, "", known, busAddress, p.Username(), password); err != nil {
return err
}
modules++
default:
skipped++
}
}
fmt.Printf("minted for %d machine(s) and %d module runtime(s); %d skipped; the bus is at %s\n",
machines, modules, skipped, busAddress)
fmt.Println(" each machine's membership and each module's credential arrive with the next push of its machine;")
fmt.Println(" push the machine running the bus first, so the bus stands with its user list before anything dials it")
return nil
}
// providesBus is whether a manifest provides the mesh's bus.
func providesBus(m catalogue.Manifest) bool {
for _, o := range m.Provides {
if o.Name == "mesh-bus" {
return true
}
}
return false
}
// rolloutHand mints a machine its credential for the new bus afresh and prints its membership
// once, for an operator to carry by hand — the rescue for a machine that cannot be reached over
// any bus: rotated while it still held the old password, or reachable only by ssh. The plaintext
// exists on this terminal and then only where it is written; the store keeps the hash, and the
// sealed copy in the machine's declaration is replaced too, so the next push says the same.
func rolloutHand(ctx context.Context, node string) error {
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("the bus's certificate is not known to this process: %w", err)
}
busAddress, _, err := broker.OnNATS()
if err != nil {
return err
}
if busAddress == "" {
return errors.New("this control plane is not on the new bus, so there is no membership to hand out")
}
_, _, bare := broker.CredentialIn(busAddress)
if _, after, has := strings.Cut(bare, "://"); has {
bare = after
}
if _, err := inv.NodeByName(ctx, node); err != nil {
return err
}
p := broker.Principal{Kind: broker.KindNode, Node: node}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: node})
if err != nil {
return err
}
membership, _ := json.Marshal(map[string]string{
"broker": bare, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
})
key, err := inv.SealingKeyOf(ctx, node)
if err != nil {
return err
}
sealed, err := secrets.Seal(key, membership)
if err != nil {
return err
}
if err := inv.PutBusMembership(ctx, node, sealed); err != nil {
return err
}
// The one line of output is the membership itself, so it can be piped to the machine without
// being read on the way. Everything else goes to stderr.
fmt.Fprintf(os.Stderr, "%s's credential is minted afresh. Write this to %s on it and restart its host; "+
"then push the machine running the bus so the user list carries the new hash.\n",
node, catalogue.BusMembershipPath)
fmt.Println(string(membership))
return nil
}
+27 -8
View File
@@ -156,18 +156,37 @@ func handOver(ctx context.Context, seatName, to string) error {
if m == nil { if m == nil {
return fmt.Errorf("%s is assigned but not in the catalogue, which should not happen", module) return fmt.Errorf("%s is assigned but not in the catalogue, which should not happen", module)
} }
if err := catalogue.CanHold(*m, seat); err != nil { var was string
return fmt.Errorf("%s cannot hold %s: %w", module, seat.Name, err) holdings, err := inv.Holdings(ctx)
if err != nil {
return err
}
for _, h := range holdings {
if hs, ok := catalogue.SeatNamed(h.Claim); ok && hs.Name == seat.Name {
was = h.Node
}
} }
var was string // **Recording who already holds the seat is not making a new holder, and is not judged like
if holdings, err := inv.Holdings(ctx); err == nil { // one.** On a mesh that predates the record, the first handover has to begin by writing down
for _, h := range holdings { // the standing holder — otherwise the next holder cannot be assigned beside it, because two
if hs, ok := catalogue.SeatNamed(h.Claim); ok && hs.Name == seat.Name { // eligible claimants with nothing on record are refused. That standing holder may no longer
was = h.Node // satisfy what the seat delivers (the row moved under it, on purpose, as ADR 0131's first step),
} // and it holds regardless: derivation never read that column. So when nothing is on record and
// the named assignment is the one holding by derivation, only the claim itself is checked here.
// Every *change* of holder is judged in full.
claimsIt := false
for _, c := range m.Claims {
if cs, ok := catalogue.SeatNamed(c.Name); ok && cs.Name == seat.Name && c.At() == seat.Scope {
claimsIt = true
} }
} }
if was == "" && claimsIt {
fmt.Printf("nothing was on record for %s; recording %s on %s as its standing holder\n",
seat.Name, module, nodeName)
} else if err := catalogue.CanHold(*m, seat); err != nil {
return fmt.Errorf("%s cannot hold %s: %w", module, seat.Name, err)
}
if err := inv.HoldSeat(ctx, seat.Name, seat.Scope, nodeName, module); err != nil { if err := inv.HoldSeat(ctx, seat.Name, seat.Scope, nodeName, module); err != nil {
return err return err
} }
+125
View File
@@ -6,6 +6,7 @@ import (
"flag" "flag"
"fmt" "fmt"
"strings" "strings"
"time"
"github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/link"
@@ -218,3 +219,127 @@ func notNow(err error) error {
} }
return err return err
} }
// SourceMoved is the forge announcing a merge: every module recorded as built from that
// repository and branch is marked as moved to the merge commit, and built — bases first, so a
// module that stands on another's artifact is built after it and not against the old one
// (novox/hq 04-ISSUES/131). Nothing is pushed here: what a finished build does to the machines
// running the module is the upgrade's decision, taken when the catalogue announces it.
func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
inv := f.open.inventory
entries, err := inv.Catalogued(ctx)
if err != nil {
return notNow(err)
}
var moved []inventory.Entry
for _, e := range entries {
if !sourceIs(e.Source, m) {
continue
}
if e.Source.BuiltFrom == m.Commit {
continue
}
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
return notNow(err)
}
moved = append(moved, e)
}
if len(moved) == 0 {
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds is built from it\n",
m.Owner, m.Repo, m.Base, m.Commit)
return nil
}
ordered := orderByBases(moved)
names := make([]string, 0, len(ordered))
for _, e := range ordered {
names = append(names, e.Manifest.Module)
}
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
var failed []string
for _, e := range ordered {
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 20*time.Minute); err != nil {
fmt.Printf(" %s: %v\n", e.Manifest.Module, err)
failed = append(failed, e.Manifest.Module)
// A base that failed is a reason to stop: what stands on it would be built against
// the old one, and report success (novox/hq 04-ISSUES/131).
if standsOn(ordered, e.Manifest.Module) {
fmt.Printf(" stopping: %s is a base of what was still to build\n", e.Manifest.Module)
break
}
}
}
if len(failed) > 0 {
fmt.Printf("%d of %d not built: %s\n", len(failed), len(ordered), strings.Join(failed, ", "))
}
return nil
}
// sourceIs is whether a recorded source is the repository and branch a merge announced. A source on
// the git seat is recorded as its path on the forge; one elsewhere as the URL it was cloned from.
// An empty recorded ref is the repository's default branch, which is what a merge into the base
// branch of the forge's default means.
func sourceIs(s inventory.Source, m link.SourceMoved) bool {
want := strings.ToLower(m.Owner + "/" + m.Repo)
repo := strings.ToLower(strings.TrimSuffix(s.Repository, ".git"))
matches := repo == want || strings.HasSuffix(repo, "/"+want) ||
(m.CloneURL != "" && strings.EqualFold(strings.TrimSuffix(s.Repository, ".git"), strings.TrimSuffix(m.CloneURL, ".git")))
if !matches {
return false
}
return s.Ref == "" || s.Ref == m.Base
}
// orderByBases is the entries with every base before what stands on it: a module whose build names
// another's artifact under build.on comes after that module. Entries outside the set are not
// waited for — they are not being rebuilt. Stable for what has no order between it.
func orderByBases(entries []inventory.Entry) []inventory.Entry {
inSet := map[string]bool{}
for _, e := range entries {
inSet[e.Manifest.Module] = true
}
var out []inventory.Entry
placed := map[string]bool{}
var place func(e inventory.Entry, seen map[string]bool)
place = func(e inventory.Entry, seen map[string]bool) {
name := e.Manifest.Module
if placed[name] || seen[name] {
return
}
seen[name] = true
if e.Manifest.Build != nil {
for _, on := range e.Manifest.Build.On {
if on.Module == "" || on.Module == name || !inSet[on.Module] {
continue
}
for _, base := range entries {
if base.Manifest.Module == on.Module {
place(base, seen)
}
}
}
}
placed[name] = true
out = append(out, e)
}
for _, e := range entries {
place(e, map[string]bool{})
}
return out
}
// standsOn is whether anything in the set is built on the named module's artifacts.
func standsOn(entries []inventory.Entry, module string) bool {
for _, e := range entries {
if e.Manifest.Build == nil {
continue
}
for _, on := range e.Manifest.Build.On {
if on.Module == module {
return true
}
}
}
return false
}
+67 -4
View File
@@ -1,8 +1,14 @@
package broker package broker
import ( import (
"crypto/sha256"
"crypto/tls"
"crypto/x509"
"encoding/hex"
"errors" "errors"
"fmt" "fmt"
"os"
"strings"
"time" "time"
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
@@ -24,21 +30,78 @@ type JetStream struct {
// Dial connects and returns the controller's JetStream handle. // Dial connects and returns the controller's JetStream handle.
func Dial(url string, opts ...nats.Option) (*JetStream, error) { func Dial(url string, opts ...nats.Option) (*JetStream, error) {
// A name, because a connection nobody can identify in the server's own monitoring is one
// nobody can attribute a problem to.
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second)) opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
// checks nothing else, and so does this). Without this, the first connection failed with
// "certificate is not valid for any names" against a bus that was answering (2026-09-28).
if path := strings.TrimSpace(os.Getenv(CertificateVar)); path != "" {
pinned, err := pinnedTo(path)
if err != nil {
return nil, err
}
opts = append(opts, nats.Secure(pinned))
}
// **Its own inbox, and nothing wider.** Every principal is granted `_INBOX.<its user>.>` and
// no other inbox; the client's default prefix is random, and the server refused the first
// subscription to it (2026-09-28). The user is in the URL, so the prefix follows from it.
if user, _, _ := CredentialIn(url); user != "" {
opts = append(opts, nats.CustomInboxPrefix("_INBOX."+user))
}
// The address in an error is the address alone. The URL carries this controller's password,
// and an error here is written on the assumption it will be logged.
where := BareAddress(url)
conn, err := nats.Connect(url, opts...) conn, err := nats.Connect(url, opts...)
if err != nil { if err != nil {
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err) return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
} }
js, err := conn.JetStream() js, err := conn.JetStream()
if err != nil { if err != nil {
conn.Close() conn.Close()
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err) return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
} }
return &JetStream{conn: conn, js: js}, nil return &JetStream{conn: conn, js: js}, nil
} }
// pinnedTo is a TLS configuration that accepts exactly the certificate in the file and no other:
// the leaf's SHA-256, compared on every handshake, with the name and the chain deliberately not
// consulted — a self-signed certificate with no names is the ordinary case for a mesh's bus.
func pinnedTo(path string) (*tls.Config, error) {
want, err := FingerprintOf(path)
if err != nil {
return nil, err
}
return PinnedToFingerprint(want), nil
}
// DialPinned is Dial with the server's certificate pinned by a fingerprint the caller already holds
// — a module or a build machine that was handed one beside its credential, and has no file.
func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) {
if strings.TrimSpace(fingerprint) != "" {
opts = append(opts, nats.Secure(PinnedToFingerprint(fingerprint)))
}
return Dial(url, opts...)
}
// PinnedToFingerprint accepts exactly the certificate with this SHA-256 and no other.
func PinnedToFingerprint(want string) *tls.Config {
return &tls.Config{
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
MinVersion: tls.VersionTLS12,
VerifyPeerCertificate: func(rawCerts [][]byte, _ [][]*x509.Certificate) error {
if len(rawCerts) == 0 {
return errors.New("the bus presented no certificate")
}
sum := sha256.Sum256(rawCerts[0])
got := "sha256:" + hex.EncodeToString(sum[:])
if got != want {
return fmt.Errorf("the bus presented a certificate this mesh does not know (%s…), expected %s…", got[:23], want[:23])
}
return nil
},
}
}
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat, // Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
// a tool call — where a lost message is answered by the next one or by a timeout the caller // a tool call — where a lost message is answered by the next one or by a timeout the caller
// already handles (design 25 §3). // already handles (design 25 §3).
+29 -4
View File
@@ -169,7 +169,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// The controller owns the mesh's own traffic and the streams. It is the only writer of // The controller owns the mesh's own traffic and the streams. It is the only writer of
// stream definitions (design 25 §3), so it alone reaches the JetStream API. // stream definitions (design 25 §3), so it alone reaches the JetStream API.
pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"} pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"}
sub = []string{"mesh.control.>", "$JS.API.>"} // **And where its consumers deliver.** A push consumer delivers on `_DELIVER.<its name>`,
// and a client bound to it subscribes exactly that; the server refused it for every
// principal the first time one bound a consumer (2026-09-28). Each kind below is granted
// its own consumers' delivery subjects and no other's.
sub = []string{"mesh.control.>", "$JS.API.>", "_DELIVER." + ControllerName, "_DELIVER." + ControllerName + ".>"}
// Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A // Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A
// build is the one today: the controller asks, and reads the answer from the seat's event // build is the one today: the controller asks, and reads the answer from the seat's event
@@ -244,8 +248,15 @@ func PermissionsFor(p Principal) (Permissions, error) {
case KindNode: case KindNode:
// A host publishes its own node's control traffic and subscribes its own declaration — // A host publishes its own node's control traffic and subscribes its own declaration —
// and nothing of any other node's. // and nothing of any other node's.
pub = []string{"mesh.control." + p.Node + ".>"} // And binding to its consumer, which asks the server about it (CONSUMER.INFO) — the one
sub = []string{"mesh.node." + p.Node + ".declare"} // thing the host does that nothing granted. Found the first time a machine dialled a
// permissioned server: "this node cannot read its declarations" (2026-09-28). The ack and
// the inbox are granted below with every principal's.
pub = []string{
"mesh.control." + p.Node + ".>",
"$JS.API.CONSUMER.INFO.NODES." + p.Node,
}
sub = []string{"mesh.node." + p.Node + ".declare", "_DELIVER." + p.Node}
case KindModule: case KindModule:
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing // 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
@@ -278,7 +289,18 @@ func PermissionsFor(p Principal) (Permissions, error) {
} }
// 3. Seats it holds: full participation. // 3. Seats it holds: full participation.
// Its consumer's name, not ConsumerFor: that asks for these permissions to build the
// consumer, and would ask forever. A subject for a consumer that turns out not to exist
// grants nothing anybody can use.
sub = append(sub, "_DELIVER."+consumerDurable(p))
for _, s := range p.Holds { for _, s := range p.Holds {
// Taking work from the role's queue: the worker consumer it binds (asked about,
// delivered on, acknowledged), each on the seat's own stream. The first machine to
// take work over the new bus was refused the asking (2026-09-28).
worker := "SEAT_" + upperSnake(s.Name) + "_worker"
stream := seatStreamName(s.Name)
sub = append(sub, "_DELIVER."+worker)
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>")
for _, a := range s.Accepts { for _, a := range s.Accepts {
sub = append(sub, seatSubject(s, "accept", a)) sub = append(sub, seatSubject(s, "accept", a))
} }
@@ -514,7 +536,10 @@ func ComposeAccounts(principals []Principal) (string, error) {
// One account for the mesh: accounts in NATS isolate subject spaces entirely, and the mesh is // One account for the mesh: accounts in NATS isolate subject spaces entirely, and the mesh is
// one space (design 25 §4). The cost of that — that permissions are the only isolation — is // one space (design 25 §4). The cost of that — that permissions are the only isolation — is
// paid in the scoping of every inbox and every ack subject. // paid in the scoping of every inbox and every ack subject.
b.WriteString("accounts {\n MESH {\n users = [\n") // JetStream is enabled per account once accounts exist at all: with only the global block set,
// a user in MESH is told "JetStream not enabled for account" the first time it binds a
// consumer, which is the first thing every host does (2026-09-28).
b.WriteString("accounts {\n MESH {\n jetstream: enabled\n users = [\n")
for _, p := range sorted { for _, p := range sorted {
perms, err := PermissionsFor(p) perms, err := PermissionsFor(p)
if err != nil { if err != nil {
+5 -8
View File
@@ -35,10 +35,6 @@ type Readiness struct {
Modules []string Modules []string
// ModuleCredentialled is which of those has one. // ModuleCredentialled is which of those has one.
ModuleCredentialled map[string]bool ModuleCredentialled map[string]bool
// StillOnTheOldBus is whether anything of the mesh's own still needs the bus it is leaving —
// which is not a reason to stop, because that broker stays as an ordinary provider of `amqp`
// (ADR 0119). Recorded so nobody reads the move as a retirement.
OldBusHasOtherClients bool
} }
// NotReady is every reason this mesh cannot move its bus yet, in the order somebody would fix them. // NotReady is every reason this mesh cannot move its bus yet, in the order somebody would fix them.
@@ -120,10 +116,11 @@ func WhatMoves(r Readiness) []string {
out = append(out, fmt.Sprintf("move %d module runtime(s), and confirm each answers", out = append(out, fmt.Sprintf("move %d module runtime(s), and confirm each answers",
len(r.Modules))) len(r.Modules)))
} }
if r.OldBusHasOtherClients { // **The old broker goes, and it goes last** (novox/hq ADR 0131). AMQP is not a provision, so once
out = append(out, "leave the old broker running: it stays an ordinary provider of `amqp` for "+ // every machine reports on the new bus nothing of the mesh is left speaking to it, and its module
"whatever else uses it (ADR 0119), and this move is not its retirement") // is unassigned. Said as a step so nobody reads the move as leaving a second bus behind.
} out = append(out, "then unassign the old broker's module: AMQP is not a provision (ADR 0131), and "+
"once every machine reports on the new bus nothing of the mesh speaks to it")
return out return out
} }
+7 -6
View File
@@ -79,9 +79,8 @@ func TestEachThingMissingNamesItsOwnRemedy(t *testing.T) {
// What the move would do is written out rather than summarised, because this is the one step with // What the move would do is written out rather than summarised, because this is the one step with
// nothing to inspect afterwards — so reading it is the last chance to disagree. // nothing to inspect afterwards — so reading it is the last chance to disagree.
func TestWhatMovesNamesEveryMachineAndSaysTheOldBrokerStays(t *testing.T) { func TestWhatMovesNamesEveryMachineAndEndsWithTheOldBrokerGoing(t *testing.T) {
r := aMeshReadyToMove() r := aMeshReadyToMove()
r.OldBusHasOtherClients = true
steps := strings.Join(WhatMoves(r), "\n") steps := strings.Join(WhatMoves(r), "\n")
for _, want := range []string{"anchor", "laptop", "user list", "module runtime"} { for _, want := range []string{"anchor", "laptop", "user list", "module runtime"} {
@@ -89,9 +88,11 @@ func TestWhatMovesNamesEveryMachineAndSaysTheOldBrokerStays(t *testing.T) {
t.Errorf("the plan does not mention %q:\n%s", want, steps) t.Errorf("the plan does not mention %q:\n%s", want, steps)
} }
} }
// Said explicitly, so nobody reads the move as switching the old broker off — it stays serving // Said explicitly, and last: AMQP is not a provision (novox/hq ADR 0131), so the move ends with
// whatever else uses it, and that is a decision already taken. // the old broker's module unassigned, not left behind as a second bus. An earlier version of this
if !strings.Contains(steps, "not its retirement") { // test pinned the opposite, under a record 0131 superseded.
t.Errorf("the plan does not say the old broker stays:\n%s", steps) lines := WhatMoves(r)
if last := lines[len(lines)-1]; !strings.Contains(last, "unassign the old broker") {
t.Errorf("the plan does not end with the old broker going:\n%s", steps)
} }
} }
+3
View File
@@ -175,6 +175,9 @@ var ControllerFollows = []string{
// message on the control branch. Same three audiences, one publish: whoever asked, this, and the // message on the control branch. Same three audiences, one publish: whoever asked, this, and the
// catalogue. // catalogue.
seatEventSubject("mesh-build-machine", "built"), seatEventSubject("mesh-build-machine", "built"),
// The forge's merges: what moved a source, so the mesh builds what that source produces
// without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name.
moduleEventSubject("gitea", "pull.merged"),
} }
// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so // moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so
+8 -7
View File
@@ -21,10 +21,11 @@ jetstream {
accounts { accounts {
MESH { MESH {
jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] } publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] } subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
@@ -32,21 +33,21 @@ accounts {
subscribe: { allow: ["_INBOX.enrol.one.>"] } subscribe: { allow: ["_INBOX.enrol.one.>"] }
} } } }
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: { { user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
publish: { allow: ["$JS.ACK.NODES.one.>", "mesh.control.one.>"] } publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] } subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} } } }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] } subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.one_telegram", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>"] } publish: { allow: ["$JS.ACK.EVENTS.two_audit.>"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] } subscribe: { allow: ["_DELIVER.two_audit", "_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
} } } }
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: { { user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] } publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>"] } subscribe: { allow: ["_DELIVER.two_shop", "_INBOX.two.shop.>"] }
} } } }
] ]
} }
+7 -1
View File
@@ -186,7 +186,13 @@ func TestWhatTheMeshWritesIsUsersAndNothingAboutTheServer(t *testing.T) {
} }
// None of the server's own settings. Each of these in the mesh's file is a value the controller // None of the server's own settings. Each of these in the mesh's file is a value the controller
// would then own, and the module could no longer change its own image without the mesh agreeing. // would then own, and the module could no longer change its own image without the mesh agreeing.
for _, absent := range []string{"port:", "http:", "jetstream", "tls {", "store_dir", "cert_file"} { // `jetstream {` is the server's block (its store, its limits); `jetstream: enabled` inside the
// account is the account's, and the mesh owns the account — a user in it is told "JetStream
// not enabled for account" without it (2026-09-28).
if !strings.Contains(got, "jetstream: enabled") {
t.Errorf("the account does not enable JetStream, so no user in it can bind a consumer")
}
for _, absent := range []string{"port:", "http:", "jetstream {", "tls {", "store_dir", "cert_file"} {
if strings.Contains(got, absent) { if strings.Contains(got, absent) {
t.Errorf("the accounts file contains %q, which belongs to the module that raises the "+ t.Errorf("the accounts file contains %q, which belongs to the module that raises the "+
"server, not to the mesh", absent) "server, not to the mesh", absent)
@@ -0,0 +1,48 @@
package catalogue
import (
"os"
"path/filepath"
"strings"
"testing"
)
// **The word does not come back through a manifest** (novox/hq ADR 0131). A module that wants
// messaging wants the mesh's bus, reached through the sdk and named by the `mesh-broker` seat. Naming
// the old wire protocol asks for the one server being retired, so both directions are refused at the
// parser — this is judged from the manifest alone, no store needed.
func TestAManifestProvidingAmqpIsRefused(t *testing.T) {
raw := []byte(`{"module":"old-broker","version":"1","provides":[{"name":"amqp","scope":"mesh"}]}`)
_, err := ParseManifest(raw)
if err == nil || !strings.Contains(err.Error(), `provides "amqp", which is not a provision`) {
t.Fatalf("a module providing amqp was not refused, or not for the reason: %v", err)
}
}
func TestAManifestRequiringAmqpIsRefused(t *testing.T) {
raw := []byte(`{"module":"forwarder","version":"1","requires":["amqp"]}`)
_, err := ParseManifest(raw)
if err == nil || !strings.Contains(err.Error(), `requires "amqp", which is not a provision`) {
t.Fatalf("a module requiring amqp was not refused, or not for the reason: %v", err)
}
}
// And the catalogue as checked out beside this repository names it nowhere — the three modules that
// did are removed under design 28 task 5.4, not converted.
func TestNoCatalogueManifestNamesAmqp(t *testing.T) {
modules, err := filepath.Glob("../../../mesh-catalog/modules/*/module.json")
if err != nil || len(modules) == 0 {
t.Skip("the catalogue is not checked out beside this repository")
}
for _, path := range modules {
raw, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
if strings.Contains(string(raw), `"amqp"`) {
t.Errorf("%s names amqp, which is not a provision (novox/hq ADR 0131)",
filepath.Base(filepath.Dir(path)))
}
}
}
+3 -35
View File
@@ -46,41 +46,9 @@ func TestTheSeatRefusesADifferentBusToo(t *testing.T) {
} }
} }
// **The old broker claims the seat until the seat is handed over, and stands beside the new one // The old broker is gone from the catalogue (novox/hq ADR 0131, design 28 task 5.4), so it is no
// while it waits** (novox/hq ADR 0131, superseding the record this test used to pin). Whoever is on // longer a fixture here. That two eligible holders stand beside each other with one on record is
// record holds it; the other eligible claimant is neither refused nor holding. This is the shape the // pinned in holdings_test.go against manifests this package owns.
// handover needs: both brokers assigned, one bus, no moment with nobody in the seat.
func TestTheOldBrokerStandsBesideTheNewOneUntilTheHandover(t *testing.T) {
was := Seats()
t.Cleanup(func() { UseSeats(was) })
UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "test"}})
lavinmq := catalogueManifest(t, "lavinmq")
if !lavinmq.ClaimsSeat("mesh-broker") {
t.Skip("the old broker no longer claims the seat: design 28 task 5.4 has removed it")
}
nats := catalogueManifest(t, "nats")
onRecord := World{Holdings: []Held{{Claim: "mesh-broker", Scope: ScopeMesh,
Node: "anchor", Module: "nats"}}}
// The same machine runs both. Without the record this is two holders and refused; with it, the
// recorded one holds and the other is silent.
anchor := workstation()
anchor.Name = "anchor"
got, err := Resolve(shelf(lavinmq, nats), []string{"lavinmq", "nats"}, anchor, onRecord)
if err != nil {
t.Fatalf("the old broker beside the recorded holder was refused: %v", err)
}
var holders []string
for _, h := range got.Claims {
if h.Claim == "mesh-broker" {
holders = append(holders, h.Module)
}
}
if len(holders) != 1 || holders[0] != "nats" {
t.Fatalf("the seat is held by %v, not by the holder on record alone", holders)
}
}
// **A seat and the interface it delivers are different names, and renaming one must not rename // **A seat and the interface it delivers are different names, and renaming one must not rename
// the other** (novox/hq ADR 0118). This nearly went wrong: the seats were renamed to the `mesh-*` // the other** (novox/hq ADR 0118). This nearly went wrong: the seats were renamed to the `mesh-*`
+19
View File
@@ -103,6 +103,11 @@ type Rendering struct {
// **Only the users, never the server's own settings**: those are the module's, in its image and // **Only the users, never the server's own settings**: those are the module's, in its image and
// its mounts (Manifest.BusUsers). // its mounts (Manifest.BusUsers).
BusUsers string BusUsers string
// BusMembership is this machine's membership for the bus the mesh is moving to, sealed to it
// (design 28, task 5.2). Empty for a machine not being moved. Written as a file the host reads
// after the declaration has applied, so the bus it names is standing before the machine leaves
// the one it is on.
BusMembership string
// MeshRange is the private network's CIDR (the range node addresses are allocated from), for a // MeshRange is the private network's CIDR (the range node addresses are allocated from), for a
// module that must name the whole mesh rather than one machine — an intrusion filter that must // module that must name the whole mesh rather than one machine — an intrusion filter that must
@@ -238,9 +243,23 @@ func (r Resolution) Compose(with Rendering) (Composed, error) {
if err != nil { if err != nil {
return Composed{}, err return Composed{}, err
} }
if with.BusMembership != "" {
// The machine's own, not any module's: how it reaches the mesh from now on. Sealed like a
// secret and placed where the host looks for exactly this (design 28, task 5.2).
resources = append(resources, map[string]any{
"id": BusMembershipID(), "type": "file", "path": BusMembershipPath,
"sealed": with.BusMembership, "mode": "0600",
})
}
return Composed{Resources: resources, Owner: owner}, nil return Composed{Resources: resources, Owner: owner}, nil
} }
// BusMembershipID names the resource carrying a machine's membership for the new bus, and
// BusMembershipPath is where the host reads it — the same constant on both sides.
func BusMembershipID() string { return "bus-membership" }
const BusMembershipPath = "/var/lib/mesh/membership-next.json"
func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) { func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) {
// Every manifest is placed first (novox/hq ADR 0112): the maps naming where its bindings, // Every manifest is placed first (novox/hq ADR 0112): the maps naming where its bindings,
// credentials and contributions land are resolved against this node's directories, so every // credentials and contributions land are resolved against this node's directories, so every
@@ -28,8 +28,10 @@ func TestTheStoreAndTheBrokerSayWhatTheMeshGuards(t *testing.T) {
if got := catalogueManifest(t, "postgres").Guards; !reflect.DeepEqual(got, []int{5432}) { if got := catalogueManifest(t, "postgres").Guards; !reflect.DeepEqual(got, []int{5432}) {
t.Errorf("postgres guards %v; the store's port must be refused from outside", got) t.Errorf("postgres guards %v; the store's port must be refused from outside", got)
} }
if got := catalogueManifest(t, "lavinmq").Guards; !reflect.DeepEqual(got, []int{15672}) { // The bus's monitoring port, not its client port: a node reaches the bus, nobody outside
t.Errorf("lavinmq guards %v; the management port must be refused from outside", got) // reads its state (novox/hq ADR 0131 — the broker that guarded 15672 has left the catalogue).
if got := catalogueManifest(t, "nats").Guards; !reflect.DeepEqual(got, []int{8222}) {
t.Errorf("nats guards %v; the monitoring port must be refused from outside", got)
} }
} }
+44
View File
@@ -114,3 +114,47 @@ func TestCanHoldJudgesClaimScopeAndWhatTheSeatDelivers(t *testing.T) {
t.Fatalf("with the row saying amqp, an amqp provider was refused: %v", err) t.Fatalf("with the row saying amqp, an amqp provider was refused: %v", err)
} }
} }
// A machine being moved is handed its membership for the new bus as a sealed file in its own
// declaration — the machine's, not any module's (design 28, task 5.2).
func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) {
r := Resolution{Node: "anchor"}
got, err := r.Compose(Rendering{BusMembership: "sealed-blob"})
if err != nil {
t.Fatal(err)
}
var found map[string]any
for _, res := range got.Resources {
if res["id"] == BusMembershipID() {
found = res
}
}
if found == nil {
t.Fatalf("no membership resource in %v", got.Resources)
}
if found["path"] != BusMembershipPath || found["sealed"] != "sealed-blob" || found["mode"] != "0600" {
t.Fatalf("the membership is not a sealed 0600 file where the host reads it: %v", found)
}
// And a machine not being moved is handed nothing.
got, _ = r.Compose(Rendering{})
for _, res := range got.Resources {
if res["id"] == BusMembershipID() {
t.Fatal("a machine with no membership on record was handed one")
}
}
}
// The store's seat rows have no protocol columns yet; loading them must not drop the protocol the
// bus is derived from, or no role's work queue is ever raised (found live, 2026-09-28).
func TestAStoreRowWithoutAProtocolKeepsTheCompiledOne(t *testing.T) {
was := Seats()
t.Cleanup(func() { UseSeats(was) })
UseSeats([]Seat{{Name: "mesh-build-machine", Scope: ScopeMesh, Decision: "row"}})
got, ok := SeatNamed("mesh-build-machine")
if !ok || len(got.Accepts) == 0 {
t.Fatalf("the build machine's seat lost what it accepts when loaded from the store: %+v", got)
}
if got.Decision != "row" {
t.Fatalf("the store's own columns were not kept: %+v", got)
}
}
+22
View File
@@ -1027,6 +1027,28 @@ func ParseManifest(raw []byte) (Manifest, error) {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%q is not a usable slug: lower-case letters, digits, dashes and dots", m.Slug)) "%q is not a usable slug: lower-case letters, digits, dashes and dots", m.Slug))
} }
// **`amqp` is not a provision, and not a requirement** (novox/hq ADR 0131). A module that wants
// messaging wants the mesh's bus — it emits and consumes through the sdk, which the mesh hands the
// bus with the module's own credential — and the bus is whatever holds `mesh-broker`, spoken in
// whatever that holder speaks. Naming the old wire protocol asks for a specific server, and the
// only one that could answer is the one being retired. Refused here so the word cannot come back
// through a manifest.
for _, offer := range m.Provides {
if offer.Name == "amqp" {
problems = append(problems, fmt.Sprintf(
"%s provides %q, which is not a provision: the mesh's bus is whatever holds "+
"mesh-broker, and a module provides mesh-bus to be it (novox/hq ADR 0131)",
m.Module, offer.Name))
}
}
for _, r := range m.Requires {
if r == "amqp" {
problems = append(problems, fmt.Sprintf(
"%s requires %q, which is not a provision: a module reaches the mesh's bus through "+
"the sdk, and depends on the mesh-broker seat, not on a protocol (novox/hq ADR 0131)",
m.Module, r))
}
}
for _, offer := range m.Provides { for _, offer := range m.Provides {
p := offer.Name p := offer.Name
if !name.MatchString(p) { if !name.MatchString(p) {
+1 -5
View File
@@ -13,7 +13,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
"type": "container", "id": "server", "name": "mesh-controller", "type": "container", "id": "server", "name": "mesh-controller",
"env": map[string]any{ "env": map[string]any{
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}", "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
}, },
@@ -27,7 +26,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
env := control["env"].(map[string]any) env := control["env"].(map[string]any)
for key, want := range map[string]string{ for key, want := range map[string]string{
"MESH_STORE_INVENTORY_PORT": "6852", "MESH_STORE_INVENTORY_PORT": "6852",
"MESH_BROKER_AMQP_PORT": "5679",
"MESH_BROKER_ADDRESS_PORT": "5671", "MESH_BROKER_ADDRESS_PORT": "5671",
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
} { } {
@@ -146,7 +144,6 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
"MESH_STORE_INVENTORY_PORT": "6852", "MESH_STORE_INVENTORY_PORT": "6852",
"MESH_STORE_IDENTITY_PORT": "6852", "MESH_STORE_IDENTITY_PORT": "6852",
"MESH_STORE_LICENCES_PORT": "6852", "MESH_STORE_LICENCES_PORT": "6852",
"MESH_BROKER_AMQP_PORT": "5679",
"MESH_BROKER_MANAGEMENT_PORT": "15673", "MESH_BROKER_MANAGEMENT_PORT": "15673",
"MESH_BROKER_ADDRESS_PORT": "5671", "MESH_BROKER_ADDRESS_PORT": "5671",
} { } {
@@ -164,7 +161,7 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any) env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any)
if env["MESH_STORE_INVENTORY_PORT"] != "" || env["MESH_BROKER_AMQP_PORT"] != "" { if env["MESH_STORE_INVENTORY_PORT"] != "" {
t.Errorf("with no settings, the control plane is told %v", env) t.Errorf("with no settings, the control plane is told %v", env)
} }
} }
@@ -175,7 +172,6 @@ var SeatPorts = map[string]string{
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}", "MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}", "MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}", "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
} }
+23 -2
View File
@@ -109,9 +109,30 @@ func DefaultSeats() []Seat { return append([]Seat(nil), defaultSeats...) }
// than running on the set the binary shipped with. So the store can only ever *replace* the set with // than running on the set the binary shipped with. So the store can only ever *replace* the set with
// a non-empty one, never erase it. // a non-empty one, never erase it.
func UseSeats(s []Seat) { func UseSeats(s []Seat) {
if len(s) > 0 { if len(s) == 0 {
seats = s return
} }
// **The store's rows carry no protocol yet, and the protocol is what the bus is derived
// from.** ADR 0129 gives a seat what it accepts, emits and serves; ADR 0122 moved the set into
// a table that has name, scope, delivers and decision and nothing else, and the columns for
// the rest are not there yet. So a row replacing a compiled entry would silently drop the
// protocol, and the roles' work queues would never be raised — found live as "no response
// from stream" the first time a build was submitted over the new bus (2026-09-28). Until the
// table gains the columns, a row without a protocol keeps the compiled one of the same name.
byName := map[string]Seat{}
for _, d := range defaultSeats {
byName[d.Name] = d
}
merged := make([]Seat, 0, len(s))
for _, row := range s {
if len(row.Accepts)+len(row.Emits)+len(row.Serves) == 0 {
if d, known := byName[row.Name]; known {
row.Accepts, row.Emits, row.Serves = d.Accepts, d.Emits, d.Serves
}
}
merged = append(merged, row)
}
seats = merged
} }
// aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from // aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from
+5 -5
View File
@@ -121,9 +121,9 @@ func TestTheParserDoesNotJudgeWhatOnlyTheStoreKnows(t *testing.T) {
t.Cleanup(func() { UseSeats(was) }) t.Cleanup(func() { UseSeats(was) })
// A store whose bus seat delivers something this module does provide. // A store whose bus seat delivers something this module does provide.
UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "amqp", Decision: "test"}}) UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "test"}})
raw := []byte(`{"module":"lavinmq","version":"1",` + raw := []byte(`{"module":"a-bus","version":"1",` +
`"provides":[{"name":"amqp","scope":"mesh"}],` + `"provides":[{"name":"mesh-bus","scope":"mesh"}],` +
`"claims":[{"name":"mesh-broker","scope":"mesh"}]}`) `"claims":[{"name":"mesh-broker","scope":"mesh"}]}`)
m, err := ParseManifest(raw) m, err := ParseManifest(raw)
@@ -135,12 +135,12 @@ func TestTheParserDoesNotJudgeWhatOnlyTheStoreKnows(t *testing.T) {
} }
// And with the store saying the seat delivers something else, registration is what refuses it. // And with the store saying the seat delivers something else, registration is what refuses it.
UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "test"}}) UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "other-bus", Decision: "test"}})
if _, err := ParseManifest(raw); err != nil { if _, err := ParseManifest(raw); err != nil {
t.Fatalf("the parser judged it the second time: %v", err) t.Fatalf("the parser judged it the second time: %v", err)
} }
got := strings.Join(CatalogueProblems(Shelf{m.Module: m}), "; ") got := strings.Join(CatalogueProblems(Shelf{m.Module: m}), "; ")
if !strings.Contains(got, `does not provide "mesh-bus"`) { if !strings.Contains(got, `does not provide "other-bus"`) {
t.Fatalf("registration did not refuse a holder that cannot answer for the seat: %q", got) t.Fatalf("registration did not refuse a holder that cannot answer for the seat: %q", got)
} }
} }
+32
View File
@@ -230,3 +230,35 @@ func (i *Inventory) ForgetPerson(ctx context.Context, name string) error {
} }
return i.ForgetBusUser(ctx, "person."+name) return i.ForgetBusUser(ctx, "person."+name)
} }
// PutBusMembership records a machine's membership for the new bus, sealed to it (design 28, 5.2).
// Replaces any earlier one: a machine has one membership per bus, and re-minting is re-telling.
func (i *Inventory) PutBusMembership(ctx context.Context, nodeName, sealed string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into bus_membership (node, sealed) values ($1, $2)
on conflict (node) do update set sealed = excluded.sealed, since = now()`, node.ID, sealed)
return err
}
// BusMemberships is every machine's sealed membership for the new bus, by node name.
func (i *Inventory) BusMemberships(ctx context.Context) (map[string]string, error) {
rows, err := i.store.Pool().Query(ctx,
`select n.name, b.sealed from bus_membership b join node n on n.id = b.node`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]string{}
for rows.Next() {
var name, sealed string
if err := rows.Scan(&name, &sealed); err != nil {
return nil, err
}
out[name] = sealed
}
return out, rows.Err()
}
+17
View File
@@ -80,3 +80,20 @@ func TestUnassigningTheHolderTakesTheHoldingWithIt(t *testing.T) {
t.Fatalf("the holding outlived the assignment it pointed at: %+v", held) t.Fatalf("the holding outlived the assignment it pointed at: %+v", held)
} }
} }
func TestAMachinesMembershipIsOneRowReplacedAndGoesWithTheMachine(t *testing.T) {
inv, ctx := twoBrokersOnTwoNodes(t)
if err := inv.PutBusMembership(ctx, "anchor", "first"); err != nil {
t.Fatal(err)
}
if err := inv.PutBusMembership(ctx, "anchor", "second"); err != nil {
t.Fatal(err)
}
got, err := inv.BusMemberships(ctx)
if err != nil {
t.Fatal(err)
}
if got["anchor"] != "second" || len(got) != 1 {
t.Fatalf("a re-told membership did not replace the first: %v", got)
}
}
@@ -0,0 +1,12 @@
-- The bus seat's holder answers for the mesh's bus, not for a wire protocol (novox/hq ADR 0131).
--
-- The row said `amqp`, which is the protocol the old broker spoke, and so only that broker could hold
-- the seat that names the mesh's bus — while the module that will carry the bus could not. The seat
-- delivers `mesh-bus`; whichever module provides that may hold it, and today that is one module.
--
-- Safe under the current holder: the control plane composes its own bus address through the seat by
-- name, and the overview derives holders by name. Only registration and provision-to-seat resolution
-- read this column. So the row changes, the current holder keeps holding by derivation, the next one
-- can register its claim, and the handover (0039) moves the seat when both are running. What must
-- not happen in between is re-registering the current holder — registration would now refuse it.
update seat set delivers = 'mesh-bus' where name = 'mesh-broker' and delivers = 'amqp';
@@ -0,0 +1,11 @@
-- A machine already enrolled is moved to the new bus by being told its membership for it
-- (novox/hq design 28, task 5.2). Until this, a membership — bus address, fingerprint, password,
-- transport — existed only in the enrolment reply, and nothing could hand one to a machine that
-- had already joined. The row is the membership sealed to that machine, composed into its
-- declaration as a file it reads after applying; the plaintext exists once, at minting, and then
-- only on the machine. One per node: the mesh moves to one bus.
create table bus_membership (
node uuid primary key references node(id) on delete cascade,
sealed text not null,
since timestamptz not null default now()
);
+73
View File
@@ -0,0 +1,73 @@
package link
import (
"context"
"fmt"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
)
// OverNats is the controller's outbound on the bus being built: the same three acts the other
// transport has, on the subjects the permissions were derived for (design 25). A declaration is a
// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on
// the subject its name derives to; a tool is asked by request and reply on the module's tool subject.
type OverNats struct{ JS *broker.JetStream }
// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the
// only subject that node's consumer delivers. The host subscribes exactly this.
func declareSubject(node string) string { return "mesh.node." + node + ".declare" }
func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error {
publish, cancel := context.WithTimeout(ctx, 15*time.Second)
defer cancel()
if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil {
return fmt.Errorf("declaring to %s: %w", node, err)
}
return nil
}
func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
// The key is the subject: the controller's own events are named in full, and what a module
// emits is derived before it reaches here. Headers carry the envelope the other transport put
// in message properties (ADR 0042), so a consumer reads who and when without the payload.
msg := nats.NewMsg(key)
msg.Data = body
msg.Header.Set("x-source", source)
msg.Header.Set("x-node", node)
msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano))
if err := b.JS.Conn().PublishMsg(msg); err != nil {
return fmt.Errorf("announcing %s: %w", key, err)
}
return nil
}
func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
ask, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args)
if err != nil {
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
}
return reply.Data, nil
}
// ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one
// JetStream connection the caller has already raised the streams on. Nothing is declared here —
// the streams and the controller's consumers are asserted by Raise, before anything is served.
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
return &Server{
inbound: Nats(js),
bus: OverNats{JS: js},
js: js,
enroller: enroller,
listener: listener,
log: newLog(),
}
}
// Bus is the controller's outbound, whichever transport it connected over. Callers that send a
// declaration or ask a tool use this rather than the channel, which one transport does not have.
func (s *Server) Bus() Bus { return s.bus }
+13
View File
@@ -117,6 +117,19 @@ type Announcement struct {
} }
// Upgraded is what the catalogue says when a module's current version moves. // Upgraded is what the catalogue says when a module's current version moves.
// SourceMoved is what the forge announces when a pull request is merged: which repository, into
// which branch, producing which commit. The mesh matches it against every module's recorded
// source and builds what moved, bases first.
type SourceMoved struct {
Owner string `json:"owner"`
Repo string `json:"repo"`
Base string `json:"base"`
Head string `json:"head"`
Commit string `json:"merge_commit_sha"`
CloneURL string `json:"clone_url"`
HTMLURL string `json:"html_url"`
}
type Upgraded struct { type Upgraded struct {
Module string `json:"module"` Module string `json:"module"`
Commit string `json:"commit"` Commit string `json:"commit"`
+3
View File
@@ -30,6 +30,9 @@ const (
KindHeartbeat = "heartbeat" KindHeartbeat = "heartbeat"
KindBuilt = "built" KindBuilt = "built"
KindModuleMoved = "module-moved" KindModuleMoved = "module-moved"
// KindSourceMoved is the forge announcing a merge: a source moved, and what it produces is
// built without anybody telling the mesh (novox/hq 04-ISSUES/131).
KindSourceMoved = "source-moved"
KindCatchUp = "catch-up" KindCatchUp = "catch-up"
) )
+4
View File
@@ -67,6 +67,10 @@ func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound {
// are events get their own queue each, and only when something is listening. // are events get their own queue each, and only when something is listening.
func (c *currentInbound) Also(kind string) error { func (c *currentInbound) Also(kind string) error {
switch kind { switch kind {
case KindSourceMoved:
// Not followed on the bus the mesh is leaving: the forge's merges are announced on the
// new one, and this transport goes with the move (design 28, task 5.5).
return nil
case KindModuleMoved: case KindModuleMoved:
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil { if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
return err return err
+3 -1
View File
@@ -48,7 +48,7 @@ func Nats(js *broker.JetStream) Inbound {
// whatever was asked for — and not at all when nothing was. // whatever was asked for — and not at all when nothing was.
func (n *natsInbound) Also(kind string) error { func (n *natsInbound) Also(kind string) error {
switch kind { switch kind {
case KindModuleMoved, KindCatchUp: case KindModuleMoved, KindCatchUp, KindSourceMoved:
n.follows[kind] = true n.follows[kind] = true
return nil return nil
default: default:
@@ -180,6 +180,8 @@ func kindOfSubject(subject string) (string, bool) {
return KindModuleMoved, true return KindModuleMoved, true
case broker.ControllerFollows[1]: case broker.ControllerFollows[1]:
return KindCatchUp, true return KindCatchUp, true
case broker.ControllerFollows[3]:
return KindSourceMoved, true
case BuildOutcome(): case BuildOutcome():
// A build's outcome is the role's event now, so it arrives on the events stream rather than // A build's outcome is the role's event now, so it arrives on the events stream rather than
// the control branch — and is acted on by the same handler, because what the controller does // the control branch — and is acted on by the same handler, because what the controller does
+17
View File
@@ -374,3 +374,20 @@ func TestNatsAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
t.Fatalf("a report was recorded by a store that never came back") t.Fatalf("a report was recorded by a store that never came back")
} }
} }
// A merge announcement is not what these tests are about; taken and forgotten.
func (t *toldAbout) SourceMoved(context.Context, SourceMoved) error { return nil }
// Every subject the controller follows decodes to a kind, and decoding never reaches past the
// list: the day the list was three long and the decoder named a fourth, every message panicked
// the control plane (2026-09-28).
func TestEverySubjectTheControllerFollowsDecodesToAKind(t *testing.T) {
for _, subject := range broker.ControllerFollows {
if _, ok := kindOfSubject(subject); !ok {
t.Errorf("%s is followed and decodes to nothing", subject)
}
}
if _, ok := kindOfSubject("mesh.mod.nobody.event.nothing"); ok {
t.Error("a subject nobody follows decoded to a kind")
}
}
+36 -1
View File
@@ -7,6 +7,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"github.com/novox/mesh-controller/internal/broker"
"log" "log"
"os" "os"
"time" "time"
@@ -57,6 +58,8 @@ type Upgrader interface {
// stop every upgrade behind it — except the store unreachable for the moment, which is asked // stop every upgrade behind it — except the store unreachable for the moment, which is asked
// again for a bounded time (novox/hq issue 083). // again for a bounded time (novox/hq issue 083).
Upgraded(ctx context.Context, u Upgraded) error Upgraded(ctx context.Context, u Upgraded) error
// SourceMoved is a merge on the forge: build what that source produces, bases first.
SourceMoved(ctx context.Context, m SourceMoved) error
} }
// Server acts on what nodes and modules say. // Server acts on what nodes and modules say.
@@ -70,6 +73,7 @@ type Server struct {
bus Bus bus Bus
conn *amqp.Connection conn *amqp.Connection
channel *amqp.Channel channel *amqp.Channel
js *broker.JetStream
enroller Enroller enroller Enroller
listener Listener listener Listener
@@ -99,6 +103,9 @@ func (s *Server) Records(r Recorder) { s.recorder = r }
// Follows says what to do about upgrades, and asks for them to be delivered. // Follows says what to do about upgrades, and asks for them to be delivered.
func (s *Server) Follows(u Upgrader) error { func (s *Server) Follows(u Upgrader) error {
if err := s.inbound.Also(KindSourceMoved); err != nil {
return err
}
if err := s.inbound.Also(KindModuleMoved); err != nil { if err := s.inbound.Also(KindModuleMoved); err != nil {
return err return err
} }
@@ -180,10 +187,12 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
channel: channel, channel: channel,
enroller: enroller, enroller: enroller,
listener: listener, listener: listener,
log: log.New(os.Stdout, "", log.LstdFlags), log: newLog(),
}, nil }, nil
} }
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
// Channel is the controller's channel, for the command line's own publishing. // Channel is the controller's channel, for the command line's own publishing.
func (s *Server) Channel() *amqp.Channel { return s.channel } func (s *Server) Channel() *amqp.Channel { return s.channel }
@@ -197,6 +206,9 @@ func (s *Server) Close() {
if s.conn != nil { if s.conn != nil {
_ = s.conn.Close() _ = s.conn.Close()
} }
if s.js != nil {
s.js.Close()
}
} }
// Serve acts on what arrives until the context ends. // Serve acts on what arrives until the context ends.
@@ -225,6 +237,8 @@ func (s *Server) act(ctx context.Context, m Control) {
s.wasBuilt(ctx, m) s.wasBuilt(ctx, m)
case KindModuleMoved: case KindModuleMoved:
s.moved(ctx, m) s.moved(ctx, m)
case KindSourceMoved:
s.sourceMoved(ctx, m)
case KindCatchUp: case KindCatchUp:
s.catchingUp(ctx, m) s.catchingUp(ctx, m)
default: default:
@@ -596,3 +610,24 @@ func short(commit string) string {
} }
return commit return commit
} }
// sourceMoved acts on the forge's announcement of a merge. Taken whatever happens: a build that
// fails is reported by the build itself, and re-delivering the merge would only re-fail it.
func (s *Server) sourceMoved(ctx context.Context, m Control) {
var moved SourceMoved
if err := json.Unmarshal(m.Body(), &moved); err != nil {
s.log.Printf("a merge announcement could not be read: %v", err)
_ = m.Took()
return
}
m.About("merge " + moved.Owner + "/" + moved.Repo + " into " + moved.Base)
if moved.Commit == "" || moved.Repo == "" {
s.log.Printf("a merge announcement named no repository or no commit; ignored")
_ = m.Took()
return
}
if err := s.upgrader.SourceMoved(ctx, moved); err != nil {
s.log.Printf("%s/%s moved to %.8s and the mesh could not act on it: %v", moved.Owner, moved.Repo, moved.Commit, err)
}
_ = m.Took()
}
+2
View File
@@ -152,3 +152,5 @@ func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to) t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to)
} }
} }
func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return nil }
+5 -4
View File
@@ -23,7 +23,8 @@
"licences": "/var/lib/mesh/mesh-controller/licences", "licences": "/var/lib/mesh/mesh-controller/licences",
"broker": "/var/lib/mesh/mesh-controller/broker", "broker": "/var/lib/mesh/mesh-controller/broker",
"broker-management": "/var/lib/mesh/mesh-controller/broker-management", "broker-management": "/var/lib/mesh/mesh-controller/broker-management",
"broker-address": "/var/lib/mesh/mesh-controller/broker-address" "broker-address": "/var/lib/mesh/mesh-controller/broker-address",
"bus": "/var/lib/mesh/mesh-controller/bus"
}, },
"secrets-owner": "65534:65534", "secrets-owner": "65534:65534",
"resources": [ "resources": [
@@ -46,15 +47,14 @@
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
"MESH_STORE_IDENTITY_FILE": "/run/secrets/identity", "MESH_STORE_IDENTITY_FILE": "/run/secrets/identity",
"MESH_STORE_LICENCES_FILE": "/run/secrets/licences", "MESH_STORE_LICENCES_FILE": "/run/secrets/licences",
"MESH_BROKER_AMQP_FILE": "/run/secrets/broker",
"MESH_BROKER_MANAGEMENT_FILE": "/run/secrets/broker-management", "MESH_BROKER_MANAGEMENT_FILE": "/run/secrets/broker-management",
"MESH_BROKER_ADDRESS_FILE": "/run/secrets/broker-address", "MESH_BROKER_ADDRESS_FILE": "/run/secrets/broker-address",
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}", "MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}", "MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}" "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
"MESH_BUS_NATS_FILE": "/run/secrets/bus"
}, },
"volumes": [ "volumes": [
"/var/lib/mesh-broker-tls:/broker-tls:ro", "/var/lib/mesh-broker-tls:/broker-tls:ro",
@@ -62,6 +62,7 @@
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro", "/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro", "/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro",
"/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro", "/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro",
"/var/lib/mesh/mesh-controller/bus:/run/secrets/bus:ro",
"/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro", "/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro",
"/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro" "/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro"
], ],