Compare commits

..
Author SHA1 Message Date
mesh-admin 1513bbaac9 Merge pull request 'An older merge does not move a source' (#120) from fix/an-older-merge-does-not-move-a-source into main 2026-09-28 03:12:40 +00:00
jschoubben 0014984116 An older merge does not move a source
The forge announces what it finds merged, and an old merge surfacing late moved the recorded head
backwards and rebuilt everything built from that repository, once per old merge. A merge made
before the source was last seen is history; one that says nothing about when is taken as news.
The catalogue now carries when each source was last seen. The bus-records test follows #116:
a module's tools are every one under its own name.
2026-09-28 05:12:37 +02:00
mesh-admin d6e49dbd68 Merge pull request 'A base the registry already holds is not pulled from upstream again' (#119) from fix/a-mirrored-base-is-not-pulled-twice into main 2026-09-28 02:48:35 +00:00
jschoubben 3756bb3460 A base the registry already holds is not pulled from upstream again
A base is named by digest, and a digest the mesh's registry holds under the module's repository
is the same bytes whatever upstream would say. Asked on every build, the public hub's anonymous
pull limit was reached on the first merge that rebuilt a whole catalogue, and every module whose
base lives there failed on a copy it did not need.
2026-09-28 04:48:32 +02:00
mesh-admin 60be9c5360 Merge pull request 'A module hears what it consumes: its consumer is raised with the bus, and it pulls it' (#118) from fix/a-module-hears-what-it-consumes into main 2026-09-28 02:29:34 +00:00
jschoubben da31bcb11e A module hears what it consumes: its consumer is raised with the bus, and it pulls it
Every module moved onto the bus by the rollout was issued on the old one, so none had a consumer
waiting; and the grant named a push delivery a runtime's client never binds, while the pull it
does make — asking about its consumer, asking it for messages — was refused. The consumers a
module's declarations imply are now raised whenever the bus is, and the grant is the pull.
2026-09-28 04:29:32 +02:00
mesh-admin f03e7b33c9 Merge pull request 'The controller may ask any module's tool' (#117) from fix/the-controller-may-ask-a-tool into main 2026-09-28 02:21:26 +00:00
jschoubben 0baf727f36 The controller may ask any module's tool
The control plane is the way in for tool calls (novox/hq ADR 0095): a person or an agent asks
through it, so it alone may publish to every module's tool subject. The first ask on the new bus
was refused the publish.
2026-09-28 04:21:25 +02:00
mesh-admin 94dd49a968 Merge pull request 'A module serves every tool under its own name, and may answer' (#116) from fix/a-module-serves-its-own-namespace into main 2026-09-28 02:14:52 +00:00
jschoubben 83a298e7e0 A module serves every tool under its own name, and may answer
Every module that served a tool was refused the subscription on the new bus: the grant listed
tools from a manifest field no module fills, because the tools a module serves are what its code
answers and a second copy of that list would be a second source of truth. The grant is now the
module's own tool namespace; nothing else may subscribe it, a caller is still granted per tool by
name, and a module may answer what it was asked.
2026-09-28 04:14:47 +02:00
mesh-admin 220b79f5cd Merge pull request 'rollout check dials the bus the way the mesh does' (#115) from fix/the-check-dials-as-the-mesh-does into main 2026-09-28 02:05:35 +00:00
jschoubben e62201e227 rollout check dials the bus the way the mesh does
The probe connected bare, and a bus that requires TLS and a user refused it at the handshake —
so the check reported the standing server as absent. It now dials with the controller's own
credential and pin, which is the one fact the check is there to report.
2026-09-28 04:05:32 +02:00
mesh-admin 0ab9b86f0a Merge pull request 'An edge is recorded by path, like the artifact it points at' (#114) from fix/an-edge-is-recorded-by-path into main 2026-09-28 01:52:10 +00:00
jschoubben c7aabd3037 An edge is recorded by path, like the artifact it points at
The test pinned what a build stood on to the address the builder pulled from; the edge names
another module's artifact and is kept the way that artifact is (novox/hq 04-ISSUES/102).
2026-09-28 03:52:08 +02:00
mesh-admin c74d990cee Merge pull request 'A build records the bases it was handed, and the mesh reads its edges from builds' (#113) from feat/build-edges-are-recorded into main 2026-09-28 01:51:16 +00:00
jschoubben 35252af665 A build records the bases it was handed, and the mesh reads its edges from builds
Bases reach a recipe as build arguments, so the digest was never in the file the builder read
edges from: no build on the mesh recorded what it stood on, and 'build --on', the bases-first
order and the merge follow-up all walked a graph with no edges (novox/hq 04-ISSUES/131). The
builder now reports every base it resolved; the controller records them by artifact path and
reads the newest build's edges from the store, since a recorded manifest carries no build.on.
2026-09-28 03:51:14 +02:00
mesh-admin aab6ded41b Merge pull request 'One bus: the AMQP transport is gone from the controller' (#112) from feat/one-bus into main 2026-09-28 01:36:27 +00:00
jschoubben aecac5bda2 One bus: the AMQP transport is gone from the controller
The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The old
transport's consume loop, build request, tool ask, management API and account scoping are
deleted, and the bus switch with them; the controller connects to the broker seat and to
nothing else. The store-window tests keep their assertions on a bus-less fake, and the tests
that only made sense for the old transport's in-memory holding go with it.
2026-09-28 03:36:16 +02:00
mesh-admin 30362118a1 Merge pull request 'The controller follows the subject it decodes' (#111) from fix/the-controller-follows-what-it-decodes into main 2026-09-28 01:15:18 +00:00
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
54 changed files with 1223 additions and 1932 deletions
+25 -77
View File
@@ -17,12 +17,7 @@ package main
import ( import (
"context" "context"
"crypto/sha256"
"crypto/tls"
"crypto/x509"
"encoding/hex"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"net/url" "net/url"
"os" "os"
@@ -30,8 +25,6 @@ import (
"strings" "strings"
"syscall" "syscall"
amqp "github.com/rabbitmq/amqp091-go"
"github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/builder" "github.com/novox/mesh-controller/internal/builder"
"github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/link"
@@ -52,7 +45,6 @@ const usage = `mesh-builder — builds modules for the mesh
It consumes build requests and answers with what it made. Nothing is listened on and nothing It consumes build requests and answers with what it made. Nothing is listened on and nothing
is dialled except the broker. is dialled except the broker.
MESH_BROKER_AMQP where the broker is, with this builder's own credential
MESH_BROKER_FILE a file the mesh sealed to this machine holding the same MESH_BROKER_FILE a file the mesh sealed to this machine holding the same
MESH_REGISTRY host:port to publish artifacts to, when the mesh has not said MESH_REGISTRY host:port to publish artifacts to, when the mesh has not said
MESH_BINDING a file the mesh wrote saying where the artifact store is MESH_BINDING a file the mesh wrote saying where the artifact store is
@@ -135,32 +127,17 @@ 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) {
address, onNATS, err := broker.OnNATS() // **The credential names the bus, and there is one** (novox/hq ADR 0131, design 28 task 5.5).
// A credential for the mesh's bus carries user, password and fingerprint beside the address,
// and that is enough to dial it, pinned.
if !credential.onTheNewBus() {
return nil, fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL)
}
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
if err != nil { if err != nil {
return nil, err return nil, err
} }
if err := broker.MustBeOneBus(credential.URL, address); err != nil { return link.MachineOverNATS(js, on), nil
return nil, err
}
if onNATS {
js, err := broker.Dial(address)
if err != nil {
return nil, fmt.Errorf("cannot reach the bus at %s: %w", address, err)
}
return link.MachineOverNATS(js, on), nil
}
conn, err := dial(credential)
if err != nil {
// Not quoted back: the URL carries this builder's broker password.
return nil, fmt.Errorf("cannot reach the broker: %w", err)
}
channel, err := conn.Channel()
if err != nil {
conn.Close()
return nil, err
}
return link.MachineOverCurrent(conn, channel, on), nil
} }
// answer does one build and says what happened, whichever way it went. // answer does one build and says what happened, whichever way it went.
@@ -442,13 +419,8 @@ func brokerFrom() (Credential, error) {
// broker is then verified against whatever this machine already trusts. // broker is then verified against whatever this machine already trusts.
return Credential{URL: said}, nil return Credential{URL: said}, nil
} }
url := strings.TrimSpace(os.Getenv("MESH_BROKER_AMQP")) return Credential{}, fmt.Errorf(
if url == "" { "no MESH_BROKER_FILE: a build machine with no credential for the bus has nothing to build")
return Credential{}, fmt.Errorf(
"neither MESH_BROKER_FILE nor MESH_BROKER_AMQP: a builder with no broker has " +
"nothing to build")
}
return Credential{URL: url}, nil
} }
// Credential is what a build machine is given so it can reach the broker. // Credential is what a build machine is given so it can reach the broker.
@@ -460,48 +432,24 @@ 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"`
} }
// dial opens the connection, pinning the broker's certificate when there is one to pin. // onTheNewBus is whether a credential is for the bus being built: its address says so, and the
func dial(held Credential) (*amqp.Connection, error) { // mesh only ever seals such a credential with the user and password beside it.
if held.Fingerprint == "" { func (c Credential) onTheNewBus() bool { return strings.HasPrefix(strings.TrimSpace(c.URL), "nats://") }
return amqp.Dial(held.URL)
}
return amqp.DialTLS(held.URL, pinning(held.Fingerprint))
}
// pinning is a TLS configuration that trusts exactly one certificate. // natsURL is the address with this machine's credential in it, for the one dial that needs it.
// func (c Credential) natsURL() string {
// InsecureSkipVerify with a VerifyPeerCertificate is **pinning, not skipping**: the standard chain rest := strings.TrimPrefix(strings.TrimSpace(c.URL), "nats://")
// check is replaced, not removed, and what replaces it is stricter — one certificate is accepted if c.User == "" {
// rather than every certificate a public authority would sign. return "nats://" + rest
//
// Its own function so a test can drive it against a real handshake. A pin check that is only ever
// exercised through a broker is a pin check nothing tests.
func pinning(fingerprint string) *tls.Config {
return &tls.Config{
InsecureSkipVerify: true,
VerifyPeerCertificate: func(raw [][]byte, _ [][]*x509.Certificate) error {
if len(raw) == 0 {
return errors.New("the broker presented no certificate")
}
// The leaf, and in the same spelling the mesh writes it — `sha256:` and 64 hex
// characters. Comparing a bare digest against a written fingerprint never matches,
// and the failure is indistinguishable from being pointed at the wrong broker.
sum := sha256.Sum256(raw[0])
got := "sha256:" + hex.EncodeToString(sum[:])
if got != fingerprint {
return fmt.Errorf(
"this is not the broker this builder was told about\n expected %s\n "+
"got %s\nEither this mesh's broker was replaced, or this builder is "+
"being pointed at something else. Retrying will not help",
fingerprint, got)
}
return nil
},
} }
return "nats://" + c.User + ":" + c.Password + "@" + rest
} }
+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)
+4 -2
View File
@@ -15,6 +15,8 @@ import (
"strings" "strings"
"testing" "testing"
"time" "time"
"github.com/novox/mesh-controller/internal/broker"
) )
// Where a builder publishes. // Where a builder publishes.
@@ -193,9 +195,9 @@ func TestThePinIsComparedInTheSpellingTheMeshWritesIt(t *testing.T) {
} }
} }
// handshakeWith runs the builder's own pin check against an address. // handshakeWith runs the pin check the builder dials with against an address.
func handshakeWith(address, pin string) error { func handshakeWith(address, pin string) error {
conn, err := tls.Dial("tcp", address, pinning(pin)) conn, err := tls.Dial("tcp", address, broker.PinnedToFingerprint(pin))
if err != nil { if err != nil {
return err return err
} }
+7 -4
View File
@@ -17,8 +17,8 @@ import (
var aDigest = "sha256:" + strings.Repeat("e", 64) var aDigest = "sha256:" + strings.Repeat("e", 64)
// **A build is recorded by digest and path**, whatever address the builder pushed to — and only // **A build is recorded by digest and path**, whatever address the builder pushed to — what it
// what the build made is rewritten: an image the module runs from elsewhere is left where it says. // made and what it stood on both; an image the module runs from elsewhere is left where it says.
func TestABuildIsRecordedWithoutTheStoresAddress(t *testing.T) { func TestABuildIsRecordedWithoutTheStoresAddress(t *testing.T) {
manifest, _ := json.Marshal(map[string]any{ manifest, _ := json.Marshal(map[string]any{
"module": "gitea", "version": "1", "module": "gitea", "version": "1",
@@ -65,8 +65,11 @@ func TestABuildIsRecordedWithoutTheStoresAddress(t *testing.T) {
if strings.Contains(string(kept.Manifest), "anchor.internal:5100") { if strings.Contains(string(kept.Manifest), "anchor.internal:5100") {
t.Errorf("the recorded manifest still carries the store's address:\n%s", kept.Manifest) t.Errorf("the recorded manifest still carries the store's address:\n%s", kept.Manifest)
} }
if kept.Against[0] != "anchor.internal:5100/mesh-tools/runtime@"+aDigest { // What the build stood on is an edge to another module's artifact, and it is recorded the way
t.Errorf("what the build stood on was rewritten: %v", kept.Against) // that artifact is: by path in the store, so the edge still names the same thing when the
// store answers at another address.
if kept.Against[0] != catalogue.ArtifactStoreScheme+"mesh-tools/runtime@"+aDigest {
t.Errorf("what the build stood on was recorded by address: %v", kept.Against)
} }
} }
+2 -2
View File
@@ -37,13 +37,13 @@ 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
} }
defer server.Close() defer server.Close()
answer, err := link.Ask(ctx, server.Channel(), module, tool, arguments, *wait) answer, err := link.Ask(ctx, server.Bus(), module, tool, arguments, *wait)
if err != nil { if err != nil {
return err return err
} }
+106 -89
View File
@@ -2,8 +2,6 @@ package main
import ( import (
"context" "context"
"crypto/rand"
"encoding/base64"
"encoding/json" "encoding/json"
"errors" "errors"
"flag" "flag"
@@ -31,6 +29,51 @@ 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
}
against, err := open.inventory.BuiltAgainst(ctx)
if err != nil {
return err
}
var on []inventory.Entry
for _, e := range held {
if standsOnModule(e, base, against) {
on = append(on, e)
}
}
if len(on) == 0 {
fmt.Printf("nothing the mesh holds stands on %s\n", base)
return nil
}
on = orderByBases(on, against)
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 +89,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 +97,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 +112,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 {
@@ -81,8 +132,9 @@ func buildCommand(ctx context.Context, args []string) error {
// //
// **By digest and path, never by where it was pushed** (novox/hq 04-ISSUES/102). The builder // **By digest and path, never by where it was pushed** (novox/hq 04-ISSUES/102). The builder
// says `<registry>:<port>/<module>/<artifact>@sha256:…`; the mesh records the artifact-store // says `<registry>:<port>/<module>/<artifact>@sha256:…`; the mesh records the artifact-store
// reference and composes the store's address back in where a reference is used. `against` is kept // reference and composes the store's address back in where a reference is used. `against` — what
// as announced: it is what the build stood on as the builder saw it, and the catalogue's edge. // the build stood on, the catalogue's edge — is recorded the same way, so an edge names a module's
// artifact and not the machine it was pulled from.
func buildFrom(result link.BuildResult) inventory.Build { func buildFrom(result link.BuildResult) inventory.Build {
kept := inventory.Build{ kept := inventory.Build{
ID: result.ID, Repository: result.Repository, Ref: result.Ref, ID: result.ID, Repository: result.Repository, Ref: result.Ref,
@@ -92,7 +144,10 @@ func buildFrom(result link.BuildResult) inventory.Build {
// edges, and it is not always listening when a build happens — on a fresh mesh it cannot // edges, and it is not always listening when a build happens — on a fresh mesh it cannot
// be, for exactly the modules it needs most. Keeping them is what makes a replay able to // be, for exactly the modules it needs most. Keeping them is what makes a replay able to
// rebuild the graph rather than a list of names. // rebuild the graph rather than a list of names.
Path: result.Path, Against: result.Against, Path: result.Path,
}
for _, ref := range result.Against {
kept.Against = append(kept.Against, catalogue.Recorded(ref))
} }
var announced []inventory.Artifact var announced []inventory.Artifact
for _, made := range result.Made { for _, made := range result.Made {
@@ -206,85 +261,45 @@ func builderCommand(ctx context.Context, args []string) error {
} }
name := positionals[1] name := positionals[1]
management, err := broker.ManagementFromEnvironment() // **The build machine's credential is a module's credential** (novox/hq ADR 0131, design 28
// task 5.5): minted into the mesh's records and sealed to the machine as the builder module's
// broker secret, usable at the next push — the same act `module issue` performs, and the same
// account the composed user list carries. Nothing is created on a server; the bus reads the list.
open, err := openStores(ctx)
if err != nil { if err != nil {
return err return err
} }
defer open.Close()
// The same shape of secret a token carries: enough entropy that guessing is not a strategy, inv := open.inventory
// and safe to put in a URL because that is where it goes. shelf, err := inv.Catalogue(ctx)
raw := make([]byte, 32) if err != nil {
if _, err := rand.Read(raw); err != nil {
return err return err
} }
password := base64.RawURLEncoding.EncodeToString(raw) m, known := shelf[*module]
if err := management.CreateBuilderAccount(ctx, name, password); err != nil { if !known {
return fmt.Errorf("%s is not in the catalogue; `module add` it first", *module)
}
fmt.Printf("build machine %s: ", name)
node := *forNode
if node == "" {
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
for _, e := range entries {
if e.Manifest.Module == *module && len(e.On) > 0 {
node = e.On[0]
}
}
}
if node == "" {
return fmt.Errorf("%s is assigned nowhere; `assign <machine> %s` first, or say --node", *module, *module)
}
address, err := broker.BusAddress()
if err != nil {
return err return err
} }
return issueOnTheNewBus(ctx, inv, m, node, address)
fmt.Printf("broker account %s created, scoped to the %s queue and the %s exchange\n\n",
name, link.BuildQueue, link.Exchange)
if *forNode != "" {
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err)
}
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode)
if err != nil {
return err
}
// The URL and what verifies the broker, together. A mesh's broker presents a certificate
// of the mesh's own, which is in no public trust store — so a URL on its own reaches only
// a broker somebody else vouches for, and the connection fails at TLS with an error about
// an unknown authority rather than about a missing pin.
//
// **The same two facts a node's token carries** (novox/hq ADR 0004), delivered the same
// way: out of band relative to the broker, so what is trusted does not come from the thing
// being trusted.
held, err := json.Marshal(struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
}{
URL: fmt.Sprintf("amqps://%s:%s@%s/", name, password, brokerAddr),
Fingerprint: known.Fingerprint,
})
if err != nil {
return err
}
if err := inv.AcceptSecretForModule(ctx, *forNode, *module, "broker", string(held)); err != nil {
return err
}
// Not printed. It is sealed to that machine and the mesh cannot read it back, which is
// the whole point — printing it here would put the one copy that matters on a terminal.
fmt.Printf(" sealed to %s, for the %s module. It arrives with the next push.\n",
*forNode, *module)
fmt.Printf(" run `push %s` to send it\n", *forNode)
return nil
}
// The whole line only when the address is known. A URL with a placeholder where the host
// should be is a URL somebody pastes and then debugs, and the placeholder is the last thing
// they look at.
if known, err := broker.FromEnvironment(); err == nil {
fmt.Printf(" MESH_BROKER_AMQP=amqps://%s:%s@%s/\n\n", name, password, known.Address)
} else {
fmt.Printf(" the password is %s\n\n", password)
fmt.Printf(" This control plane has no %s, so it cannot say where the broker is.\n"+
" Put the password in MESH_BROKER_AMQP on the build machine.\n\n",
broker.AddressVar)
}
// Shown once, like a token, and for the same reason: what is stored is the broker's own hash
// of it, and a control plane that could show it back would be a control plane that holds it.
fmt.Println("This is the only time it is shown.")
return nil
} }
// buildBehind builds every module the mesh holds older than its source has. // buildBehind builds every module the mesh holds older than its source has.
@@ -329,6 +344,14 @@ 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).
against, err := inv.BuiltAgainst(ctx)
if err != nil {
return err
}
stale = orderByBases(stale, against)
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 +391,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 +495,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
} }
@@ -575,16 +598,10 @@ func heldBy(ctx context.Context) map[string]string {
// **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus // **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus
// the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus // the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus
// being built it dials, because a build request is a one-shot and holds nothing else. // being built it dials, because a build request is a one-shot and holds nothing else.
func askOver(server *link.Server) (link.Builders, error) { func askOver(_ *link.Server) (link.Builders, error) {
address, onNATS, err := broker.OnNATS() address, err := broker.BusAddress()
if err != nil { if err != nil {
return nil, err return nil, err
} }
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), address); err != nil { return link.BuildsOverNATS(address)
return nil, err
}
if onNATS {
return link.BuildsOverNATS(address)
}
return link.BuildsOverCurrent(server.Channel()), nil
} }
+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)
+2 -70
View File
@@ -2,8 +2,6 @@ package main
import ( import (
"context" "context"
"crypto/rand"
"encoding/base64"
"encoding/json" "encoding/json"
"errors" "errors"
"flag" "flag"
@@ -261,77 +259,11 @@ func moduleCommand(ctx context.Context, args []string) error {
// made in entirely different ways: on the bus the mesh runs on today an account is a // made in entirely different ways: on the bus the mesh runs on today an account is a
// management call, and on the bus being built it is a row the next composition writes into // management call, and on the bus being built it is a row the next composition writes into
// the server's user list (novox/hq design 25 §4). // the server's user list (novox/hq design 25 §4).
busAddress, onNATS, err := broker.OnNATS() busAddress, err := broker.BusAddress()
if err != nil { if err != nil {
return err return err
} }
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil { return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress)
return err
}
if onNATS {
return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress)
}
management, err := broker.ManagementFromEnvironment()
if err != nil {
return err
}
// The foundation owns the bus; make sure it exists before a module binds onto it.
if err := management.EnsureEventExchanges(ctx); err != nil {
return err
}
secret := make([]byte, 32)
if _, err := rand.Read(secret); err != nil {
return err
}
password := base64.RawURLEncoding.EncodeToString(secret)
account, err := management.CreateModuleAccount(ctx, *forNode, module, password, m.Emits, m.Consumes)
if err != nil {
return err
}
// A consumer's queue, with its dead-letter, is the foundation's to declare — its own account
// may not (ADR 0043). Made now, so it exists before the module binds onto it.
if len(m.Consumes) > 0 {
if err := management.EnsureModuleQueue(ctx, *forNode, module); err != nil {
return err
}
}
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err)
}
brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode)
if err != nil {
return err
}
// The URL and what verifies the broker, together — a mesh's broker presents its own
// certificate, in no public trust store, so a URL alone fails at TLS (as `builder issue`).
held, err := json.Marshal(struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
Node string `json:"node"`
Module string `json:"module"`
}{
URL: fmt.Sprintf("amqps://%s:%s@%s/", account, password, brokerAddr),
Fingerprint: known.Fingerprint,
// The node and module the account is for, so the runtime names its queue as the mesh
// scoped it (<node>.<module>.events) without a manifest having to interpolate a node.
Node: *forNode,
Module: module,
})
if err != nil {
return err
}
if err := inv.AcceptSecretForModule(ctx, *forNode, module, "broker", string(held)); err != nil {
return err
}
fmt.Printf("broker account %s created for %s, scoped to what it emits and consumes\n",
account, module)
fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n",
*forNode, *forNode)
return nil
default: default:
return fmt.Errorf("module has no %q; it has add, list, moved, forget and issue", args[0]) return fmt.Errorf("module has no %q; it has add, list, moved, forget and issue", args[0])
-10
View File
@@ -10,7 +10,6 @@ import (
"github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/broker"
"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/token" "github.com/novox/mesh-controller/internal/token"
) )
@@ -258,15 +257,6 @@ func tokenCommand(ctx context.Context, args []string) error {
// chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can // chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can
// exist before it does. The one-time secret IS the password, so a node's first connection is // exist before it does. The one-time secret IS the password, so a node's first connection is
// already authenticated and enrolment is what happens over it. // already authenticated and enrolment is what happens over it.
if management, err := broker.ManagementFromEnvironment(); err == nil {
if err := management.CreateNodeAccount(ctx, issued.Node.Name, issued.Secret); err != nil {
return err
}
fmt.Printf("broker account %s created, scoped to %s and the %s exchange\n\n",
issued.Node.Name, link.QueueFor(issued.Node.Name), link.Exchange)
} else if !errors.Is(err, broker.ErrNotConfigured) {
return err
}
made := token.Token{Node: issued.Node.Name, Signer: key.Public, Secret: issued.Secret, made := token.Token{Node: issued.Node.Name, Signer: key.Public, Secret: issued.Secret,
Adopted: issued.Node.Adopted} Adopted: issued.Node.Adopted}
+108
View File
@@ -0,0 +1,108 @@
package main
import (
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// entry is a module as the catalogue holds it: built, so its manifest carries no `build` any more.
func entry(module string, _ ...string) inventory.Entry {
return inventory.Entry{Manifest: catalogue.Manifest{Module: module}}
}
// stoodOn is what each module's newest build recorded it was handed.
func stoodOn(edges map[string][]string) map[string][]string {
out := map[string][]string{}
for module, bases := range edges {
for _, b := range bases {
out[module] = append(out[module], catalogue.ArtifactStoreScheme+b+"/runtime@sha256:"+strings.Repeat("0", 64))
}
}
return out
}
// 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"), entry("runtime"), entry("other"), entry("base")}
edges := stoodOn(map[string][]string{"app": {"runtime"}, "runtime": {"base"}})
got := orderByBases(in, edges)
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")}, stoodOn(map[string][]string{"app": {"elsewhere"}}))
if len(got) != 1 {
t.Fatalf("a dependency outside the set changed the set: %v", got)
}
}
// A module registered from its manifest and never built still names its bases there; once built,
// the recorded edge is what says so. Both are read, and a module never stands on itself.
func TestWhatStandsOnAModuleIsReadFromItsBuildOrItsManifest(t *testing.T) {
built := entry("gitea")
edges := stoodOn(map[string][]string{"gitea": {"mesh-tools"}})
if !standsOnModule(built, "mesh-tools", edges) {
t.Fatal("a recorded edge was not read")
}
if standsOnModule(built, "gitea", edges) || standsOnModule(built, "postgres", edges) {
t.Fatal("an edge was invented")
}
fresh := inventory.Entry{Manifest: catalogue.Manifest{Module: "plex", Build: &catalogue.Build{
On: []catalogue.BuildsOn{{Arg: "RUNTIME_BASE", Module: "mesh-tools", Artifact: "runtime"}},
}}}
if !standsOnModule(fresh, "mesh-tools", nil) {
t.Fatal("a manifest's own base was not read")
}
}
// 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)
}
}
}
// A merge made before the source was last seen is history: it does not move the source, and a
// merge that says nothing about when it was made is taken as news.
func TestAMergeOlderThanTheLastLookIsHistory(t *testing.T) {
seen := time.Date(2026, 9, 28, 3, 0, 0, 0, time.UTC)
if !isHistory("2026-09-28T02:00:00Z", seen) {
t.Fatal("an older merge was taken as news")
}
if isHistory("2026-09-28T04:00:00Z", seen) {
t.Fatal("a newer merge was taken as history")
}
if isHistory("", seen) || isHistory("2026-09-28T02:00:00Z", time.Time{}) {
t.Fatal("a merge or a source with no time on it was refused")
}
}
+100 -31
View File
@@ -38,6 +38,27 @@ 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, err := broker.BusAddress()
if err != nil {
return nil, err
}
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 {
@@ -61,11 +82,6 @@ func serve(ctx context.Context) error {
} }
fmt.Printf("signing as %s\n", key.Fingerprint()[:16]) fmt.Printf("signing as %s\n", key.Fingerprint()[:16])
management, err := broker.ManagementFromEnvironment()
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
return err
}
// Where the broker is and what to expect there, so a node can be told how to come back // Where the broker is and what to expect there, so a node can be told how to come back
// without a person and a new token. // without a person and a new token.
known, err := broker.FromEnvironment() known, err := broker.FromEnvironment()
@@ -80,17 +96,10 @@ func serve(ctx context.Context) error {
// **Which bus this mesh is on, read once** (novox/hq ADR 0116 step 5). Both clients ship; both // **Which bus this mesh is on, read once** (novox/hq ADR 0116 step 5). Both clients ship; both
// being live is refused, because a mesh half on each is one where a declaration goes out on one // being live is refused, because a mesh half on each is one where a declaration goes out on one
// and the report comes back on the other, and every component logs success while it happens. // and the report comes back on the other, and every component logs success while it happens.
busAddress, onNATS, err := broker.OnNATS()
if err != nil {
return err
}
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
return err
}
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known, work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known,
OnNATS: onNATS} OnNATS: true}
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 +109,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 +162,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 +275,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 +336,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 +419,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 +632,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 +643,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 +708,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 +743,75 @@ 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 return err
} }
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n", if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
address, len(names)) return err
}
// And how every module hears what it consumes. Derived from the same records the user list is
// composed from, so a module the mesh grants a consumer's subjects has that consumer waiting.
// Done on every raise, not only when a credential is issued: every module moved onto this bus
// by the rollout was issued on the old one, and came up with nothing to bind to (2026-09-28).
records, err := inv.BusRecords(ctx)
if err != nil {
return err
}
users, err := broker.Users(records)
if err != nil {
return err
}
hearing := 0
for _, p := range users {
consumer, needed := broker.ConsumerFor(p)
if !needed {
continue
}
if err := js.EnsureConsumer(consumer); err != nil {
return fmt.Errorf("how %s on %s hears what it consumes: %w", p.Module, p.Node, err)
}
hearing++
}
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, and %d module(s) "+
"can hear what they consume\n", broker.BareAddress(address), len(names), hearing)
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
}
+77 -3
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"os"
"strings" "strings"
"time" "time"
@@ -33,7 +34,7 @@ import (
// 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 — measured on 2026-09-27, when
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up. // a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
const rolloutUsage = "rollout check | rollout mint [--again] | 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 {
@@ -41,6 +42,8 @@ func rolloutCommand(ctx context.Context, args []string) error {
return rolloutCheck(ctx) return rolloutCheck(ctx)
case len(args) == 1 && args[0] == "mint": case len(args) == 1 && args[0] == "mint":
return rolloutMint(ctx, false) 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": 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 // 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, // before anything was pushed. Afterwards nothing that received the old one still works,
@@ -124,9 +127,13 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines
if address != "" { if address != "" {
// One dial, briefly. "Is it answering" is the one fact records cannot hold, and a mesh about // One dial, briefly. "Is it answering" is the one fact records cannot hold, and a mesh about
// to move onto a server that is not there should hear it here rather than afterwards. // to move onto a server that is not there should hear it here rather than afterwards.
if conn, err := nats.Connect(broker.BareAddress(address), nats.Timeout(5*time.Second)); err == nil { //
// **Dialled the way the mesh dials it** — credential and pin — because a bare connect to a
// bus that requires TLS and a user fails at the handshake, and the check then reported a
// standing server as absent (seen live, 2026-09-28).
if js, err := broker.Dial(address, nats.Timeout(5*time.Second)); err == nil {
state.ServerStanding = true state.ServerStanding = true
conn.Close() js.Close()
} }
} }
@@ -340,6 +347,14 @@ func rolloutMint(ctx context.Context, again bool) error {
machines++ machines++
case broker.KindModule: 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] m, inShelf := shelf[p.Module]
if !inShelf { if !inShelf {
skipped++ skipped++
@@ -379,3 +394,62 @@ func providesBus(m catalogue.Manifest) bool {
} }
return false 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
}
+164
View File
@@ -6,7 +6,9 @@ import (
"flag" "flag"
"fmt" "fmt"
"strings" "strings"
"time"
"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/link" "github.com/novox/mesh-controller/internal/link"
) )
@@ -218,3 +220,165 @@ 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
}
// **A merge older than the last look at the source is history, not a move.** The forge
// announces what it finds merged, and an old merge surfacing late would otherwise move the
// recorded head backwards and rebuild everything built from that repository, once per old
// merge (2026-09-28).
if isHistory(m.MergedAt, e.Source.Seen) {
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
}
against, err := inv.BuiltAgainst(ctx)
if err != nil {
return notNow(err)
}
ordered := orderByBases(moved, against)
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, against) {
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 stood
// on another's artifact 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.
//
// `against` is what each module's newest build stood on (inventory.BuiltAgainst): the edges are
// derived from builds, not declared, because a recorded manifest no longer carries `build.on`.
func orderByBases(entries []inventory.Entry, against map[string][]string) []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
for _, base := range entries {
if base.Manifest.Module != name && inSet[base.Manifest.Module] && standsOnModule(e, base.Manifest.Module, against) {
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, against map[string][]string) bool {
for _, e := range entries {
if standsOnModule(e, module, against) {
return true
}
}
return false
}
// standsOnModule is whether an entry's build stood on the named module: by what its newest build
// recorded it was handed (`artifact-store://<module>/<artifact>@…`, the module's own artifact), or
// — for a module registered from a manifest and not yet built — by the base its manifest names.
func standsOnModule(e inventory.Entry, module string, against map[string][]string) bool {
if e.Manifest.Module == module {
return false
}
if e.Manifest.Build != nil {
for _, on := range e.Manifest.Build.On {
if on.Module == module {
return true
}
}
}
prefix := catalogue.ArtifactStoreScheme + module + "/"
for _, ref := range against[e.Manifest.Module] {
if strings.HasPrefix(ref, prefix) {
return true
}
}
return false
}
// isHistory is whether a merge made at mergedAt predates the last time the source was seen. A merge
// with no time on it is taken as news: refusing it would silence a forge that says less.
func isHistory(mergedAt string, seen time.Time) bool {
if mergedAt == "" || seen.IsZero() {
return false
}
at, err := time.Parse(time.RFC3339, mergedAt)
if err != nil {
return false
}
return at.Before(seen)
}
+1 -2
View File
@@ -4,7 +4,7 @@ go 1.26.0
require ( require (
github.com/jackc/pgx/v5 v5.10.0 github.com/jackc/pgx/v5 v5.10.0
github.com/rabbitmq/amqp091-go v1.14.0 github.com/nats-io/nats.go v1.54.0
golang.org/x/crypto v0.57.0 golang.org/x/crypto v0.57.0
) )
@@ -13,7 +13,6 @@ require (
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/klauspost/compress v1.20.0 // indirect github.com/klauspost/compress v1.20.0 // indirect
github.com/nats-io/nats.go v1.54.0 // indirect
github.com/nats-io/nkeys v0.4.16 // indirect github.com/nats-io/nkeys v0.4.16 // indirect
github.com/nats-io/nuid v1.0.1 // indirect github.com/nats-io/nuid v1.0.1 // indirect
golang.org/x/net v0.58.0 // indirect golang.org/x/net v0.58.0 // indirect
-4
View File
@@ -19,15 +19,11 @@ github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek=
github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M= golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
-67
View File
@@ -1,67 +0,0 @@
package broker_test
import (
"regexp"
"testing"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/link"
)
// The names are written twice, so a test keeps them agreeing.
//
// `link` imports `broker`, so `broker` cannot import `link` — the queue and exchange names
// therefore exist in both. A scoped account naming a queue nothing publishes to produces a
// builder that takes no work and says nothing about why, which is the worst kind of silence.
//
// An external test package, because it may import both without either importing the other.
func TestTheNamesTheBrokerScopesAreTheNamesTheLinkUses(t *testing.T) {
for _, agreed := range []struct {
what string
scoped string
actually string
}{
{"the build queue", broker.BuildQueueName, link.BuildQueue},
{"the exchange", broker.ExchangeName, link.Exchange},
{"a node's queue", broker.QueueFor("somewhere"), link.QueueFor("somewhere")},
} {
if agreed.scoped != agreed.actually {
t.Errorf("%s: the broker scopes %q and the link uses %q",
agreed.what, agreed.scoped, agreed.actually)
}
}
}
func TestABuilderMayWriteToAReplyQueueAndReadNoNodesDeclarations(t *testing.T) {
// The scoping, checked as patterns rather than by connecting: what it may write must include
// the queue an asker actually waits on, and what it may read must not include any node's.
//
// This exists because the first version scoped writes to `amq.gen-*` — the name one broker
// happens to generate — and the builder built, could not answer, and the connection simply
// closed saying only "not allowed to publish to exchange \'\'". Answering through the
// exchange is what removed the need for any of that.
write := regexp.MustCompile("^" + regexp.QuoteMeta(broker.ExchangeName) + "$")
if !write.MatchString(broker.ExchangeName) {
t.Error("a builder may not write to the exchange, so it can take work and never answer")
}
// And not the default exchange, where permission is per exchange rather than per queue — a
// builder allowed to use it could publish into any node's queue.
//
// Confirmed against a real broker as well, and worth recording how that nearly went wrong:
// an unconfirmed publish is asynchronous, so a refusal arrives as a channel close afterwards
// and a naive check reports success. With publisher confirms the broker's refusal is
// immediate. **A negative security assertion made against an asynchronous call is not an
// assertion.**
if write.MatchString("") {
t.Error("a builder may publish to the default exchange, and so into any node's queue")
}
read := regexp.MustCompile("^" + regexp.QuoteMeta(broker.BuildQueueName) + "$")
if !read.MatchString(broker.BuildQueueName) {
t.Error("a builder may not read the build queue")
}
if read.MatchString(link.QueueFor("someone-else")) {
// A build machine is not a node, and a node's queue carries its declarations.
t.Error("a builder may read another machine's declarations")
}
}
-13
View File
@@ -194,16 +194,3 @@ func TestTheAddressPortFollowsThePortTwin(t *testing.T) {
t.Fatalf("the address is %q; the node put the bus on 5679", b.Address) t.Fatalf("the address is %q; the node put the bus on 5679", b.Address)
} }
} }
func TestTheManagementPortFollowsThePortTwin(t *testing.T) {
t.Setenv(ManagementVar, "http://guest:guest@127.0.0.1:15672")
t.Setenv(ManagementVar+"_FILE", "")
t.Setenv(ManagementVar+"_PORT", "15673")
m, err := ManagementFromEnvironment()
if err != nil {
t.Fatal(err)
}
if m.base.Host != "127.0.0.1:15673" {
t.Fatalf("the management API is at %q; the node put it on 15673", m.base.Host)
}
}
+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).
-366
View File
@@ -1,366 +0,0 @@
package broker
import (
"bytes"
"context"
"encoding/json"
"fmt"
"github.com/novox/mesh-controller/internal/envfile"
"io"
"net/http"
"net/url"
"regexp"
"strings"
"time"
)
// The mesh runs the broker, so there is no chicken-and-egg in a node needing an account before it
// can connect: the account is created when the token is issued, and the one-time secret in that
// token IS the password. A node's first connection is already authenticated, and enrolment is
// what happens over it.
//
// novox/hq ADR 0004's *a node holds its own identity and nothing else* is why the account is per
// node rather than shared. A shared enrolment account would let any node consume another's queue,
// which is the shared-credential fault that record exists to remove, reappearing at the transport.
// ManagementVar holds the broker's management API, credentials included.
const ManagementVar = "MESH_BROKER_MANAGEMENT"
// safeName is what a node may be called at the broker.
//
// The name goes into a URL path and into permission patterns, which are regular expressions. A
// name carrying a `.` or a `*` would silently widen what that node may reach — so it is
// constrained here rather than escaped later, because an escape that is forgotten once is a node
// reading everybody's queues.
var safeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,62}$`)
// Management is the broker's administrative interface.
type Management struct {
base *url.URL
client *http.Client
}
// ManagementFromEnvironment reads where the management API is, if it is configured.
//
// On the port MESH_BROKER_MANAGEMENT_PORT names when the node moved it (novox/hq 04-ISSUES/102);
// the URL's own port otherwise.
func ManagementFromEnvironment() (*Management, error) {
raw, err := envfile.Placed(ManagementVar)
if err != nil {
return nil, err
}
if raw == "" {
return nil, ErrNotConfigured
}
base, err := url.Parse(raw)
if err != nil || base.Host == "" {
// The value carries a password, so it is not quoted back.
return nil, fmt.Errorf("%s is not a usable URL", ManagementVar)
}
return &Management{base: base, client: &http.Client{Timeout: 15 * time.Second}}, nil
}
// QueueFor is the queue a node consumes from. One per node, named after it.
func QueueFor(node string) string { return "node." + node }
// ExchangeName is where nodes publish what they have to say. One exchange, and the control plane
// is the only consumer behind it (novox/hq ADR 0006 — one consumer, so two cannot silently split
// the traffic between them).
const ExchangeName = "mesh"
// BuildQueueName is where build work waits. Duplicated from `link` rather than imported, for the
// same reason QueueFor above is: this package must not depend on the one that uses it, and a
// constant that differed would be caught by the test that asserts they agree.
const BuildQueueName = "builds"
// The events bus (novox/hq ADR 0042): one topic exchange every event rides, a second for tool
// RPC kept apart, and a dead-letter home for a poison event. The foundation owns these — a module's
// account cannot declare them, only bind its own queue to the events one.
const (
EventsExchangeName = "mesh.events"
RPCExchangeName = "mesh.rpc"
DeadExchangeName = "mesh.events.dead"
)
// ModuleQueueFor is the durable queue a module consumes its events from — one per module per node
// (novox/hq ADR 0042), named so the account that may read it is exactly this module's.
func ModuleQueueFor(node, module string) string { return node + "." + module + ".events" }
// modulePermissions is a module's authority on the bus, derived from its manifest (novox/hq
// ADR 0043): what it consumes and what it emits, and nothing else. Pure, so the scope is tested as
// patterns without a broker — the way a builder's is.
//
// A note on the limit: the broker's write permission is per exchange, not per routing key (LavinMQ
// has no topic permissions), so an emitting module is granted the events exchange whole. ADR 0042's
// origin reservation — a module publishes only under `module.<self>.*` — is stamped by the sdk, not
// enforced here; that gap is the broker's, and is recorded rather than hidden. A pure consumer like
// the audit logger is unaffected: it is granted no write to the exchange at all.
func modulePermissions(node, module string, emits, consumes []string) (configure, write, read string) {
queue := regexp.QuoteMeta(ModuleQueueFor(node, module))
events := regexp.QuoteMeta(EventsExchangeName)
rpc := regexp.QuoteMeta(RPCExchangeName)
// A module serves each of its tools on its own queue, namespaced by the module (novox/hq
// ADR 0047) — serve.<module>.<tool> — so the account may declare, bind and read exactly its own,
// and no other module's.
serve := "serve\\." + regexp.QuoteMeta(module) + "\\..*"
// Declare its own events queue and its own tool serve queues.
configure = "^(" + queue + "|" + serve + ")$"
// Write to bind its queue and serve queues (binding is a write on the queue), and to the RPC
// exchange to publish replies (ADR 0047: replies ride mesh.rpc, never the default exchange, which
// would let it publish into any queue). To the events exchange only if it emits.
writes := []string{queue, serve, rpc}
if len(emits) > 0 {
writes = append(writes, events)
}
write = "^(" + strings.Join(writes, "|") + ")$"
// Read its own queue and serve queues to consume them, and the RPC exchange to bind its serve
// queues onto. The events exchange to bind onto only if it consumes.
reads := []string{queue, serve, rpc}
if len(consumes) > 0 {
reads = append(reads, events)
}
read = "^(" + strings.Join(reads, "|") + ")$"
return configure, write, read
}
// CreateModuleAccount gives an assigned module its own broker account, scoped by what it emits and
// consumes (novox/hq ADR 0043). The account name carries the node so the same module on two machines
// holds two accounts, each sealed to its own; the permissions carry the module so one module cannot
// read another's queue. Generic — the builder is one instance of this rule, not a separate kind.
func (m *Management) CreateModuleAccount(ctx context.Context, node, module, password string, emits, consumes []string) (string, error) {
if !safeName.MatchString(node) {
return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", node)
}
if !safeName.MatchString(module) {
return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", module)
}
account := node + "-" + module
if !safeName.MatchString(account) {
return "", fmt.Errorf("%q is not a usable broker account name", account)
}
if err := m.put(ctx, "/api/users/"+url.PathEscape(account),
map[string]string{"password": password, "tags": ""}); err != nil {
return "", fmt.Errorf("cannot create the broker account for %s on %s: %w", module, node, err)
}
configure, write, read := modulePermissions(node, module, emits, consumes)
if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(account), map[string]string{
"configure": configure, "write": write, "read": read,
}); err != nil {
return "", fmt.Errorf("cannot scope the broker account for %s on %s: %w", module, node, err)
}
return account, nil
}
// EnsureModuleQueue declares a consuming module's queue with its dead-letter exchange, idempotently.
// The foundation declares it because a scoped module account may not: the broker refuses a queue with
// a dead-letter exchange to a non-administrator (novox/hq ADR 0043), so a consumer passively checks
// the queue the mesh made rather than declaring its own.
func (m *Management) EnsureModuleQueue(ctx context.Context, node, module string) error {
queue := ModuleQueueFor(node, module)
if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(queue), map[string]any{
"durable": true,
"arguments": map[string]any{"x-dead-letter-exchange": DeadExchangeName},
}); err != nil {
return fmt.Errorf("cannot declare the queue for %s on %s: %w", module, node, err)
}
return nil
}
// EnsureEventExchanges declares the bus's exchanges and the dead-letter home, idempotently. The
// foundation owns them (a module's account may not declare an exchange), and a dead-letter exchange
// with no queue behind it drops what it receives — so a durable queue bound to `#` retains a poison
// event for inspection, which is the whole reason the trail exists.
func (m *Management) EnsureEventExchanges(ctx context.Context) error {
for _, exchange := range []string{EventsExchangeName, RPCExchangeName, DeadExchangeName} {
if err := m.put(ctx, "/api/exchanges/%2f/"+url.PathEscape(exchange),
map[string]any{"type": "topic", "durable": true}); err != nil {
return fmt.Errorf("cannot declare the %s exchange: %w", exchange, err)
}
}
if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(DeadExchangeName),
map[string]any{"durable": true}); err != nil {
return fmt.Errorf("cannot declare the dead-letter queue: %w", err)
}
if err := m.post(ctx, "/api/bindings/%2f/e/"+url.PathEscape(DeadExchangeName)+
"/q/"+url.PathEscape(DeadExchangeName), map[string]string{"routing_key": "#"}); err != nil {
return fmt.Errorf("cannot bind the dead-letter queue: %w", err)
}
return nil
}
// CreateNodeAccount gives a node its own broker account, with the token's secret as the password.
//
// Scoped so a node can reach its own queue and the one exchange, and nothing else. The patterns
// are anchored: a node called `laptop` must not be able to read `laptop-of-somebody-else`.
func (m *Management) CreateNodeAccount(ctx context.Context, node, password string) error {
if !safeName.MatchString(node) {
return fmt.Errorf(
"%q cannot be a broker account name: it becomes part of a permission pattern, so it "+
"is lower-case letters, digits and dashes", node)
}
if err := m.put(ctx, "/api/users/"+url.PathEscape(node),
map[string]string{"password": password, "tags": ""}); err != nil {
return fmt.Errorf("cannot create the broker account for %s: %w", node, err)
}
queue := regexp.QuoteMeta(QueueFor(node))
if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(node), map[string]string{
"configure": "^" + queue + "$",
"write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + queue + ")$",
"read": "^" + queue + "$",
}); err != nil {
return fmt.Errorf("cannot scope the broker account for %s: %w", node, err)
}
return nil
}
// CreateBuilderAccount scopes an account to taking build work and answering it.
//
// **A build machine is not a node**, and giving it a node's account would let it read another
// machine's declarations. What it needs is narrower and different: read the build queue, and
// write to the exchange and to whatever temporary queue an asker is waiting on.
//
// The reply queues are the reason `write` is not simply the exchange. `RequestBuild` declares an
// exclusive queue with a generated name and waits on it, so a builder that could not write to it
// could take work and never answer — which is the failure that looks like a builder that is not
// running.
func (m *Management) CreateBuilderAccount(ctx context.Context, name, password string) error {
if !safeName.MatchString(name) {
return fmt.Errorf(
"%q cannot be a broker account name: it becomes part of a permission pattern, so it "+
"is lower-case letters, digits and dashes", name)
}
if err := m.put(ctx, "/api/users/"+url.PathEscape(name),
map[string]string{"password": password, "tags": ""}); err != nil {
return fmt.Errorf("cannot create the broker account for %s: %w", name, err)
}
builds := regexp.QuoteMeta(BuildQueueName)
if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(name), map[string]string{
// It declares the build queue, because whichever builder starts first must be able to —
// and a queue nobody may declare is a queue that exists only if the control plane has
// already run, which makes the order they start in matter.
"configure": "^" + builds + "$",
// Two exchanges, and nothing else. **Not the default exchange**: permission there is
// granted per exchange rather than per queue, so a builder allowed to use it could
// publish into any node's queue — the privilege a build machine most obviously should
// not have. Answers go through the node exchange; announcing what was built goes through
// the events exchange, which is a different act with a different audience (novox/hq
// ADR 0072). A builder that could answer and not announce would leave the module graph
// knowing less than the registry does.
"write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + regexp.QuoteMeta(EventsExchangeName) + ")$",
// The build queue and nothing else. Not another machine's declarations.
"read": "^" + builds + "$",
}); err != nil {
return fmt.Errorf("cannot scope the broker account for %s: %w", name, err)
}
return nil
}
// RemoveNodeAccount withdraws a node's access.
func (m *Management) RemoveNodeAccount(ctx context.Context, node string) error {
if !safeName.MatchString(node) {
return fmt.Errorf("%q is not a broker account name", node)
}
return m.do(ctx, http.MethodDelete, "/api/users/"+url.PathEscape(node), nil)
}
// Accounts lists the broker's users, so a picture can be read from the system rather than assumed
// (novox/hq ADR 0018).
func (m *Management) Accounts(ctx context.Context) ([]string, error) {
body, err := m.get(ctx, "/api/users")
if err != nil {
return nil, err
}
var users []struct {
Name string `json:"name"`
}
if err := json.Unmarshal(body, &users); err != nil {
return nil, err
}
names := make([]string, 0, len(users))
for _, u := range users {
names = append(names, u.Name)
}
return names, nil
}
func (m *Management) put(ctx context.Context, path string, body any) error {
return m.do(ctx, http.MethodPut, path, body)
}
func (m *Management) post(ctx context.Context, path string, body any) error {
return m.do(ctx, http.MethodPost, path, body)
}
func (m *Management) get(ctx context.Context, path string) ([]byte, error) {
request, err := m.request(ctx, http.MethodGet, path, nil)
if err != nil {
return nil, err
}
response, err := m.client.Do(request)
if err != nil {
return nil, err
}
defer response.Body.Close()
if response.StatusCode >= 300 {
return nil, fmt.Errorf("the broker's management API answered %s to GET %s",
response.Status, path)
}
return io.ReadAll(io.LimitReader(response.Body, 1<<20))
}
func (m *Management) do(ctx context.Context, method, path string, body any) error {
request, err := m.request(ctx, method, path, body)
if err != nil {
return err
}
response, err := m.client.Do(request)
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode >= 300 {
detail, _ := io.ReadAll(io.LimitReader(response.Body, 4096))
return fmt.Errorf("the broker's management API answered %s to %s %s: %s",
response.Status, method, path, strings.TrimSpace(string(detail)))
}
return nil
}
func (m *Management) request(ctx context.Context, method, path string, body any) (*http.Request, error) {
var payload io.Reader
if body != nil {
raw, err := json.Marshal(body)
if err != nil {
return nil, err
}
payload = bytes.NewReader(raw)
}
// Path joined by hand rather than through url.Parse: %2f is the default vhost and must reach
// the broker still encoded. Parsing would decode it to a slash and address a different route.
target := strings.TrimSuffix(m.base.String(), "/")
if user := m.base.User; user != nil {
target = strings.TrimSuffix(m.base.Scheme+"://"+m.base.Host, "/")
}
request, err := http.NewRequestWithContext(ctx, method, target+path, payload)
if err != nil {
return nil, err
}
if user := m.base.User; user != nil {
password, _ := user.Password()
request.SetBasicAuth(user.Username(), password)
}
if body != nil {
request.Header.Set("Content-Type", "application/json")
}
return request, nil
}
-96
View File
@@ -1,96 +0,0 @@
package broker_test
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"regexp"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/broker"
)
// The audit logger consumes everything and emits nothing. Its account must let it declare and read
// its own queue and read the events exchange to bind onto — and must not reach another module's
// queue, nor grant any write to the events exchange (novox/hq ADR 0043).
func TestAConsumerReadsItsOwnQueueAndTheEventsExchangeAndNoOthers(t *testing.T) {
// Rebuilt through CreateModuleAccount's own path by asking for the same scope it would apply.
// The queue this module reads:
mine := broker.ModuleQueueFor("anchor", "audit-logger")
other := broker.ModuleQueueFor("anchor", "plex")
read := scope(t, "read", "anchor", "audit-logger", nil, []string{"#"})
if !read.MatchString(mine) {
t.Error("the audit logger may not read its own queue, so it consumes nothing")
}
if !read.MatchString(broker.EventsExchangeName) {
t.Error("the audit logger may not read the events exchange, so it cannot bind onto it")
}
if read.MatchString(other) {
t.Error("the audit logger may read another module's queue")
}
// It emits nothing, so it is granted no write to the events exchange — only its own queue, to bind.
write := scope(t, "write", "anchor", "audit-logger", nil, []string{"#"})
if write.MatchString(broker.EventsExchangeName) {
t.Error("a pure consumer was granted write to the events exchange")
}
if !write.MatchString(mine) {
t.Error("the audit logger may not write to its own queue, so it cannot bind it")
}
}
// An emitter is granted the events exchange to write; a consumer is not.
func TestAnEmitterMayWriteTheEventsExchangeAndAConsumerMayNot(t *testing.T) {
emitter := scope(t, "write", "anchor", "umami", []string{"module.umami.site.created"}, nil)
if !emitter.MatchString(broker.EventsExchangeName) {
t.Error("an emitting module may not write the events exchange, so it cannot emit")
}
}
// scope reconstructs one of the three permission patterns CreateModuleAccount would apply, by
// reading it back from a captured request against a stub management API.
func scope(t *testing.T, which, node, module string, emits, consumes []string) *regexp.Regexp {
t.Helper()
pat := capturePermission(t, which, node, module, emits, consumes)
re, err := regexp.Compile(pat)
if err != nil {
t.Fatalf("the %s pattern does not compile: %v", which, err)
}
return re
}
// capturePermission runs CreateModuleAccount against a stub management API and returns the pattern
// it set for `which` ("configure"/"write"/"read"). The scope is tested where it is applied, not
// reconstructed by the test — so a change to the mapping cannot pass a test that hard-codes the old
// one.
func capturePermission(t *testing.T, which, node, module string, emits, consumes []string) string {
t.Helper()
var captured map[string]string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.HasPrefix(r.URL.Path, "/api/permissions/") {
_ = json.NewDecoder(r.Body).Decode(&captured)
}
w.WriteHeader(http.StatusNoContent)
}))
defer server.Close()
t.Setenv(broker.ManagementVar, server.URL)
m, err := broker.ManagementFromEnvironment()
if err != nil {
t.Fatalf("stub management not usable: %v", err)
}
if _, err := m.CreateModuleAccount(context.Background(), node, module, "pw", emits, consumes); err != nil {
t.Fatalf("CreateModuleAccount: %v", err)
}
if captured == nil {
t.Fatal("no permissions were set")
}
pattern, ok := captured[which]
if !ok {
t.Fatalf("no %s permission was set; got %v", which, captured)
}
return pattern
}
+52 -10
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
@@ -177,6 +181,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
for _, seat := range meshSeatsTheControllerUses { for _, seat := range meshSeatsTheControllerUses {
pub = append(pub, "mesh.seat."+seat+".accept.>") pub = append(pub, "mesh.seat."+seat+".accept.>")
} }
// Every module's tools: **the control plane is the way in** (novox/hq ADR 0095). A person
// or an agent asks through it and every question passes one process where an audit
// belongs — so it, alone among principals, may call any tool by name. The first `ask` on
// the new bus was refused the publish (2026-09-28).
pub = append(pub, "mesh.mod.*.tool.>")
// The two events it reacts to, and its ack subject on the stream they arrive from // The two events it reacts to, and its ack subject on the stream they arrive from
// (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the // (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the
@@ -244,8 +253,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
@@ -255,9 +271,13 @@ func PermissionsFor(p Principal) (Permissions, error) {
for _, e := range p.Emits { for _, e := range p.Emits {
pub = append(pub, own+".event."+e) pub = append(pub, own+".event."+e)
} }
for _, t := range p.Serves { // Every tool under its own name, not a list: the tools a module serves are what its code
sub = append(sub, own+".tool."+t) // answers, and a second copy of that list in the manifest would be a second source of
} // truth for the mesh to keep in step (2026-09-28: every module that served a tool was
// refused the subscription, because none had written the list twice). Nothing is given
// away — no other principal may subscribe this namespace, and a caller's authority is
// still granted per tool, by name, on the publish side.
sub = append(sub, own+".tool.>")
// 2. What it consumes, by the emitter's own subject — an event is addressed to its // 2. What it consumes, by the emitter's own subject — an event is addressed to its
// emitter, because the emitter's identity is the meaning (ADR 0118). // emitter, because the emitter's identity is the meaning (ADR 0118).
@@ -277,8 +297,26 @@ func PermissionsFor(p Principal) (Permissions, error) {
} }
} }
// 2c. Its own consumer, which it **pulls**: the runtime asks for the next message and is
// answered on its own inbox, so what it needs is to ask about the consumer and to ask it
// for messages — its own consumer's name, and no other's. Pulled rather than pushed
// because that is the one shape a runtime's client binds without creating anything; the
// controller and the hosts are pushed to. Named here rather than through ConsumerFor,
// which 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.
pub = append(pub,
"$JS.API.CONSUMER.INFO."+consumerStream(p)+"."+consumerDurable(p),
"$JS.API.CONSUMER.MSG.NEXT."+consumerStream(p)+"."+consumerDurable(p))
// 3. Seats it holds: full participation. // 3. Seats it holds: full participation.
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))
} }
@@ -326,9 +364,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
return Permissions{ return Permissions{
Publish: pub, Publish: pub,
Subscribe: sub, Subscribe: sub,
// Only something that serves is ever answering. A pure consumer is granted nothing here. // A module answers what it was asked — a tool call reaches it on its own namespace, so the
AllowResponses: p.Kind == KindModule && (len(p.Serves) > 0 || len(p.Holds) > 0) || // authority is bounded by having been asked — and so does the controller. A node and a
p.Kind == KindController, // person are never asked anything, and are granted nothing here.
AllowResponses: p.Kind == KindModule || p.Kind == KindController,
}, nil }, nil
} }
@@ -514,7 +553,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 {
+45 -10
View File
@@ -1,6 +1,7 @@
package broker package broker
import ( import (
"slices"
"strings" "strings"
"testing" "testing"
) )
@@ -89,17 +90,30 @@ func TestAnInboxIsScopedToItsOwner(t *testing.T) {
} }
// A responder answers on the caller's inbox, which it has no permission for. allow_responses is // A responder answers on the caller's inbox, which it has no permission for. allow_responses is
// what makes a scoped inbox workable at all — the authority is bounded by having been asked. // what makes a scoped inbox workable at all — the authority is bounded by having been asked. A
func TestOnlySomethingThatServesMayAnswer(t *testing.T) { // module is asked on its own namespace and may answer; a node and a person are never asked.
serving, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "billing", func TestOnlyWhatCanBeAskedMayAnswer(t *testing.T) {
Serves: []string{"status"}, PasswordHash: "x"}) module, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
if !serving.AllowResponses {
t.Fatal("a module serving a tool cannot answer the caller's inbox")
}
consumer, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"}) Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
if consumer.AllowResponses { if !module.AllowResponses {
t.Fatal("a pure consumer was granted the right to answer, which nothing asked it to do") t.Fatal("a module cannot answer a tool call on its own namespace")
}
node, _ := PermissionsFor(Principal{Kind: KindNode, Node: "one", PasswordHash: "x"})
if node.AllowResponses {
t.Fatal("a node was granted the right to answer, and nothing asks a node anything")
}
}
// A module serves every tool under its own name, and no other module's.
func TestAModuleServesItsOwnNamespaceAndNoOthers(t *testing.T) {
p, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "gitea", PasswordHash: "x"})
if !slices.Contains(p.Subscribe, "mesh.mod.gitea.tool.>") {
t.Fatalf("a module may not serve its own tools: %v", p.Subscribe)
}
for _, s := range p.Subscribe {
if strings.HasPrefix(s, "mesh.mod.") && !strings.HasPrefix(s, "mesh.mod.gitea.") {
t.Fatalf("a module may subscribe another's namespace: %s", s)
}
} }
} }
@@ -324,3 +338,24 @@ func admits(pattern, subject []string) bool {
} }
return len(pattern) == len(subject) return len(pattern) == len(subject)
} }
// A module pulls its own consumer — asks about it, asks it for messages — and no other module's.
func TestAModulePullsItsOwnConsumerAndNoOthers(t *testing.T) {
p, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
for _, want := range []string{"$JS.API.CONSUMER.INFO.EVENTS.one_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_audit"} {
if !slices.Contains(p.Publish, want) {
t.Errorf("a module cannot bind its own consumer: %v lacks %s", p.Publish, want)
}
}
for _, s := range p.Publish {
if strings.Contains(s, "CONSUMER.") && !strings.HasSuffix(s, ".one_audit") {
t.Errorf("a module may reach another consumer: %s", s)
}
}
for _, s := range p.Subscribe {
if strings.HasPrefix(s, "_DELIVER.") {
t.Errorf("a module is granted a push delivery it never binds: %s", s)
}
}
}
+11 -19
View File
@@ -68,24 +68,16 @@ func BareAddress(address string) string {
return "nats://" + bare return "nats://" + bare
} }
// MustBeOneBus refuses a configuration that names both buses for the mesh's own traffic. // BusAddress is where the mesh's bus is, with this process's credential, and refuses to be empty:
// // there is one bus, and a control plane without it can hold records and answer nothing (novox/hq
// **Both clients ship and that is the point; both being live is not.** The rollout moves every node // ADR 0131, design 28 task 5.5).
// at once (ADR 0116 step 5): a mesh half on each is one where a declaration goes out on one bus and func BusAddress() (string, error) {
// the report comes back on the other, and nothing anywhere says so — every component would log address, _, err := OnNATS()
// success. Refused at start, where it can be said in one sentence. if err != nil {
func MustBeOneBus(amqp, nats string) error { return "", err
if strings.TrimSpace(amqp) != "" && strings.TrimSpace(nats) != "" {
return fmt.Errorf(
"this control plane is told about both buses (%s and %s) and can only be on one. A mesh "+
"half on each is one where a declaration goes out on one and the report comes back "+
"on the other, and every component reports success while it happens. The rollout "+
"moves every node at once: unset %s to stay, or unset %s to move",
AMQPVarName, NATSVar, NATSVar, AMQPVarName)
} }
return nil if address == "" {
return "", fmt.Errorf("this control plane has no %s, so it cannot reach the mesh's bus", NATSVar)
}
return address, nil
} }
// AMQPVarName is the variable naming the bus the mesh runs on today. Named here rather than
// imported from the link package, for the one direction of dependency.
const AMQPVarName = "MESH_BROKER_AMQP"
-32
View File
@@ -1,43 +1,11 @@
package broker package broker
import ( import (
"strings"
"testing" "testing"
) )
// Which bus the mesh is on is one fact, and being told about both is refused. // Which bus the mesh is on is one fact, and being told about both is refused.
// //
// **Not a warning.** A mesh half on each bus is one where a declaration goes out on one and the
// report comes back on the other, and every component reports success while it happens — which is
// the exact failure ADR 0074 exists to catch, arriving through configuration instead of through code.
func TestBeingToldAboutBothBusesIsRefused(t *testing.T) {
err := MustBeOneBus("amqps://broker:5671/", "nats://bus:4222")
if err == nil {
t.Fatal("a control plane told about both buses was allowed to start")
}
// The remedy is in the words, because whoever reads this has to choose one and the wrong choice
// is a rollout half done.
for _, want := range []string{AMQPVarName, NATSVar, "unset"} {
if !strings.Contains(err.Error(), want) {
t.Errorf("the refusal does not mention %s: %v", want, err)
}
}
}
// One bus, or none, is ordinary. None is a control plane that publishes nothing and holds records,
// which several of its own commands are.
func TestOneBusOrNeitherIsAllowed(t *testing.T) {
for _, c := range []struct{ what, amqp, nats string }{
{"the bus the mesh runs on today", "amqps://broker:5671/", ""},
{"the bus being built", "", "nats://bus:4222"},
{"neither", "", ""},
{"neither, with whitespace for an address", " ", "\t"},
} {
if err := MustBeOneBus(c.amqp, c.nats); err != nil {
t.Errorf("%s was refused: %v", c.what, err)
}
}
}
// The controller's own credential arrives in its address, and has to be readable out of it — its user // The controller's own credential arrives in its address, and has to be readable out of it — its user
// is created by the installer at a bootstrap password, before the controller exists to mint one. // is created by the installer at a bootstrap password, before the controller exists to mint one.
+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
+13 -10
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.mod.*.tool.>", "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,23 @@ 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.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "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", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "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.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] } subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
allow_responses: { max: 1, ttl: "1m" }
} } } }
{ 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.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>"] } subscribe: { allow: ["_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} } } }
] ]
} }
+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)
+34 -15
View File
@@ -169,6 +169,8 @@ func Build(ctx context.Context, run Runner, publish Publisher,
} }
var built []catalogue.Built var built []catalogue.Built
// stoodOn is every base the build was handed, as resolved — the edges the catalogue derives.
var stoodOn []string
if manifest.Build != nil { if manifest.Build != nil {
// What this module said it stands on, answered with what this mesh actually holds. Done // What this module said it stands on, answered with what this mesh actually holds. Done
// before anything is built, so a missing base is refused in front of the person who can // before anything is built, so a missing base is refused in front of the person who can
@@ -186,11 +188,12 @@ func Build(ctx context.Context, run Runner, publish Publisher,
} }
return from, nil return from, nil
} }
args, err := standingOn(ctx, manifest, held, mirror) args, bases, err := standingOn(ctx, manifest, held, mirror)
if err != nil { if err != nil {
say("bases", "UNMET: %v", err) say("bases", "UNMET: %v", err)
return Result{}, err return Result{}, err
} }
stoodOn = bases
if len(args) > 0 { if len(args) > 0 {
say("bases", "%d resolved from what the mesh holds", len(args)/2) say("bases", "%d resolved from what the mesh holds", len(args)/2)
} }
@@ -217,7 +220,7 @@ func Build(ctx context.Context, run Runner, publish Publisher,
} }
say("done", "%s at %s — %d artifact(s) pinned", manifest.Module, short(commit), len(built)) say("done", "%s at %s — %d artifact(s) pinned", manifest.Module, short(commit), len(built))
return Result{Manifest: resolved, Commit: commit, Built: built, return Result{Manifest: resolved, Commit: commit, Built: built,
Against: against(within, manifest)}, nil Against: against(within, manifest, stoodOn)}, nil
} }
// Log is where a build says what it is doing, step by step. Nil is silent — the tests pass none, // Log is where a build says what it is doing, step by step. Nil is silent — the tests pass none,
@@ -339,14 +342,27 @@ func describe(path string) string {
// allowed to name — a tag is something somebody else can move under you. // allowed to name — a tag is something somebody else can move under you.
var pinnedImage = regexp.MustCompile(`[A-Za-z0-9][A-Za-z0-9._/:-]*@sha256:[0-9a-f]{64}`) var pinnedImage = regexp.MustCompile(`[A-Za-z0-9][A-Za-z0-9._/:-]*@sha256:[0-9a-f]{64}`)
// against reads what this module's image artifacts are built on top of, out of the files that // against is what this module's image artifacts are built on top of: every base the mesh resolved
// build them. Nothing is guessed: a reference that is not written down is not reported. // and handed the recipe as a build argument (`build.on`), and any image a recipe pins by digest
func against(within string, manifest catalogue.Manifest) []string { // itself. Nothing is guessed: a reference that was neither resolved nor written down is not
// reported.
//
// **The resolved bases are the edges.** A recipe reads its base from an argument (`FROM
// ${RUNTIME_BASE}`), so the digest is never in the file, and a derivation that read files alone
// recorded no edge for any module on the mesh — which is why nothing knew what a changed base
// meant to rebuild (novox/hq 04-ISSUES/131).
func against(within string, manifest catalogue.Manifest, resolved []string) []string {
if manifest.Build == nil { if manifest.Build == nil {
return nil return nil
} }
seen := map[string]bool{} seen := map[string]bool{}
var out []string var out []string
for _, r := range resolved {
if r != "" && !seen[r] {
seen[r] = true
out = append(out, r)
}
}
for _, a := range manifest.Build.Artifacts { for _, a := range manifest.Build.Artifacts {
if a.Kind != catalogue.ArtifactImage || a.From == "" { if a.Kind != catalogue.ArtifactImage || a.From == "" {
continue continue
@@ -704,40 +720,42 @@ var _ io.Writer = (*stringWriter)(nil)
// built cannot be built here yet, and the useful sentence names which module is missing — not the // built cannot be built here yet, and the useful sentence names which module is missing — not the
// one a container runtime produces when a recipe's first line refers to an image nobody has. // one a container runtime produces when a recipe's first line refers to an image nobody has.
// //
// The order is fixed so two builds of one commit invoke the same command. // The order is fixed so two builds of one commit invoke the same command. Returned alongside the
// arguments is every reference they resolved to, which is what the build stood on.
func standingOn(ctx context.Context, manifest catalogue.Manifest, held map[string]string, func standingOn(ctx context.Context, manifest catalogue.Manifest, held map[string]string,
mirror func(ctx context.Context, from, repository string) (string, error)) ([]string, error) { mirror func(ctx context.Context, from, repository string) (string, error)) ([]string, []string, error) {
if manifest.Build == nil || len(manifest.Build.On) == 0 { if manifest.Build == nil || len(manifest.Build.On) == 0 {
return nil, nil return nil, nil, nil
} }
on := append([]catalogue.BuildsOn{}, manifest.Build.On...) on := append([]catalogue.BuildsOn{}, manifest.Build.On...)
sort.Slice(on, func(i, j int) bool { return on[i].Arg < on[j].Arg }) sort.Slice(on, func(i, j int) bool { return on[i].Arg < on[j].Arg })
var args []string var args, resolved []string
for _, base := range on { for _, base := range on {
if base.Image != "" { if base.Image != "" {
// A vendor's image, declared (novox/hq 04-ISSUES/064, ADR 0097). Pinned, because a tag // A vendor's image, declared (novox/hq 04-ISSUES/064, ADR 0097). Pinned, because a tag
// is what somebody else can move; copied into the mesh's registry, because a build // is what somebody else can move; copied into the mesh's registry, because a build
// that reaches a public registry on its own is a build that works sometimes. // that reaches a public registry on its own is a build that works sometimes.
if base.Arg == "" || base.Module != "" || base.Artifact != "" { if base.Arg == "" || base.Module != "" || base.Artifact != "" {
return nil, fmt.Errorf( return nil, nil, fmt.Errorf(
"%s stands on the image %s, and a base is either a module's artifact or an "+ "%s stands on the image %s, and a base is either a module's artifact or an "+
"image — never both — read from one build argument", manifest.Module, base.Image) "image — never both — read from one build argument", manifest.Module, base.Image)
} }
if !strings.Contains(base.Image, "@sha256:") { if !strings.Contains(base.Image, "@sha256:") {
return nil, fmt.Errorf( return nil, nil, fmt.Errorf(
"%s stands on the image %q, which is not pinned by digest. A tag is what "+ "%s stands on the image %q, which is not pinned by digest. A tag is what "+
"somebody else can move; name it as <image>@sha256:…", manifest.Module, base.Image) "somebody else can move; name it as <image>@sha256:…", manifest.Module, base.Image)
} }
reference, err := mirror(ctx, base.Image, manifest.Module+"/on-"+strings.ToLower(base.Arg)) reference, err := mirror(ctx, base.Image, manifest.Module+"/on-"+strings.ToLower(base.Arg))
if err != nil { if err != nil {
return nil, fmt.Errorf("%s stands on %s: %w", manifest.Module, base.Image, err) return nil, nil, fmt.Errorf("%s stands on %s: %w", manifest.Module, base.Image, err)
} }
args = append(args, "--build-arg", base.Arg+"="+reference) args = append(args, "--build-arg", base.Arg+"="+reference)
resolved = append(resolved, reference)
continue continue
} }
if base.Arg == "" || base.Module == "" || base.Artifact == "" { if base.Arg == "" || base.Module == "" || base.Artifact == "" {
return nil, fmt.Errorf( return nil, nil, fmt.Errorf(
"%s says its build stands on something, and does not say all of what: a base "+ "%s says its build stands on something, and does not say all of what: a base "+
"needs the module, the artifact, and the build argument the recipe reads it "+ "needs the module, the artifact, and the build argument the recipe reads it "+
"from", manifest.Module) "from", manifest.Module)
@@ -745,14 +763,15 @@ func standingOn(ctx context.Context, manifest catalogue.Manifest, held map[strin
key := base.Module + "/" + base.Artifact key := base.Module + "/" + base.Artifact
reference, has := held[key] reference, has := held[key]
if !has { if !has {
return nil, fmt.Errorf( return nil, nil, fmt.Errorf(
"%s builds on %s, and this mesh has not built it. Build %s first — every module "+ "%s builds on %s, and this mesh has not built it. Build %s first — every module "+
"in this toolchain stands on it, so it is the thing to have before anything "+ "in this toolchain stands on it, so it is the thing to have before anything "+
"else", manifest.Module, key, base.Module) "else", manifest.Module, key, base.Module)
} }
args = append(args, "--build-arg", base.Arg+"="+reference) args = append(args, "--build-arg", base.Arg+"="+reference)
resolved = append(resolved, reference)
} }
return args, nil return args, resolved, nil
} }
// compile runs a module's own code through its toolchain, and says where the result is. // compile runs a module's own code through its toolchain, and says where the result is.
+14
View File
@@ -180,6 +180,20 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
if err != nil { if err != nil {
return "", err return "", err
} }
// **Already held is already mirrored.** A base is named by digest, and a digest this registry
// holds under the module's repository is the same bytes whatever upstream would say — so
// upstream is not asked. Asked every build, the public hub's anonymous pull limit was reached
// on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base
// lives there failed on a copy it did not need.
if strings.HasPrefix(where.reference, "sha256:") {
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference)
if err != nil {
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
}
if held {
return r.Address + "/" + repository + "@" + where.reference, nil
}
}
src := &source{client: r.client()} src := &source{client: r.client()}
digest, err := r.copyManifest(ctx, src, where, where.reference, repository) digest, err := r.copyManifest(ctx, src, where, where.reference, repository)
if err != nil { if err != nil {
+30
View File
@@ -103,6 +103,12 @@ func (m *theMeshsRegistry) handler() http.Handler {
m.mu.Lock() m.mu.Lock()
defer m.mu.Unlock() defer m.mu.Unlock()
switch { switch {
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"):
if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK)
} else {
w.WriteHeader(http.StatusNotFound)
}
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"): case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"):
if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok { if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
@@ -219,3 +225,27 @@ func TestATagBeforeTheDigestIsNotPartOfTheRepository(t *testing.T) {
t.Fatalf("got %+v", got) t.Fatalf("got %+v", got)
} }
} }
// A base this registry already holds by digest is not asked of upstream at all: the public hub
// limits anonymous pulls, and a catalogue rebuilt on one merge asked it once per module.
func TestABaseAlreadyHeldIsNotAskedOfUpstream(t *testing.T) {
src, indexDigest, _ := anUpstreamRegistry(t)
dst := &theMeshsRegistry{blobs: map[string][]byte{}, manifests: map[string][]byte{}}
dstServer := httptest.NewServer(dst.handler())
defer dstServer.Close()
address := strings.TrimPrefix(dstServer.URL, "http://")
r := Registry{Address: address, HTTP: src.Client()}
host := strings.TrimPrefix(src.URL, "http://")
if _, err := r.MirrorImage(context.Background(), host+"/library/thing:latest", "hello-web/server"); err != nil {
t.Fatal(err)
}
// Upstream gone: the pinned base is answered from what the mesh holds.
src.Close()
reference, err := r.MirrorImage(context.Background(), host+"/library/thing@"+indexDigest, "hello-web/server")
if err != nil {
t.Fatalf("a base the registry holds was asked of an upstream that is gone: %v", err)
}
if reference != address+"/hello-web/server@"+indexDigest {
t.Fatalf("pinned as %q", reference)
}
}
+34 -6
View File
@@ -22,7 +22,7 @@ func TestABaseTheMeshHasNotBuiltIsRefused(t *testing.T) {
On: []catalogue.BuildsOn{{Arg: "RUNTIME_BASE", Module: "mesh-tools", Artifact: "runtime"}}, On: []catalogue.BuildsOn{{Arg: "RUNTIME_BASE", Module: "mesh-tools", Artifact: "runtime"}},
}, },
} }
_, err := standingOn(context.Background(), manifest, map[string]string{}, noMirror) _, _, err := standingOn(context.Background(), manifest, map[string]string{}, noMirror)
if err == nil { if err == nil {
t.Fatal("a base nothing has built was accepted; the build would have failed on its first line") t.Fatal("a base nothing has built was accepted; the build would have failed on its first line")
} }
@@ -42,7 +42,7 @@ func TestABaseTheMeshHoldsBecomesABuildArgument(t *testing.T) {
}, },
} }
held := map[string]string{"mesh-tools/runtime": "127.0.0.1:5000/mesh-tools/runtime@sha256:" + strings.Repeat("a", 64)} held := map[string]string{"mesh-tools/runtime": "127.0.0.1:5000/mesh-tools/runtime@sha256:" + strings.Repeat("a", 64)}
args, err := standingOn(context.Background(), manifest, held, noMirror) args, _, err := standingOn(context.Background(), manifest, held, noMirror)
if err != nil { if err != nil {
t.Fatalf("a base this mesh holds was refused: %v", err) t.Fatalf("a base this mesh holds was refused: %v", err)
} }
@@ -54,7 +54,7 @@ func TestABaseTheMeshHoldsBecomesABuildArgument(t *testing.T) {
// A module naming no base asks for nothing, which is most modules. // A module naming no base asks for nothing, which is most modules.
func TestAModuleNamingNoBaseAddsNoArguments(t *testing.T) { func TestAModuleNamingNoBaseAddsNoArguments(t *testing.T) {
args, err := standingOn(context.Background(), catalogue.Manifest{Module: "hello-web", Build: &catalogue.Build{}}, nil, noMirror) args, _, err := standingOn(context.Background(), catalogue.Manifest{Module: "hello-web", Build: &catalogue.Build{}}, nil, noMirror)
if err != nil || args != nil { if err != nil || args != nil {
t.Fatalf("a module naming no base produced %v, %v", args, err) t.Fatalf("a module naming no base produced %v, %v", args, err)
} }
@@ -66,7 +66,7 @@ func TestAnIncompleteBaseIsRefused(t *testing.T) {
Module: "postgres", Module: "postgres",
Build: &catalogue.Build{On: []catalogue.BuildsOn{{Module: "mesh-tools", Artifact: "runtime"}}}, Build: &catalogue.Build{On: []catalogue.BuildsOn{{Module: "mesh-tools", Artifact: "runtime"}}},
} }
if _, err := standingOn(context.Background(), manifest, map[string]string{"mesh-tools/runtime": "x"}, noMirror); err == nil { if _, _, err := standingOn(context.Background(), manifest, map[string]string{"mesh-tools/runtime": "x"}, noMirror); err == nil {
t.Fatal("a base with no build argument was accepted; nothing would have read it") t.Fatal("a base with no build argument was accepted; nothing would have read it")
} }
} }
@@ -86,7 +86,7 @@ func TestADeclaredVendorImageIsCopiedInAndHandedToTheRecipe(t *testing.T) {
}, },
} }
var asked []string var asked []string
args, err := standingOn(context.Background(), manifest, nil, func(_ context.Context, from, repository string) (string, error) { args, _, err := standingOn(context.Background(), manifest, nil, func(_ context.Context, from, repository string) (string, error) {
asked = append(asked, from+" -> "+repository) asked = append(asked, from+" -> "+repository)
return "127.0.0.1:5000/" + repository + "@sha256:" + strings.Repeat("d", 64), nil return "127.0.0.1:5000/" + repository + "@sha256:" + strings.Repeat("d", 64), nil
}) })
@@ -101,7 +101,7 @@ func TestADeclaredVendorImageIsCopiedInAndHandedToTheRecipe(t *testing.T) {
} }
// Unpinned, it is refused: a tag is what somebody else can move. // Unpinned, it is refused: a tag is what somebody else can move.
manifest.Build.On[0].Image = "quay.io/minio/mc:latest" manifest.Build.On[0].Image = "quay.io/minio/mc:latest"
if _, err := standingOn(context.Background(), manifest, nil, noMirror); err == nil || !strings.Contains(err.Error(), "not pinned") { if _, _, err := standingOn(context.Background(), manifest, nil, noMirror); err == nil || !strings.Contains(err.Error(), "not pinned") {
t.Fatalf("an unpinned vendor image was accepted: %v", err) t.Fatalf("an unpinned vendor image was accepted: %v", err)
} }
} }
@@ -151,3 +151,31 @@ func TestARecipeIsReadAsInstructions(t *testing.T) {
t.Fatalf("a heredoc line or a continued stage was read as a base: %v", bases) t.Fatalf("a heredoc line or a continued stage was read as a base: %v", bases)
} }
} }
// What a build was handed as its bases is what it stood on — recorded, so a changed base knows what
// to rebuild (novox/hq 04-ISSUES/131). A recipe reads the base from an argument, so nothing else
// could know.
func TestTheBasesABuildWasHandedAreWhatItStoodOn(t *testing.T) {
manifest := catalogue.Manifest{
Module: "gitea",
Build: &catalogue.Build{
On: []catalogue.BuildsOn{
{Arg: "RUNTIME_BASE", Module: "mesh-tools", Artifact: "runtime"},
{Arg: "BUILD_BASE", Module: "mesh-tools", Artifact: "build"},
},
Artifacts: []catalogue.Artifact{{Name: "runtime", Kind: catalogue.ArtifactImage, From: "Dockerfile"}},
},
}
held := map[string]string{
"mesh-tools/runtime": "127.0.0.1:5000/mesh-tools/runtime@sha256:" + strings.Repeat("a", 64),
"mesh-tools/build": "127.0.0.1:5000/mesh-tools/build@sha256:" + strings.Repeat("b", 64),
}
_, resolved, err := standingOn(context.Background(), manifest, held, noMirror)
if err != nil {
t.Fatal(err)
}
got := against(t.TempDir(), manifest, resolved)
if len(got) != 2 || got[0] != held["mesh-tools/build"] || got[1] != held["mesh-tools/runtime"] {
t.Fatalf("the bases the build was handed were not what it stood on: %v", got)
}
}
+15
View File
@@ -143,3 +143,18 @@ func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) {
} }
} }
} }
// 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)
}
}
+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
+35
View File
@@ -164,6 +164,41 @@ func (i *Inventory) Held(ctx context.Context) (map[string]string, error) {
return held, rows.Err() return held, rows.Err()
} }
// BuiltAgainst is what each module's newest successful build stood on, as recorded — the build
// edges (ADR 0009). A module whose last build recorded no bases is absent, which is also what a
// module standing on nothing looks like: an edge the mesh has not derived is not an edge.
func (i *Inventory) BuiltAgainst(ctx context.Context) (map[string][]string, error) {
rows, err := i.store.Pool().Query(ctx,
`select distinct on (module) module, built_against
from build
where module is not null and module <> '' and failed = ''
order by module, at desc`)
if err != nil {
return nil, err
}
defer rows.Close()
against := map[string][]string{}
for rows.Next() {
var module string
var raw []byte
if err := rows.Scan(&module, &raw); err != nil {
return nil, err
}
if len(raw) == 0 {
continue
}
var refs []string
if err := json.Unmarshal(raw, &refs); err != nil {
continue
}
if len(refs) > 0 {
against[module] = refs
}
}
return against, rows.Err()
}
// manifestOrNil keeps the difference between "declared nothing" and "predates this being kept". // manifestOrNil keeps the difference between "declared nothing" and "predates this being kept".
// //
// A build recorded before the mesh kept manifests has no manifest, and that is not the same as one // A build recorded before the mesh kept manifests has no manifest, and that is not the same as one
+3 -1
View File
@@ -86,9 +86,11 @@ func TestAnAssignedModuleBecomesAUserWithWhatItDeclared(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
// Its tools are every one under its own name — the list in the manifest is a person's
// vocabulary for asking, not the module's permission to answer.
if !granted(perms.Publish, "mesh.mod.shop.event.order.placed") || if !granted(perms.Publish, "mesh.mod.shop.event.order.placed") ||
!granted(perms.Publish, "mesh.seat.telegram-sender.accept.send") || !granted(perms.Publish, "mesh.seat.telegram-sender.accept.send") ||
!granted(perms.Subscribe, "mesh.mod.shop.tool.price") { !granted(perms.Subscribe, "mesh.mod.shop.tool.>") {
t.Fatalf("one.shop's authority is not what it declared: %+v", perms) t.Fatalf("one.shop's authority is not what it declared: %+v", perms)
} }
} }
+10 -2
View File
@@ -8,6 +8,7 @@ import (
"sort" "sort"
"strconv" "strconv"
"strings" "strings"
"time"
"github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5"
"github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/catalogue"
@@ -38,6 +39,9 @@ type Source struct {
BuiltFrom string BuiltFrom string
// Head is the newest commit the source is known to have. // Head is the newest commit the source is known to have.
Head string Head string
// Seen is when the source was last looked at — by a build, by hand, or by the forge saying it
// moved. What a late report of an older move is judged against.
Seen time.Time
} }
// Current reports whether what the mesh holds is what the source last had. // Current reports whether what the mesh holds is what the source last had.
@@ -936,12 +940,13 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
`select m.name, m.manifest, `select m.name, m.manifest,
coalesce(m.source, ''), m.source_path, m.source_seat, coalesce(m.ref, ''), coalesce(m.source, ''), m.source_path, m.source_seat, coalesce(m.ref, ''),
coalesce(m.built_from, ''), coalesce(m.source_head, ''), coalesce(m.built_from, ''), coalesce(m.source_head, ''),
coalesce(m.source_seen, to_timestamp(0)),
coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}') coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}')
from module m from module m
left join assignment a on a.module = m.name left join assignment a on a.module = m.name
left join node n on n.id = a.node left join node n on n.id = a.node
group by m.name, m.manifest, m.source, m.source_path, m.source_seat, m.ref, m.built_from, group by m.name, m.manifest, m.source, m.source_path, m.source_seat, m.ref, m.built_from,
m.source_head m.source_head, m.source_seen
order by m.name`) order by m.name`)
if err != nil { if err != nil {
return nil, err return nil, err
@@ -955,9 +960,12 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
var source Source var source Source
var on []string var on []string
if err := rows.Scan(&name, &raw, &source.Repository, &source.Path, &source.Seat, &source.Ref, if err := rows.Scan(&name, &raw, &source.Repository, &source.Path, &source.Seat, &source.Ref,
&source.BuiltFrom, &source.Head, &on); err != nil { &source.BuiltFrom, &source.Head, &source.Seen, &on); err != nil {
return nil, err return nil, err
} }
if source.Seen.Unix() == 0 {
source.Seen = time.Time{}
}
var m catalogue.Manifest var m catalogue.Manifest
if err := json.Unmarshal(raw, &m); err != nil { if err := json.Unmarshal(raw, &m); err != nil {
return nil, err return nil, err
+22 -60
View File
@@ -3,16 +3,13 @@ package link
import ( import (
"context" "context"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"time" "time"
amqp "github.com/rabbitmq/amqp091-go" "github.com/nats-io/nats.go"
) )
// RPCExchange is where a module's tools are asked over the broker, keyed `<module>.<tool>`, and
// where the answer comes back, keyed by the asker's reply queue (novox/hq ADR 0047).
const RPCExchange = "mesh.rpc"
// Answer is what a module's tool replies: one of the two, never both. // Answer is what a module's tool replies: one of the two, never both.
type Answer struct { type Answer struct {
Result json.RawMessage `json:"result,omitempty"` Result json.RawMessage `json:"result,omitempty"`
@@ -27,70 +24,35 @@ type Answer struct {
// grants, and should not. The control plane already holds a connection that may, so a person or // grants, and should not. The control plane already holds a connection that may, so a person or
// an agent asks through it, and every question passes one process where an audit belongs. // an agent asks through it, and every question passes one process where an audit belongs.
// //
// The reply queue is the caller's own, server-named and exclusive, bound to the RPC exchange under // The answer comes back on the asker's own inbox, which only the asker may read; the serving
// its own name: a serving module answers through that exchange and never the default one, whose // module answers there and nowhere else. The bus refuses a request nothing serves at once, so a
// permission is per exchange rather than per queue. The correlation is checked rather than // module that is down or a tool that does not exist is said now rather than after the whole wait.
// assumed, as every RPC here is. func Ask(ctx context.Context, bus Bus, module, tool string, args json.RawMessage,
func Ask(ctx context.Context, channel *amqp.Channel, module, tool string, args json.RawMessage,
timeout time.Duration) (Answer, error) { timeout time.Duration) (Answer, error) {
if len(args) == 0 { if len(args) == 0 {
args = json.RawMessage(`{}`) args = json.RawMessage(`{}`)
} }
replies, err := channel.QueueDeclare("", false, true, true, false, nil)
if err != nil {
return Answer{}, err
}
if err := channel.QueueBind(replies.Name, replies.Name, RPCExchange, false, nil); err != nil {
return Answer{}, fmt.Errorf("cannot bind a reply queue to %s: %w", RPCExchange, err)
}
answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
if err != nil {
return Answer{}, err
}
// Mandatory, so a request nothing consumes comes straight back: a module that is down, or a
// tool that does not exist, is said at once rather than after the whole wait.
returned := channel.NotifyReturn(make(chan amqp.Return, 1))
id := fmt.Sprintf("ask-%d", time.Now().UnixNano())
key := module + "." + tool key := module + "." + tool
if err := channel.PublishWithContext(ctx, RPCExchange, key, true, false, amqp.Publishing{ reply, err := bus.AskTool(ctx, module, tool, args, timeout)
ContentType: "application/json", if err != nil {
CorrelationId: id, if errors.Is(err, nats.ErrNoResponders) {
ReplyTo: replies.Name, return Answer{}, fmt.Errorf(
Body: args, "nothing serves %s: no runtime has bound %q on the bus. The module is not "+
}); err != nil { "assigned, its runtime is not up, or it serves no such tool — `status` says "+
return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err) "whether the machine carrying it has applied", module, key)
} }
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) {
waiting, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
for {
select {
case back := <-returned:
if back.CorrelationId == id {
return Answer{}, fmt.Errorf(
"nothing serves %s: no runtime has bound %q on the broker. The module is not "+
"assigned, its runtime is not up, or it serves no such tool — `status` "+
"says whether the machine carrying it has applied", module, key)
}
case <-waiting.Done():
return Answer{}, fmt.Errorf( return Answer{}, fmt.Errorf(
"%s did not answer within %s. Its runtime serves %q when it is up and has bound "+ "%s did not answer within %s. Its runtime serves %q when it is up and has bound "+
"the broker — `status` says whether the machine carrying it has applied", "the bus — `status` says whether the machine carrying it has applied",
module, timeout, key) module, timeout, key)
case delivery, ok := <-answers:
if !ok {
return Answer{}, fmt.Errorf("the connection closed while waiting for %s", key)
}
if delivery.CorrelationId != id {
continue
}
var answer Answer
if err := json.Unmarshal(delivery.Body, &answer); err != nil {
return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err)
}
return answer, nil
} }
return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err)
} }
var answer Answer
if err := json.Unmarshal(reply, &answer); err != nil {
return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err)
}
return answer, nil
} }
-94
View File
@@ -1,12 +1,7 @@
package link package link
import ( import (
"context"
"encoding/json" "encoding/json"
"fmt"
"time"
amqp "github.com/rabbitmq/amqp091-go"
) )
// Asking a machine to build a module, and hearing what came out. // Asking a machine to build a module, and hearing what came out.
@@ -22,21 +17,6 @@ import (
// holds no opinion about what they contain, and a host that also built things would be a host // holds no opinion about what they contain, and a host that also built things would be a host
// with a container runtime requirement and a git dependency (novox/hq ADR 0005). // with a container runtime requirement and a git dependency (novox/hq ADR 0005).
// BuildQueue is where build requests wait. One queue, so several build machines can share the
// work and each request is done exactly once — which is what a queue is for and what a
// per-machine routing key would not give.
const BuildQueue = "builds"
// KeyBuilt is what a builder publishes when it has finished, successfully or not.
const KeyBuilt = "built"
// ReplyQueue is where the answer to one request goes.
//
// **Named here rather than left to the broker**, so it can be scoped and reasoned about. A broker
// generates its own name for an unnamed queue, and a builder permitted to write to whatever that
// convention happens to produce works on one broker and silently cannot answer on another.
func ReplyQueue(id string) string { return BuildQueue + ".reply." + id }
// BuildRequest is one module to build. // BuildRequest is one module to build.
type BuildRequest struct { type BuildRequest struct {
// ID correlates the answer with the asking. Not the module name: two builds of one module can // ID correlates the answer with the asking. Not the module name: two builds of one module can
@@ -118,77 +98,3 @@ type MadeArtifact struct {
Kind string `json:"kind"` Kind string `json:"kind"`
Reference string `json:"reference"` Reference string `json:"reference"`
} }
// RequestBuild asks for a module to be built and waits for the answer.
//
// Waiting rather than returning immediately, because the thing a person wants after asking for a
// build is to know whether it worked. A build that is dispatched and forgotten needs somewhere to
// look afterwards, and there is nowhere yet.
func RequestBuild(ctx context.Context, channel *amqp.Channel, request BuildRequest,
timeout time.Duration) (BuildResult, error) {
// Its own queue for the answer, declared before the ask. Consuming from the shared exchange
// would mean competing with the control plane's own consumer for a message meant for this
// caller — which is the fault this package's own doc comment records having had.
replies, err := channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil)
if err != nil {
return BuildResult{}, err
}
// Bound to the exchange, and the answer comes back through it.
//
// **A builder never publishes to the default exchange**, because permission there is per
// exchange and not per queue — a builder allowed to use it could publish into any node's
// queue, which is the privilege a build machine most obviously should not have. Found by
// running it: the builder built, could not answer, and the connection closed saying only
// "not allowed to publish to exchange ''".
//
// The cost is that every asker sees every result, which is why the correlation is checked
// below rather than assumed.
if err := channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil {
return BuildResult{}, err
}
answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
if err != nil {
return BuildResult{}, err
}
body, err := json.Marshal(request)
if err != nil {
return BuildResult{}, err
}
if err := channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
CorrelationId: request.ID,
ReplyTo: replies.Name,
Body: body,
}); err != nil {
return BuildResult{}, err
}
waiting, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
for {
select {
case <-waiting.Done():
return BuildResult{}, fmt.Errorf(
"no builder answered within %s. Something must be consuming %q, and nothing is "+
"— or it is building something that takes longer than this",
timeout, BuildQueue)
case delivery, ok := <-answers:
if !ok {
return BuildResult{}, fmt.Errorf("the connection closed while waiting for a build")
}
var result BuildResult
if err := json.Unmarshal(delivery.Body, &result); err != nil {
return BuildResult{}, fmt.Errorf("a builder answered with something unreadable: %w", err)
}
if result.ID != request.ID {
// Somebody else's answer on this queue. Ignored rather than returned, because
// returning it would attribute one build's outcome to another's.
continue
}
return result, nil
}
}
}
-202
View File
@@ -1,202 +0,0 @@
package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
// The build flow on the bus the mesh runs on today.
//
// Moved behind the seam rather than changed. The queue, the reply binding and the correlation are
// what they were, because the mesh is running on this.
// currentBuilds asks for builds over a channel.
type currentBuilds struct{ channel *amqp.Channel }
// BuildsOverCurrent is the asking side on the bus the mesh has.
func BuildsOverCurrent(channel *amqp.Channel) Builders { return currentBuilds{channel: channel} }
func (b currentBuilds) Close() {}
func (b currentBuilds) Submit(ctx context.Context, request BuildRequest,
wait time.Duration) (BuildResult, error) {
// Its own queue for the answer, declared before the ask. Consuming from the shared exchange
// would mean competing with the controller's own consumer for a message meant for this caller.
replies, err := b.channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil)
if err != nil {
return BuildResult{}, err
}
// **A builder never publishes to the default exchange**, because permission there is per
// exchange and not per queue — a builder allowed to use it could publish into any node's queue,
// which is the privilege a build machine most obviously should not have. The cost is that every
// asker sees every result, which is why the correlation is checked below rather than assumed.
if err := b.channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil {
return BuildResult{}, err
}
answers, err := b.channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
if err != nil {
return BuildResult{}, err
}
body, err := json.Marshal(request)
if err != nil {
return BuildResult{}, err
}
if err := b.channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
CorrelationId: request.ID,
ReplyTo: replies.Name,
Body: body,
}); err != nil {
return BuildResult{}, err
}
waiting, cancel := context.WithTimeout(ctx, wait)
defer cancel()
for {
select {
case <-waiting.Done():
return BuildResult{}, waitingFor(wait)
case delivery, ok := <-answers:
if !ok {
return BuildResult{}, errors.New("the connection closed while waiting for a build")
}
result, mine, err := theOutcomeOf(delivery.Body, request.ID)
if err != nil {
return BuildResult{}, err
}
if mine {
return result, nil
}
}
}
}
// --- the machine's side ---------------------------------------------------------------------
type currentMachine struct {
conn *amqp.Connection
channel *amqp.Channel
on string
}
// MachineOverCurrent takes build work over a channel.
func MachineOverCurrent(conn *amqp.Connection, channel *amqp.Channel, on string) BuildMachine {
return &currentMachine{conn: conn, channel: channel, on: on}
}
func (m *currentMachine) Close() {}
func (m *currentMachine) Take(ctx context.Context, do func(context.Context, Build)) error {
if _, err := m.channel.QueueDeclare(BuildQueue, true, false, false, false, nil); err != nil {
return err
}
// One at a time. A machine that took five requests at once would run five container builds
// against one runtime and finish all of them slower than it would have finished the first — and
// the queue is what shares work between machines, so nothing is lost by it.
if err := m.channel.Qos(1, 0, false); err != nil {
return err
}
// Not auto-acknowledged: a request acknowledged on arrival is a build that vanishes if this
// process dies mid-way, with nobody waiting on it ever hearing why.
requests, err := m.channel.ConsumeWithContext(ctx, BuildQueue, "mesh-builder",
false, false, false, false, nil)
if err != nil {
return err
}
for {
select {
case <-ctx.Done():
return nil
case delivery, ok := <-requests:
if !ok {
return errors.New("the broker closed the connection")
}
var request BuildRequest
if err := json.Unmarshal(delivery.Body, &request); err != nil {
// Unreadable: rejected rather than retried, because the next attempt reads the same
// bytes. Nobody waiting hears an answer, which is correct — there was no request.
_ = delivery.Reject(false)
continue
}
do(ctx, &currentBuild{request: request, delivery: delivery, on: m.on, channel: m.channel})
}
}
}
type currentBuild struct {
request BuildRequest
delivery amqp.Delivery
on string
channel *amqp.Channel
}
func (b *currentBuild) Request() BuildRequest { return b.request }
// Announce answers and announces, which on this bus are two publishes to two exchanges.
//
// The reply goes to whoever asked, correlated to their request; the announcement says to the whole
// mesh that a module now exists at a commit (novox/hq ADR 0072). Only a successful build is
// announced: a failed one produced no module version, and announcing one would put something in the
// graph that was never made.
func (b *currentBuild) Announce(ctx context.Context, result BuildResult) error {
body, err := json.Marshal(result)
if err != nil {
return err
}
if err := b.channel.PublishWithContext(ctx, Exchange, KeyBuilt, false, false, amqp.Publishing{
ContentType: "application/json",
CorrelationId: result.ID,
Body: body,
}); err != nil {
return fmt.Errorf("cannot answer a build request: %w", err)
}
if result.Failed != "" || result.Commit == "" {
return nil
}
// **Announced under both names on this bus, for exactly as long as this bus lives.**
//
// A build's outcome belongs to the role now (novox/hq ADR 0121), so a catalogue built from the
// current manifests listens for the role's name. A catalogue that is *already running* listens for
// the module's, because that is what it was told when it was installed. A rename on a live bus
// needs the publisher and the subscriber to change together, and a merge cannot promise that: one
// of them is deployed first, and in that window the graph silently stops being updated — which is
// the failure this whole change was cleaning up after.
//
// So both, and the order stops mattering. The module's own name goes with the bus, in step 5's
// retirement list; nothing has ever run on the bus being built, so there is no legacy name there
// and this doubling has no counterpart.
announced := announcementOf(result)
if err := EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyModuleBuilt, "builder", b.on,
announced); err != nil {
return err
}
return EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyRoleBuilt, TheBuildMachine, b.on,
announced)
}
func (b *currentBuild) Done() error { return b.delivery.Ack(false) }
func (b *currentBuild) Hold(time.Duration) error {
// No delayed redelivery on this bus: handed back at once, which is what it has always done.
return b.delivery.Nack(false, true)
}
// announcementOf is what the mesh is told about a finished build. One function, so the two
// transports cannot describe the same build differently.
func announcementOf(result BuildResult) map[string]any {
return map[string]any{
"module": ModuleOf(result.Manifest), "commit": result.Commit,
"repository": result.Repository, "path": result.Path, "ref": result.Ref,
"manifest": json.RawMessage(result.Manifest), "against": result.Against,
"made": result.Made,
}
}
-46
View File
@@ -7,7 +7,6 @@ import (
"time" "time"
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
amqp "github.com/rabbitmq/amqp091-go"
) )
// Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a // Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a
@@ -39,51 +38,6 @@ type Bus interface {
// --- The bus the mesh runs on today ----------------------------------------------------- // --- The bus the mesh runs on today -----------------------------------------------------
// OverCurrent is the bus the mesh runs on today, until the rollout.
type OverCurrent struct{ Channel *amqp.Channel }
func (b OverCurrent) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
id, err := eventID()
if err != nil {
return err
}
return b.Channel.PublishWithContext(ctx, EventsExchange, key, false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
MessageId: id,
Timestamp: time.Now().UTC(),
Body: body,
Headers: amqp.Table{
"x-event-id": id,
"x-source": source,
"x-node": node,
"x-time": time.Now().UTC().Format(time.RFC3339),
"content-type": "application/json",
},
})
}
// AskTool is implemented over the existing reply-queue machinery in ask.go; this seam does not
// change how it works today.
func (b OverCurrent) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
answer, err := Ask(ctx, b.Channel, module, tool, args, timeout)
if err != nil {
return nil, err
}
return answer.Result, nil
}
func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error {
// To the queue directly rather than through an exchange: a declaration is for one node, and
// routing it by name through a shared exchange would mean a binding per node that nothing
// removes when a node is retired.
return b.Channel.PublishWithContext(ctx, "", QueueFor(node), false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
Body: body,
})
}
// --- NATS, the bus being built ---------------------------------------------------------------- // --- NATS, the bus being built ----------------------------------------------------------------
// OverNATS is the bus as a JetStream context. // OverNATS is the bus as a JetStream context.
+23
View File
@@ -0,0 +1,23 @@
package link
import (
"github.com/novox/mesh-controller/internal/broker"
)
// 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.Context(), Conn: js.Conn()},
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 }
+3 -18
View File
@@ -24,10 +24,9 @@ import (
// — it holds both grants and asks each for its part, which is what the process running them is // — it holds both grants and asks each for its part, which is what the process running them is
// for. // for.
type Enrolment struct { type Enrolment struct {
Inventory *inventory.Inventory Inventory *inventory.Inventory
Identity *identity.Identity Identity *identity.Identity
Management *broker.Management Broker broker.Broker
Broker broker.Broker
// OnNATS says the mesh's own traffic is on the bus being built, so a node's credential is // OnNATS says the mesh's own traffic is on the bus being built, so a node's credential is
// minted into the mesh's records and composed into the bus's user list rather than pushed // minted into the mesh's records and composed into the bus's user list rather than pushed
@@ -194,17 +193,6 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol
reply.Password = password reply.Password = password
} }
case e.Management != nil:
password, err := freshPassword()
if err != nil {
return EnrolReply{}, err
}
if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil {
log.Printf("%s is enrolled and its broker password could not be replaced, so it keeps "+
"the token's secret as its password: %v", node.Name, err)
} else {
reply.Password = password
}
} }
if profile != nil { if profile != nil {
@@ -261,9 +249,6 @@ func freshPassword() (string, error) {
var _ Enroller = Enrolment{} var _ Enroller = Enrolment{}
// ErrNoBrokerManagement is returned when an account cannot be made because nothing was configured.
var ErrNoBrokerManagement = errors.New("no broker management configured")
// Heard records what a node reported about itself. // Heard records what a node reported about itself.
// //
// A node states; the owning context writes (novox/hq ADR 0006). What a node says it applied is // A node states; the owning context writes (novox/hq ADR 0006). What a node says it applied is
+15
View File
@@ -117,6 +117,21 @@ 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"`
// MergedAt is when the forge merged it, RFC 3339. What decides whether this is news.
MergedAt string `json:"merged_at"`
}
type Upgraded struct { type Upgraded struct {
Module string `json:"module"` Module string `json:"module"`
Commit string `json:"commit"` Commit string `json:"commit"`
+3 -13
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"
) )
@@ -75,19 +78,6 @@ type Control interface {
// Took settles the message: acted on, or understood and needing no action. // Took settles the message: acted on, or understood and needing no action.
Took() error Took() error
// About names what this message is about — a node's report, one module's move, one build's
// outcome — and is said before the store is asked.
//
// A transport that holds messages **in memory** uses it to set aside anything older it is
// holding about the same thing: the older is the past, and letting it come back after the
// newer was acted on would undo the newer.
//
// **This is the one thing holding-in-memory can do that holding-in-the-server cannot**, and
// naming it here rather than hiding it is deliberate. On the bus being built the message
// belongs to the server and comes back whatever happened meanwhile, so this is ignored and the
// digest a report carries answers the same question instead (window.go, design 25 §3).
About(what string)
// Hold keeps the message and asks for it again after the delay — the store window. // Hold keeps the message and asks for it again after the delay — the store window.
Hold(after time.Duration) error Hold(after time.Duration) error
-317
View File
@@ -1,317 +0,0 @@
package link
import (
"context"
"errors"
"fmt"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
// The consume side on the bus the mesh runs on today.
//
// Everything here was the serving loop's until the seam went in: the queues, the binds, the
// prefetch, and the list of messages the store could not take yet. It moved rather than changed —
// the behaviour this transport has is the behaviour it had, because the mesh is running on it and
// a bus nothing speaks yet is no reason to alter the one every node is on (ADR 0116).
// Prefetch is how many messages the bus hands the controller before it has settled them.
//
// More than one because a message the store could not take is held, unsettled, while the loop
// goes on answering others — an enrolment above all, which a host is waiting on (novox/hq issue
// 083). Bounded, because what is held is also what the bus has not kept on its own disk as
// pending.
const Prefetch = 64
// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to
// answer — an enrolment above all — while others wait for the store.
const PrefetchHeadroom = 8
// TryAgainAfter is how often the held are looked at. A store comes back in seconds, and a report a
// few seconds late is still current.
const TryAgainAfter = 2 * time.Second
// currentInbound consumes what nodes say over the bus the mesh has.
type currentInbound struct {
conn *amqp.Connection
channel *amqp.Channel
// upgrades and catchups are bound only when something is listening (Also).
upgrades bool
catchups bool
// held is every message the store could not take, by the bus's own delivery tag. Kept here
// rather than in the serving loop because holding a delivery unacknowledged is this
// transport's way of keeping it, and the other's is to hand it back to the server.
held map[uint64]*holding
// again is how often the held are looked at; zero means TryAgainAfter. Set by tests.
again time.Duration
}
// holding is one message kept for the store, and when to try it again.
type holding struct {
message *currentControl
due time.Time
about string
}
// Current is the consume side of the bus the mesh runs on today.
func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound {
return &currentInbound{conn: conn, channel: channel, held: map[uint64]*holding{}}
}
// Also binds the queue one more kind arrives on.
//
// The kinds nodes publish all share one queue and are bound at Connect, because a node may
// publish any of them and binding one while forgetting another is a message the bus accepts, finds
// no queue for, and drops — the publisher sees success and the consumer sees nothing. The two that
// are events get their own queue each, and only when something is listening.
func (c *currentInbound) Also(kind string) error {
switch kind {
case KindModuleMoved:
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
return err
}
c.upgrades = true
case KindCatchUp:
if err := c.bindEvent(CatchUpQueue, KeyCatchingUp); err != nil {
return err
}
c.catchups = true
default:
return fmt.Errorf("nothing binds a queue for %s on this bus", kind)
}
return nil
}
func (c *currentInbound) bindEvent(queue, key string) error {
if _, err := c.channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
return fmt.Errorf("cannot declare the %s queue: %w", queue, err)
}
if err := c.channel.QueueBind(queue, key, EventsExchange, false, nil); err != nil {
return fmt.Errorf("cannot bind %s to %s/%s: %w", queue, EventsExchange, key, err)
}
return nil
}
func (c *currentInbound) Close() {}
// Receive consumes until the context ends.
//
// One consumer per queue, deliberately: with two on one queue the bus would round-robin between
// them and each would receive half of what it expects — a fault this project has already had,
// between a module's daemon and its capability server.
func (c *currentInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
// A bounded prefetch rather than one. The loop still takes messages one at a time; what the
// prefetch buys is that a message the store could not take can be held while the loop goes on
// to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything
// held goes back to the bus if the controller stops, because nothing held is acknowledged.
if err := c.channel.Qos(Prefetch, 0, false); err != nil {
return err
}
deliveries, err := c.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane",
false, false, false, false, nil)
if err != nil {
return err
}
// Its own queue and its own consumer for each event, for the reason above: two consumers on
// one queue split its messages between them, and an upgrade or a catch-up request going to
// whichever half was not listening is a gap that looks like a working mesh.
var upgrades, catchups <-chan amqp.Delivery
if c.upgrades {
upgrades, err = c.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades",
false, false, false, false, nil)
if err != nil {
return err
}
}
if c.catchups {
catchups, err = c.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup",
false, false, false, false, nil)
if err != nil {
return err
}
}
closed := c.conn.NotifyClose(make(chan *amqp.Error, 1))
again := c.again
if again == 0 {
again = TryAgainAfter
}
ticker := time.NewTicker(again)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case <-ticker.C:
if ctx.Err() != nil {
return nil
}
c.retryHeld(ctx, act)
case delivery, ok := <-catchups:
if !ok {
if catchups != nil {
return errors.New("the bus stopped delivering catch-up requests")
}
continue
}
act(ctx, c.wrap(KindCatchUp, delivery))
case delivery, ok := <-upgrades:
// A nil channel blocks for ever, so this case simply never fires when nothing is
// listening for upgrades. Closed is different, and means the bus stopped.
if !ok {
if upgrades != nil {
return errors.New("the bus stopped delivering upgrades")
}
continue
}
act(ctx, c.wrap(KindModuleMoved, delivery))
case reason := <-closed:
// Said rather than returned quietly. A controller whose bus connection dropped is a
// mesh where nothing can be told anything, and the reason is the first thing anybody
// will want.
return fmt.Errorf("the bus connection closed: %v", reason)
case delivery, ok := <-deliveries:
if !ok {
return errors.New("the bus stopped delivering")
}
kind, known := kindOfKey[delivery.RoutingKey]
if !known {
// Rejected without requeue: a message nothing understands will not be understood
// on the next attempt either, and requeuing it would spin.
_ = delivery.Reject(false)
continue
}
act(ctx, c.wrap(kind, delivery))
}
}
}
// kindOfKey is how this transport's addressing becomes what the mesh calls a message.
var kindOfKey = map[string]string{
KeyEnrol: KindEnrolment,
KeyReport: KindReport,
KeyAlive: KindHeartbeat,
KeyBuilt: KindBuilt,
KeyModuleUpgraded: KindModuleMoved,
KeyCatchingUp: KindCatchUp,
}
func (c *currentInbound) wrap(kind string, delivery amqp.Delivery) *currentControl {
return &currentControl{kind: kind, delivery: delivery, on: c}
}
// retryHeld hands every message whose delay has passed back to the loop. Each handler holds it
// again, settles it, or lets it go past the bound.
func (c *currentInbound) retryHeld(ctx context.Context, act func(context.Context, Control)) {
now := time.Now()
due := make([]*currentControl, 0, len(c.held))
for _, h := range c.held {
if !h.due.After(now) {
due = append(due, h.message)
}
}
for _, m := range due {
if ctx.Err() != nil {
return
}
act(ctx, m)
}
}
// currentControl is one delivery from the bus the mesh has, as the controller reads it.
type currentControl struct {
kind string
delivery amqp.Delivery
on *currentInbound
// about is what this message is about, as the handler named it; empty until it does.
about string
// first is when this message was first held for the store; zero while it has not been.
first time.Time
}
func (m *currentControl) Kind() string { return m.kind }
func (m *currentControl) Body() []byte { return m.delivery.Body }
func (m *currentControl) Redelivered() bool { return m.delivery.Redelivered }
func (m *currentControl) HeldFor() time.Duration {
if m.first.IsZero() {
return 0
}
return time.Since(m.first)
}
// Answer publishes to the reply queue the request named.
func (m *currentControl) Answer(ctx context.Context, body []byte) error {
if m.delivery.ReplyTo == "" {
return errors.New("that request named no reply queue, so nothing can be told the answer")
}
return m.on.channel.PublishWithContext(ctx, "", m.delivery.ReplyTo, false, false,
amqp.Publishing{
ContentType: "application/json",
CorrelationId: m.delivery.CorrelationId,
Body: body,
})
}
func (m *currentControl) Took() error {
m.forget()
return m.delivery.Ack(false)
}
// Drop rejects without requeue: on this bus that is what "understood, and not worth another
// attempt" is spelled as, and it is what feeds a dead-letter queue where one is configured.
func (m *currentControl) Drop() error {
m.forget()
return m.delivery.Reject(false)
}
// About names what this message is about, and lets go of whatever is held about the same thing:
// the held one is the past, and acting on it after this one would undo this one. Acknowledged
// rather than left to come back, because a held message nothing will act on is a place in the
// prefetch nothing gets back.
func (m *currentControl) About(what string) {
m.about = what
if what == "" {
return
}
for tag, h := range m.on.held {
if h.about != what || tag == m.delivery.DeliveryTag {
continue
}
delete(m.on.held, tag)
_ = h.message.delivery.Ack(false)
}
}
// Hold keeps the message unacknowledged and sets it aside to be handed back after the delay.
//
// Held no further than the prefetch leaves room: past that the bus would hand the loop nothing
// new — enrolments included — until something held was let go. A message that cannot be held says
// so, and the handler settles it its own way.
func (m *currentControl) Hold(after time.Duration) error {
if m.on.held == nil {
m.on.held = map[uint64]*holding{}
}
if _, already := m.on.held[m.delivery.DeliveryTag]; !already {
if len(m.on.held) >= Prefetch-PrefetchHeadroom {
return fmt.Errorf("%d messages are already held for the store, and holding more "+
"would stop the queue", len(m.on.held))
}
m.first = time.Now()
}
m.on.held[m.delivery.DeliveryTag] = &holding{
message: m, due: time.Now().Add(after), about: m.about,
}
return nil
}
func (m *currentControl) forget() {
if m.on != nil {
delete(m.on.held, m.delivery.DeliveryTag)
}
}
-71
View File
@@ -1,71 +0,0 @@
package link
import (
"context"
"encoding/json"
"io"
"log"
"testing"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
// The harness for the consume side on the bus the mesh runs on today.
//
// Messages arrive through the seam, so what these tests exercise is the controller's decision
// about a message and this transport's way of keeping one — which is what the seam separated. A
// fake acknowledger stands in for the bus, because what is asserted is how a message was settled
// and that needs no server.
// settled is how the bus was told to settle one message.
type settled struct{ acked, nacked, requeued, rejected bool }
func (a *settled) Ack(uint64, bool) error { a.acked = true; return nil }
func (a *settled) Nack(_ uint64, _ bool, requeue bool) error {
a.nacked, a.requeued = true, requeue
return nil
}
func (a *settled) Reject(uint64, bool) error { a.rejected = true; return nil }
// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will
// hand it to whatever consumes next if the controller stops.
func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected }
var tag uint64
func quiet() *log.Logger { return log.New(io.Discard, "", 0) }
// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder,
// an upgrader or a replayer to be set on it.
func serving() (*Server, *currentInbound) {
in := &currentInbound{held: map[uint64]*holding{}}
return &Server{inbound: in, bus: OverCurrent{}, log: quiet()}, in
}
// sends is one message arriving over this transport, as the controller reads it.
func (c *currentInbound) sends(t *testing.T, to *settled, kind string, v any) Control {
t.Helper()
body, err := json.Marshal(v)
if err != nil {
t.Fatal(err)
}
tag++
return &currentControl{kind: kind, on: c, delivery: amqp.Delivery{
Acknowledger: to, Body: body, DeliveryTag: tag,
}}
}
// dueNow brings every held message forward, so a test need not wait out the backoff a real store
// restart would be given (RedeliverAfter).
func (c *currentInbound) dueNow() {
for _, h := range c.held {
h.due = time.Now().Add(-time.Second)
}
}
// retries hands every held message back to the controller, the way the ticker does.
func (c *currentInbound) retries(ctx context.Context, s *Server) {
c.dueNow()
c.retryHeld(ctx, s.act)
}
+94
View File
@@ -0,0 +1,94 @@
package link
import (
"context"
"encoding/json"
"testing"
"time"
)
// A harness for the controller's decisions about a message, with no bus behind it.
//
// What these tests exercise is the server's verdict — taken, dropped, held for the store, let go
// once the store has been away too long — and the transport's part in that is only to keep a held
// message and hand it back. A fake that does exactly that stands in for the bus, so what is
// asserted is how a message was settled and the decision needs no server.
// settled is how the bus was told to settle one message.
type settled struct{ acked, nacked, requeued, rejected bool }
// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will
// hand it to whatever consumes next if the controller stops.
func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected }
// fakeInbound keeps the messages the controller held, the way a stream would.
type fakeInbound struct{ held map[uint64]*fakeControl }
var tag uint64
// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder,
// an upgrader or a replayer to be set on it.
func serving() (*Server, *fakeInbound) {
in := &fakeInbound{held: map[uint64]*fakeControl{}}
return &Server{inbound: in, log: quiet()}, in
}
func (c *fakeInbound) Also(string) error { return nil }
func (c *fakeInbound) Close() {}
func (c *fakeInbound) Receive(context.Context, func(context.Context, Control)) error {
return nil
}
// sends is one message arriving, as the controller reads it.
func (c *fakeInbound) sends(t *testing.T, to *settled, kind string, v any) Control {
t.Helper()
body, err := json.Marshal(v)
if err != nil {
t.Fatal(err)
}
tag++
return &fakeControl{kind: kind, body: body, tag: tag, to: to, on: c}
}
// retries hands every held message back to the controller, the way a stream redelivers.
func (c *fakeInbound) retries(ctx context.Context, s *Server) {
for tag, m := range c.held {
delete(c.held, tag)
m.redelivered = true
s.act(ctx, m)
}
}
type fakeControl struct {
kind string
body []byte
tag uint64
to *settled
on *fakeInbound
redelivered bool
since time.Time
}
func (m *fakeControl) Kind() string { return m.kind }
func (m *fakeControl) Body() []byte { return m.body }
func (m *fakeControl) Redelivered() bool { return m.redelivered }
func (m *fakeControl) About(string) {}
func (m *fakeControl) Answer(context.Context, []byte) error { return nil }
func (m *fakeControl) Took() error { m.to.acked = true; return nil }
func (m *fakeControl) Drop() error { m.to.rejected = true; return nil }
// HeldFor is how long the store has been waited on for this message — zero on a first delivery.
func (m *fakeControl) HeldFor() time.Duration {
if m.since.IsZero() {
return 0
}
return time.Since(m.since)
}
func (m *fakeControl) Hold(time.Duration) error {
if m.since.IsZero() {
m.since = time.Now()
}
m.on.held[m.tag] = m
return nil
}
+7 -9
View File
@@ -23,6 +23,10 @@ import (
// it first could not take it, and a controller that restarts mid-window has nothing to lose. // it first could not take it, and a controller that restarts mid-window has nothing to lose.
// natsInbound consumes what nodes and modules say over NATS. // natsInbound consumes what nodes and modules say over NATS.
// Prefetch is how many controls the controller holds unacknowledged at once: the store window
// (ADR 0083) is the server's, and this is what it may hand this process ahead of its acting.
const Prefetch = 64
type natsInbound struct { type natsInbound struct {
js *broker.JetStream js *broker.JetStream
// follows is the kinds asked for beyond what nodes say (Also). The events those are are the // follows is the kinds asked for beyond what nodes say (Also). The events those are are the
@@ -48,7 +52,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 +184,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
@@ -220,14 +226,6 @@ func (m *natsControl) HeldFor() time.Duration {
return time.Since(first) return time.Since(first)
} }
// About is nothing here, and that is the point.
//
// Setting a held message aside when a newer one about the same thing arrives is what a controller
// holding deliveries in memory can do. A naked message belongs to the server and comes back
// whatever happened meanwhile, so the question "is this the past?" is answered by what the message
// says instead — the digest of the declaration a report is about (window.go, design 25 §3).
func (m *natsControl) About(string) {}
// Answer publishes to the reply subject the request carries **in its payload**. // Answer publishes to the reply subject the request carries **in its payload**.
// //
// Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had // Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had
+21
View File
@@ -4,6 +4,8 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"errors" "errors"
"io"
"log"
"os" "os"
"sync" "sync"
"testing" "testing"
@@ -120,6 +122,8 @@ func eventually(t *testing.T, what string, is func() bool) {
// A report published by a node reaches the controller, is recorded, and is acknowledged — so the // A report published by a node reaches the controller, is recorded, and is acknowledged — so the
// stream does not hold it. A work queue is the check: what is acknowledged leaves it. // stream does not hold it. A work queue is the check: what is acknowledged leaves it.
func quiet() *log.Logger { return log.New(io.Discard, "", 0) }
func TestNatsAReportIsHeardAndLeavesTheStream(t *testing.T) { func TestNatsAReportIsHeardAndLeavesTheStream(t *testing.T) {
js := aBus(t) js := aBus(t)
store := &counted{} store := &counted{}
@@ -374,3 +378,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")
}
}
-18
View File
@@ -46,24 +46,6 @@ func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) {
} }
} }
// A newer report from the same node supersedes one of its reports still held: recorded after the
// newer, the older would overwrite what the node is doing now.
func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) {
s, in := serving()
s.listener = heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}
older, newer, other := &settled{}, &settled{}, &settled{}
s.act(context.Background(), in.sends(t, older, KindReport, aReport("anchor", "d1")))
s.act(context.Background(), in.sends(t, other, KindReport, aReport("laptop", "d7")))
s.act(context.Background(), in.sends(t, newer, KindReport, aReport("anchor", "d2")))
if !older.acked {
t.Fatalf("the older report was not set aside by the newer: %+v", older)
}
if !newer.unsettled() || !other.unsettled() || len(in.held) != 2 {
t.Fatalf("the newer report and another node's were not both held: newer %+v other %+v, %d held",
newer, other, len(in.held))
}
}
// A store that has not come back within the bound is not restarting: the report is let go, loudly, // A store that has not come back within the bound is not restarting: the report is let go, loudly,
// rather than held for ever. // rather than held for ever.
func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) { func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
+36 -81
View File
@@ -7,13 +7,10 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"github.com/novox/mesh-controller/internal/broker"
"log" "log"
"os" "os"
"time" "time"
amqp "github.com/rabbitmq/amqp091-go"
"github.com/novox/mesh-controller/internal/envfile"
) )
// AMQPVar is the controller's own connection to the bus the mesh runs on today. // AMQPVar is the controller's own connection to the bus the mesh runs on today.
@@ -57,6 +54,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.
@@ -68,8 +67,7 @@ type Upgrader interface {
type Server struct { type Server struct {
inbound Inbound inbound Inbound
bus Bus bus Bus
conn *amqp.Connection js *broker.JetStream
channel *amqp.Channel
enroller Enroller enroller Enroller
listener Listener listener Listener
@@ -99,6 +97,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
} }
@@ -120,82 +121,20 @@ func (s *Server) Answers(r Replayer) error {
// On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq // On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq
// 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may // 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may
// have moved since. // have moved since.
func Connect(enroller Enroller, listener Listener) (*Server, error) { // Connect is ConnectNats: the mesh has one bus (novox/hq ADR 0131, design 28 task 5.5). Kept as
url, err := envfile.Placed(AMQPVar) // the name callers know; the transport it opened before is gone with the bus it spoke to.
if err != nil { func Connect(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
return nil, err return ConnectNats(js, enroller, listener)
}
if url == "" {
return nil, fmt.Errorf(
"this control plane has no %s, so it cannot reach its broker. Nodes talk to it over "+
"the broker and nowhere else, so without this it can hold records and answer "+
"nothing", AMQPVar)
}
conn, err := amqp.Dial(url)
if err != nil {
// Not quoted back: the URL carries the controller's own bus password.
return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err)
}
channel, err := conn.Channel()
if err != nil {
conn.Close()
return nil, err
}
// Declared here rather than assumed. The controller is the only thing that may create them — a
// node's account can write to this exchange and read its own queue, and configure nothing
// else, so a node arriving before the controller has ever run finds nothing and says so,
// rather than quietly creating a topology nobody designed.
if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err)
}
// The events exchange too. The controller is not the only publisher on it — modules announce
// onto it with their own accounts — but it is the only thing permitted to create it, for the
// same reason it is the only thing permitted to create the direct one.
if err := channel.ExchangeDeclare(EventsExchange, "topic", true, false, false, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s exchange: %w", EventsExchange, err)
}
if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err)
}
// Every key a node may publish. Binding one and forgetting another is a message the bus
// accepts, finds no queue for, and drops — the publisher sees success and the consumer sees
// nothing. That is exactly what happened to reports: `report` was left unbound while `enrol`
// worked, so nodes announced what they had applied into a void for an afternoon.
for _, key := range []string{KeyEnrol, KeyReport, KeyAlive, KeyBuilt} {
if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil {
conn.Close()
return nil, fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err)
}
}
return &Server{
inbound: Current(conn, channel),
bus: OverCurrent{Channel: channel},
conn: conn,
channel: channel,
enroller: enroller,
listener: listener,
log: log.New(os.Stdout, "", log.LstdFlags),
}, nil
} }
// Channel is the controller's channel, for the command line's own publishing. func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
func (s *Server) Channel() *amqp.Channel { return s.channel }
func (s *Server) Close() { func (s *Server) Close() {
if s.inbound != nil { if s.inbound != nil {
s.inbound.Close() s.inbound.Close()
} }
if s.channel != nil { if s.js != nil {
_ = s.channel.Close() s.js.Close()
}
if s.conn != nil {
_ = s.conn.Close()
} }
} }
@@ -225,6 +164,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:
@@ -335,7 +276,6 @@ func (s *Server) reported(ctx context.Context, m Control) {
_ = m.Drop() _ = m.Drop()
return return
} }
m.About("report " + report.Node)
what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared)
if s.listener != nil { if s.listener != nil {
@@ -468,9 +408,6 @@ func (s *Server) wasBuilt(ctx context.Context, m Control) {
_ = m.Drop() _ = m.Drop()
return return
} }
// Each build result its own subject: none supersedes another, and recording one twice is
// harmless — the build is kept by its id.
m.About("build " + digest(m.Body()))
err := s.recorder.Built(ctx, result) err := s.recorder.Built(ctx, result)
switch s.decide(ctx, m, fmt.Sprintf("a build result from %s", result.On), "", "", err) { switch s.decide(ctx, m, fmt.Sprintf("a build result from %s", result.On), "", "", err) {
case Hold: case Hold:
@@ -502,7 +439,6 @@ func (s *Server) wasBuilt(ctx context.Context, m Control) {
// the catalogue next restarts (issue 083). One request stands for all: a newer one supersedes one // the catalogue next restarts (issue 083). One request stands for all: a newer one supersedes one
// still held. // still held.
func (s *Server) catchingUp(ctx context.Context, m Control) { func (s *Server) catchingUp(ctx context.Context, m Control) {
m.About("catch-up")
if s.replayer == nil { if s.replayer == nil {
s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay") s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay")
_ = m.Took() _ = m.Took()
@@ -555,7 +491,6 @@ func (s *Server) moved(ctx context.Context, m Control) {
_ = m.Took() _ = m.Took()
return return
} }
m.About("upgrade " + u.Module)
if u.Module == "" { if u.Module == "" {
s.log.Printf("an upgrade announcement named no module; ignored") s.log.Printf("an upgrade announcement named no module; ignored")
_ = m.Took() _ = m.Took()
@@ -596,3 +531,23 @@ 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
}
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 -48
View File
@@ -3,7 +3,6 @@ package link
import ( import (
"context" "context"
"errors" "errors"
"fmt"
"testing" "testing"
"github.com/jackc/pgx/v5/pgconn" "github.com/jackc/pgx/v5/pgconn"
@@ -18,7 +17,8 @@ func (r recordsWith) Built(context.Context, BuildResult) error { return r.err }
type upgradesWith struct{ err error } type upgradesWith struct{ err error }
func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err } func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err }
func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return u.err }
type replaysWith struct{ err error } type replaysWith struct{ err error }
@@ -106,49 +106,3 @@ func TestAMessageHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to) t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to)
} }
} }
// Two identical build results: the newer sets the older aside rather than leaving it unsettled
// for ever, holding a place in the prefetch.
func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) {
s, in := serving()
s.recorder = recordsWith{err: restarting}
built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}
first, second := &settled{}, &settled{}
s.act(context.Background(), in.sends(t, first, KindBuilt, built))
s.act(context.Background(), in.sends(t, second, KindBuilt, built))
if !first.acked || !second.unsettled() || len(in.held) != 1 {
t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held",
*first, *second, len(in.held))
}
}
// What is held stops short of the prefetch, so the loop always has room to answer an enrolment.
func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) {
s, in := serving()
s.recorder = recordsWith{err: restarting}
var last *settled
for i := 0; i < Prefetch; i++ {
last = &settled{}
s.act(context.Background(), in.sends(t, last, KindBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)}))
}
if len(in.held) != Prefetch-PrefetchHeadroom {
t.Fatalf("%d messages were held; the ceiling is %d", len(in.held), Prefetch-PrefetchHeadroom)
}
if last.unsettled() {
t.Fatalf("a message past the ceiling was held: %+v", *last)
}
}
// An upgrade handled during shutdown is left for the bus too — the upgrader's error is the
// cancelled context, which is no answer about the announcement.
func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
s, in := serving()
s.upgrader = upgradesWith{err: context.Canceled}
to := &settled{}
s.act(ctx, in.sends(t, to, KindModuleMoved, Upgraded{Module: "gitea", Commit: "abcdef0123"}))
if !to.unsettled() {
t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to)
}
}
+1
View File
@@ -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"
], ],