Compare commits
43
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d1e488efaf | ||
|
|
c5dc7e732a | ||
|
|
7efcccd013 | ||
|
|
964285f08c | ||
|
|
5698dda11f | ||
|
|
6005a8471f | ||
|
|
e2ee0dfe98 | ||
|
|
70341cfbc7 | ||
|
|
77643aa3f4 | ||
|
|
1fd6194ff8 | ||
|
|
c37018fdd2 | ||
|
|
3907ea0db0 | ||
|
|
40f5e9a41c | ||
|
|
2c2eb51878 | ||
|
|
aa2d0b51ea | ||
|
|
64d154d9d7 | ||
|
|
ffa390f916 | ||
|
|
4d62e6caf1 | ||
|
|
386ae676ca | ||
|
|
f8a9c3d6bc | ||
|
|
83671fae5f | ||
|
|
9b715524a2 | ||
|
|
e06fc1ed16 | ||
|
|
13d7c5c5dd | ||
|
|
7b02feaebb | ||
|
|
bd10e2c695 | ||
|
|
4b209d944d | ||
|
|
84024cbdb6 | ||
|
|
ad2eed2f71 | ||
|
|
e8aa7ed9e7 | ||
|
|
337aaea123 | ||
|
|
4c41628b20 | ||
|
|
ae7fb520d7 | ||
|
|
f325073982 | ||
|
|
33c4e4be34 | ||
|
|
585a6abbdd | ||
|
|
8d52a2cfb0 | ||
|
|
1cfe6be9c4 | ||
|
|
4e4481b6f2 | ||
|
|
63ca073938 | ||
|
|
b244a768a3 | ||
|
|
81d52719f4 | ||
|
|
f6a93fe74c |
@@ -135,6 +135,16 @@ func run() error {
|
||||
// machine told about both would take work from one and answer on the other, and every log line would
|
||||
// say it was fine.
|
||||
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
||||
// **The credential decides, before any variable does.** A machine moved to the new bus was
|
||||
// handed a credential for it and nothing else changed in its environment; that credential
|
||||
// names the bus by scheme, so it is enough to know which bus to take work from.
|
||||
if credential.onTheNewBus() {
|
||||
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
}
|
||||
address, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -460,10 +470,26 @@ func brokerFrom() (Credential, error) {
|
||||
// **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels
|
||||
// out of band — here, sealed with the credential — and the endpoint is verified once at connect.
|
||||
type Credential struct {
|
||||
URL string `json:"url"`
|
||||
// Fingerprint is SHA-256 over the broker certificate's DER bytes, or empty to verify the
|
||||
// ordinary way.
|
||||
URL string `json:"url"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
// User and Password ride beside the address on the bus being built (design 25): a credential
|
||||
// embedded in a URL leaks into every log line that prints a connection, so the mesh seals them
|
||||
// as two fields and this machine joins them once, here, to dial.
|
||||
User string `json:"user,omitempty"`
|
||||
Password string `json:"password,omitempty"`
|
||||
}
|
||||
|
||||
// onTheNewBus is whether a credential is for the bus being built: its address says so, and the
|
||||
// mesh only ever seals such a credential with the user and password beside it.
|
||||
func (c Credential) onTheNewBus() bool { return strings.HasPrefix(strings.TrimSpace(c.URL), "nats://") }
|
||||
|
||||
// natsURL is the address with this machine's credential in it, for the one dial that needs it.
|
||||
func (c Credential) natsURL() string {
|
||||
rest := strings.TrimPrefix(strings.TrimSpace(c.URL), "nats://")
|
||||
if c.User == "" {
|
||||
return "nats://" + rest
|
||||
}
|
||||
return "nats://" + c.User + ":" + c.Password + "@" + rest
|
||||
}
|
||||
|
||||
// dial opens the connection, pinning the broker's certificate when there is one to pin.
|
||||
|
||||
@@ -14,8 +14,8 @@ func TestBuilderDiagnosticsStayOffStdout(t *testing.T) {
|
||||
allowed := map[string]bool{
|
||||
"string(body)": true, // once.go: the result JSON, which IS stdout
|
||||
"version)": true, // --version
|
||||
`"stopping")`: true, // the loop.s shutdown line
|
||||
"usage)": true, // --help text, for a human
|
||||
`"stopping")`: true, // the loop.s shutdown line
|
||||
"usage)": true, // --help text, for a human
|
||||
}
|
||||
for _, file := range []string{"once.go", "main.go"} {
|
||||
src, err := os.ReadFile(file)
|
||||
|
||||
@@ -37,7 +37,7 @@ func askCommand(ctx context.Context, args []string) error {
|
||||
arguments = json.RawMessage(positionals[2])
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -368,7 +368,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
|
||||
}
|
||||
defer ident.Close()
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -472,7 +472,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
|
||||
return err
|
||||
}
|
||||
defer ident.Close()
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -614,6 +614,15 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return issueWith(ctx, inv, m, node, busAddress, known, reachable, user, password)
|
||||
}
|
||||
|
||||
// issueWith is the delivery half: the minted password sealed to the machine as the module's broker
|
||||
// secret, and the module's consumer created where the bus can be reached. Split from the minting
|
||||
// so the move can issue every module against a bus whose address it worked out itself
|
||||
// (`rollout mint`, design 28 task 5.2) rather than the one in this process's environment.
|
||||
func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manifest,
|
||||
node, busAddress string, known broker.Broker, reachable, user, password string) error {
|
||||
held, err := json.Marshal(struct {
|
||||
URL string `json:"url"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
@@ -639,14 +648,19 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
|
||||
Kind: broker.KindModule, Node: node, Module: m.Module,
|
||||
Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools,
|
||||
}); needed {
|
||||
js, err := broker.Dial(busAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
|
||||
"how %s hears what it consumes: %w", m.Module, err)
|
||||
}
|
||||
defer js.Close()
|
||||
if err := js.EnsureConsumer(consumer); err != nil {
|
||||
return err
|
||||
if busAddress == "" {
|
||||
fmt.Printf(" %s consumes; its consumer is created when the bus is reachable (`push`, then "+
|
||||
"`rollout mint` again is harmless)\n", m.Module)
|
||||
} else {
|
||||
js, err := broker.Dial(busAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
|
||||
"how %s hears what it consumes: %w", m.Module, err)
|
||||
}
|
||||
defer js.Close()
|
||||
if err := js.EnsureConsumer(consumer); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -537,6 +537,15 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// **With the seat holders on record**, or a machine running the next holder of a seat beside
|
||||
// the current one resolves as two holders, is refused, and drops out of the map — taking the
|
||||
// address every other machine composes for what it offers (novox/hq ADR 0131). Found live:
|
||||
// the control node vanished from the private network the moment the new bus was assigned
|
||||
// beside the old one.
|
||||
holdings, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []inventory.Overlay
|
||||
for _, p := range places {
|
||||
if p.Address == "" {
|
||||
@@ -549,7 +558,7 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
|
||||
caps, _ := inv.ProfileOf(ctx, p.Name)
|
||||
got, err := catalogue.Resolve(shelf, assigned,
|
||||
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps},
|
||||
catalogue.World{Unchecked: true})
|
||||
catalogue.World{Unchecked: true, Holdings: holdings})
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -185,6 +185,15 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
|
||||
|
||||
// Every node, not only the placed ones. A machine that was never put on the private network
|
||||
// still runs modules, still holds claims, and still offers whatever it offers.
|
||||
// **Who holds each seat on record, before anything is resolved** (novox/hq ADR 0131). Both
|
||||
// passes below need it: without it, the assignment standing beside a seat's holder — the next
|
||||
// holder, waiting for the handover — is refused as a second holder, and its node's whole set
|
||||
// with it.
|
||||
holdings, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
return catalogue.World{}, err
|
||||
}
|
||||
|
||||
nodes, err := inv.Nodes(ctx)
|
||||
if err != nil {
|
||||
return catalogue.World{}, err
|
||||
@@ -227,7 +236,7 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
|
||||
offered := map[string][]catalogue.Provider{}
|
||||
var firstHeld []catalogue.Held
|
||||
for _, o := range others {
|
||||
got, err := catalogue.Resolve(shelf, o.assigned, o.node, catalogue.World{Unchecked: true})
|
||||
got, err := catalogue.Resolve(shelf, o.assigned, o.node, catalogue.World{Unchecked: true, Holdings: holdings})
|
||||
if err != nil {
|
||||
// Their set does not resolve for some other reason. Not this node's problem to
|
||||
// report, and nothing of theirs is running, so it offers nothing.
|
||||
@@ -258,7 +267,7 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
|
||||
// with several providers (novox/hq ADR 0110), so a node consuming one resolves only once the
|
||||
// holder is known. Without them its set is refused here, and a refused node's own claims drop
|
||||
// out of what the mesh holds — so a second holder of one of its seats would pass unrefused.
|
||||
world := catalogue.World{Offered: offered, Held: firstHeld}
|
||||
world := catalogue.World{Offered: offered, Held: firstHeld, Holdings: holdings}
|
||||
var held []catalogue.Held
|
||||
for _, o := range others {
|
||||
got, err := catalogue.Resolve(shelf, o.assigned, o.node, world)
|
||||
@@ -627,8 +636,13 @@ func renderingFor(ctx context.Context, open *stores, node string,
|
||||
if err != nil {
|
||||
return catalogue.Rendering{}, inventory.Node{}, err
|
||||
}
|
||||
memberships, err := inv.BusMemberships(ctx)
|
||||
if err != nil {
|
||||
return catalogue.Rendering{}, inventory.Node{}, err
|
||||
}
|
||||
return catalogue.Rendering{
|
||||
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
|
||||
BusMembership: memberships[node],
|
||||
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
|
||||
Certificate: certificate, Authority: authority, Mesh: private, Names: names,
|
||||
Machines: machines,
|
||||
Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation,
|
||||
|
||||
+80
-16
@@ -38,6 +38,33 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
|
||||
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
||||
|
||||
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
||||
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
|
||||
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
|
||||
// so nothing served here finds them missing.
|
||||
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
|
||||
busAddress, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !onNATS {
|
||||
return link.Connect(enroller, listener)
|
||||
}
|
||||
if inv != nil {
|
||||
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
js, err := broker.Dial(busAddress)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
||||
broker.BareAddress(busAddress), err)
|
||||
}
|
||||
return link.ConnectNats(js, enroller, listener), nil
|
||||
}
|
||||
|
||||
func serve(ctx context.Context) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
@@ -90,7 +117,7 @@ func serve(ctx context.Context) error {
|
||||
|
||||
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
|
||||
OnNATS: onNATS}
|
||||
server, err := link.Connect(work, work)
|
||||
server, err := connectLink(ctx, inv, work, work)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -100,11 +127,6 @@ func serve(ctx context.Context) error {
|
||||
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
|
||||
// replaced all have records and no objects — and a node whose consumer is missing hears nothing
|
||||
// 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`
|
||||
// would otherwise be reported into the void, which is the same as not reporting it.
|
||||
server.Records(builds{inv})
|
||||
@@ -158,13 +180,13 @@ func declare(ctx context.Context, args []string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
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
|
||||
}
|
||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||
@@ -271,7 +293,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -332,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
if err != nil {
|
||||
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
|
||||
}
|
||||
// After it is away, not before. A digest recorded for something that failed to send would
|
||||
@@ -415,7 +437,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||
},
|
||||
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 {
|
||||
return err
|
||||
}
|
||||
@@ -628,7 +650,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
||||
len(refusals), strings.Join(refusals, "\n\n"))
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -639,7 +661,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
||||
if err != nil {
|
||||
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
|
||||
}
|
||||
record, err := inv.NodeByName(ctx, s.node)
|
||||
@@ -704,7 +726,7 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
js, err := broker.Dial(address)
|
||||
if err != nil {
|
||||
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()
|
||||
|
||||
@@ -739,10 +761,52 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
// The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder,
|
||||
// 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.
|
||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), nil); err != nil {
|
||||
// With the seats' holders, so each role's work queue gets the consumer its holder takes
|
||||
// work from. Passed as nil until the first live raise, which left the build machine bound to a
|
||||
// consumer nothing had created (2026-09-28).
|
||||
holders, err := seatHolders(ctx, inv)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
|
||||
address, len(names))
|
||||
broker.BareAddress(address), len(names))
|
||||
return nil
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
@@ -12,6 +13,7 @@ import (
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/secrets"
|
||||
)
|
||||
|
||||
// Moving the mesh's own traffic to the bus being built (novox/hq ADR 0116 step 5).
|
||||
@@ -25,18 +27,25 @@ import (
|
||||
// against a mesh that is serving. It answers from records: what is missing, and what would happen.
|
||||
// `rollout` itself refuses unless the check is clean.
|
||||
//
|
||||
// **The old broker is not switched off by this.** It stays an ordinary provider of `amqp` for whatever
|
||||
// else uses it — on this installation, a whole automation layer that has nothing to do with the mesh
|
||||
// ([ADR 0119](../../02-DECISIONS/0119-amqp-is-a-provision-not-the-bus.md)). Only the mesh's own
|
||||
// traffic moves, which is why this is survivable at all: what breaks if it goes wrong is the mesh's
|
||||
// ability to change things, not the services its modules are serving.
|
||||
// **The old broker goes with the move, and goes last** (novox/hq ADR 0131): AMQP is not a provision,
|
||||
// so once every machine reports on the new bus its module is unassigned. Only the mesh's own traffic
|
||||
// is what moves, which is why this is survivable at all: what breaks if it goes wrong is the mesh's
|
||||
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when
|
||||
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
|
||||
|
||||
const rolloutUsage = "rollout check | rollout --confirm"
|
||||
const rolloutUsage = "rollout check | rollout mint [--again] | rollout --confirm"
|
||||
|
||||
func rolloutCommand(ctx context.Context, args []string) error {
|
||||
switch {
|
||||
case len(args) == 1 && args[0] == "check":
|
||||
return rolloutCheck(ctx)
|
||||
case len(args) == 1 && args[0] == "mint":
|
||||
return rolloutMint(ctx, false)
|
||||
case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
|
||||
// Every credential minted afresh, whether or not one exists — for a mint that was wrong
|
||||
// before anything was pushed. Afterwards nothing that received the old one still works,
|
||||
// which is fine exactly when nothing received it.
|
||||
return rolloutMint(ctx, true)
|
||||
case len(args) == 1 && args[0] == "--confirm":
|
||||
return errors.New(
|
||||
"the rollout itself is not built yet: `rollout check` answers whether it could run, and " +
|
||||
@@ -105,7 +114,6 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines
|
||||
ModuleCredentialled: map[string]bool{},
|
||||
// The old broker keeps its other clients on this installation, and saying so is how the plan
|
||||
// stops reading as a retirement.
|
||||
OldBusHasOtherClients: true,
|
||||
}
|
||||
|
||||
address, _, err := broker.OnNATS()
|
||||
@@ -191,3 +199,191 @@ func wasSentTheUserList(ctx context.Context, inv *inventory.Inventory, node stri
|
||||
// notReadyOf is the readiness reasoning, named here so a test can reach it without the command's
|
||||
// printing. The reasoning itself is the broker package's, where it is pure.
|
||||
func notReadyOf(state broker.Readiness) []string { return broker.NotReady(state) }
|
||||
|
||||
// rolloutMint gives every principal the new bus will have a credential it does not yet have, and
|
||||
// puts each where its owner reads it (novox/hq design 28, task 5.2): a machine's as a membership
|
||||
// sealed into its declaration, a module's as its broker secret, the control plane's own as its
|
||||
// `bus` secret. Idempotent: what already has a hash is left alone, so running it again is harmless.
|
||||
//
|
||||
// **Before anything moves, and it is what makes moving possible.** A machine moved without a
|
||||
// credential cannot come back, and afterwards there is no bus to tell it anything over — which is
|
||||
// why `rollout check` refuses until this has run. The bus's address is worked out here, from where
|
||||
// the module that provides it is assigned, rather than read from this process's environment: this
|
||||
// process is still on the old bus when this runs, and must be.
|
||||
func rolloutMint(ctx context.Context, again bool) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
inv := open.inventory
|
||||
|
||||
known, err := broker.FromEnvironment()
|
||||
if err != nil {
|
||||
return fmt.Errorf("the bus's certificate is not known to this process, and every membership "+
|
||||
"must carry its fingerprint: %w", err)
|
||||
}
|
||||
shelf, err := inv.Catalogue(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var busNode, controllerNode string
|
||||
for _, e := range entries {
|
||||
switch {
|
||||
case e.Manifest.ClaimsSeat("mesh-broker") && providesBus(e.Manifest) && len(e.On) > 0:
|
||||
busNode = e.On[0]
|
||||
case e.Manifest.Module == "mesh-controller" && len(e.On) > 0:
|
||||
controllerNode = e.On[0]
|
||||
}
|
||||
}
|
||||
if busNode == "" {
|
||||
return errors.New("no assigned module provides mesh-bus and claims mesh-broker, so there is no " +
|
||||
"bus to mint credentials for — register and assign it first")
|
||||
}
|
||||
onNetwork, err := whereEveryoneIs(ctx, inv, shelf)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
busHost := onNetwork[busNode]
|
||||
if busHost == "" {
|
||||
// **The hub is not in that map.** The machine that took over the tunnel is where the current
|
||||
// bus already answers, and every machine dials it at the address the mesh handed them — so
|
||||
// when the new bus runs on the same machine, that address is the one to tell them, with the
|
||||
// new port. Found live: the control node is the hub, and the map lists the machines placed
|
||||
// around it.
|
||||
// The host alone: no scheme (BareAddress adds one where none was, which is the wrong
|
||||
// direction here — every URL built below adds its own) and no port.
|
||||
_, _, host := broker.CredentialIn(known.Address)
|
||||
if host == "" {
|
||||
host = known.Address
|
||||
}
|
||||
if _, after, hasScheme := strings.Cut(host, "://"); hasScheme {
|
||||
host = after
|
||||
}
|
||||
host = strings.TrimSpace(host)
|
||||
if i := strings.LastIndex(host, ":"); i > 0 && !strings.Contains(host[i:], "]") {
|
||||
host = host[:i]
|
||||
}
|
||||
if host == "" {
|
||||
return fmt.Errorf("%s runs the new bus and has no address on the private network, and the "+
|
||||
"current bus's address is unknown too, so no machine could be told where it is", busNode)
|
||||
}
|
||||
busHost = host
|
||||
}
|
||||
busAddress := busHost + ":4222"
|
||||
|
||||
records, err := inv.BusRecords(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
users, err := broker.Users(records)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
kept, err := inv.BusUsers(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
hashes := make(map[string]string, len(kept))
|
||||
for name, u := range kept {
|
||||
hashes[name] = u.PasswordHash
|
||||
}
|
||||
_, missing := broker.WithPasswords(users, hashes)
|
||||
wanted := map[string]bool{}
|
||||
for _, m := range missing {
|
||||
wanted[m] = true
|
||||
}
|
||||
|
||||
var machines, modules, skipped int
|
||||
for _, p := range users {
|
||||
if !again && !wanted[p.Username()] {
|
||||
continue
|
||||
}
|
||||
switch p.Kind {
|
||||
case broker.KindController:
|
||||
if controllerNode == "" {
|
||||
return errors.New("the control plane is not assigned anywhere, so its credential has nowhere to go")
|
||||
}
|
||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusController})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
url := "nats://" + p.Username() + ":" + password + "@" + busAddress
|
||||
if err := inv.AcceptSecretForModule(ctx, controllerNode, "mesh-controller", "bus", url); err != nil {
|
||||
return fmt.Errorf("the control plane's credential is minted and could not be sealed to %s: %w", controllerNode, err)
|
||||
}
|
||||
fmt.Printf("control plane: credential minted, sealed to %s as its `bus` secret\n", controllerNode)
|
||||
|
||||
case broker.KindNode:
|
||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: p.Node})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
membership, _ := json.Marshal(map[string]string{
|
||||
"broker": busAddress, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
|
||||
})
|
||||
key, err := inv.SealingKeyOf(ctx, p.Node)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%s has no sealing key, so its membership cannot be sealed to it: %w", p.Node, err)
|
||||
}
|
||||
sealed, err := secrets.Seal(key, membership)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := inv.PutBusMembership(ctx, p.Node, sealed); err != nil {
|
||||
return err
|
||||
}
|
||||
machines++
|
||||
|
||||
case broker.KindModule:
|
||||
if p.Module == "mesh-controller" {
|
||||
// The control plane is a module too, and its `broker` secret is the old bus's
|
||||
// credential it is still using while this runs. Writing the new bus's blob there
|
||||
// cut the mesh off from its own old bus mid-move (2026-09-28). Its new-bus credential
|
||||
// is the controller principal's `bus` secret above; nothing else is needed here.
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
m, inShelf := shelf[p.Module]
|
||||
if !inShelf {
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
if _, reads := m.OwnSecrets["broker"]; !reads {
|
||||
fmt.Printf(" %s on %s speaks on the bus but declares no `broker` secret to receive a credential in; skipped\n", p.Module, p.Node)
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusModule, Node: p.Node, Module: p.Module})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := issueWith(ctx, inv, m, p.Node, "", known, busAddress, p.Username(), password); err != nil {
|
||||
return err
|
||||
}
|
||||
modules++
|
||||
|
||||
default:
|
||||
skipped++
|
||||
}
|
||||
}
|
||||
fmt.Printf("minted for %d machine(s) and %d module runtime(s); %d skipped; the bus is at %s\n",
|
||||
machines, modules, skipped, busAddress)
|
||||
fmt.Println(" each machine's membership and each module's credential arrive with the next push of its machine;")
|
||||
fmt.Println(" push the machine running the bus first, so the bus stands with its user list before anything dials it")
|
||||
return nil
|
||||
}
|
||||
|
||||
// providesBus is whether a manifest provides the mesh's bus.
|
||||
func providesBus(m catalogue.Manifest) bool {
|
||||
for _, o := range m.Provides {
|
||||
if o.Name == "mesh-bus" {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"os"
|
||||
"slices"
|
||||
"sort"
|
||||
"strings"
|
||||
"text/tabwriter"
|
||||
@@ -85,7 +86,8 @@ func seatsHeld(seats []catalogue.Seat, held []catalogue.Held) ([]seatRow, []cata
|
||||
return rows, outside
|
||||
}
|
||||
|
||||
// seatCommand changes the set — the whole point of it being data (novox/hq ADR 0122).
|
||||
// seatCommand changes the set — the whole point of it being data (novox/hq ADR 0122) — and, since
|
||||
// ADR 0131, changes who holds a seat.
|
||||
func seatCommand(ctx context.Context, args []string) error {
|
||||
if len(args) == 3 && args[0] == "rename" {
|
||||
from, to := args[1], args[2]
|
||||
@@ -101,7 +103,101 @@ func seatCommand(ctx context.Context, args []string) error {
|
||||
"re-registered or frozen (novox/hq ADR 0122)\n", from, to)
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("seat rename <from> <to>")
|
||||
if len(args) == 3 && args[1] == "--to" {
|
||||
return handOver(ctx, args[0], args[2])
|
||||
}
|
||||
return fmt.Errorf("seat rename <from> <to> | seat <name> --to <node>/<module>")
|
||||
}
|
||||
|
||||
// handOver makes one assignment the holder of a seat, as one act, so the seat is never without a
|
||||
// holder in between (novox/hq ADR 0131, design 28 task 5.3). The control plane finds its own bus
|
||||
// through one of these seats; the day it was left empty mid-change is why this exists.
|
||||
//
|
||||
// Everything that could make the new holder wrong is refused here, before the row is written: the
|
||||
// seat must exist, the assignment must exist, and the module must be able to hold the seat —
|
||||
// claim it at its scope and provide what it delivers, judged against the store's row. What is
|
||||
// **not** checked is whether the module is running yet: that is what `push` confirms afterwards,
|
||||
// and refusing to record a handover to a module the node has not started would make the handover
|
||||
// impossible to do before the switch instead of as the switch.
|
||||
func handOver(ctx context.Context, seatName, to string) error {
|
||||
nodeName, module, ok := strings.Cut(to, "/")
|
||||
if !ok || nodeName == "" || module == "" {
|
||||
return fmt.Errorf("the new holder is named <node>/<module>, not %q", to)
|
||||
}
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
inv := open.inventory
|
||||
|
||||
seat, known := catalogue.SeatNamed(seatName)
|
||||
if !known {
|
||||
return fmt.Errorf("%q is not a seat this mesh defines — `seats` lists them", seatName)
|
||||
}
|
||||
assigned, err := inv.Assigned(ctx, nodeName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !slices.Contains(assigned, module) {
|
||||
return fmt.Errorf("%s is not assigned to %s, so it cannot hold anything there — "+
|
||||
"`assign %s %s` first", module, nodeName, nodeName, module)
|
||||
}
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var m *catalogue.Manifest
|
||||
for i := range entries {
|
||||
if entries[i].Manifest.Module == module {
|
||||
m = &entries[i].Manifest
|
||||
}
|
||||
}
|
||||
if m == nil {
|
||||
return fmt.Errorf("%s is assigned but not in the catalogue, which should not happen", module)
|
||||
}
|
||||
var was string
|
||||
holdings, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, h := range holdings {
|
||||
if hs, ok := catalogue.SeatNamed(h.Claim); ok && hs.Name == seat.Name {
|
||||
was = h.Node
|
||||
}
|
||||
}
|
||||
|
||||
// **Recording who already holds the seat is not making a new holder, and is not judged like
|
||||
// one.** On a mesh that predates the record, the first handover has to begin by writing down
|
||||
// the standing holder — otherwise the next holder cannot be assigned beside it, because two
|
||||
// eligible claimants with nothing on record are refused. That standing holder may no longer
|
||||
// satisfy what the seat delivers (the row moved under it, on purpose, as ADR 0131's first step),
|
||||
// and it holds regardless: derivation never read that column. So when nothing is on record and
|
||||
// the named assignment is the one holding by derivation, only the claim itself is checked here.
|
||||
// Every *change* of holder is judged in full.
|
||||
claimsIt := false
|
||||
for _, c := range m.Claims {
|
||||
if cs, ok := catalogue.SeatNamed(c.Name); ok && cs.Name == seat.Name && c.At() == seat.Scope {
|
||||
claimsIt = true
|
||||
}
|
||||
}
|
||||
if was == "" && claimsIt {
|
||||
fmt.Printf("nothing was on record for %s; recording %s on %s as its standing holder\n",
|
||||
seat.Name, module, nodeName)
|
||||
} else if err := catalogue.CanHold(*m, seat); err != nil {
|
||||
return fmt.Errorf("%s cannot hold %s: %w", module, seat.Name, err)
|
||||
}
|
||||
if err := inv.HoldSeat(ctx, seat.Name, seat.Scope, nodeName, module); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("%s is held by %s on %s\n", seat.Name, module, nodeName)
|
||||
if was != "" && was != nodeName {
|
||||
fmt.Printf(" `push %s` and `push %s` send both machines what changed\n", was, nodeName)
|
||||
} else {
|
||||
fmt.Printf(" `push %s` sends the machine what changed; every other machine that reads the "+
|
||||
"seat is re-declared by `push --behind`\n", nodeName)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func seatsCommand(ctx context.Context, args []string) error {
|
||||
|
||||
@@ -28,24 +28,14 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu
|
||||
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.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
|
||||
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
|
||||
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/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE=
|
||||
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU=
|
||||
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
|
||||
golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU=
|
||||
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
|
||||
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sync v0.23.0 h1:KameEIfc1IkluZyXWLn39Wd4tURc6GbCiISGiZm2bQk=
|
||||
golang.org/x/sync v0.23.0/go.mod h1:sUUOizhqBxiL6pEWpqNLUiaJn1ShEbZ6BBqskPbjZm0=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
|
||||
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
|
||||
golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8=
|
||||
golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M=
|
||||
golang.org/x/text v0.42.0 h1:JbOZXgfeCPU9gacVtYliJqOhD+zhrEqK4LfdpmlUZqI=
|
||||
golang.org/x/text v0.42.0/go.mod h1:ojzP1Z+2QtioaF8DTtO8K5q7JWVVYwZKenzujK0Zd0E=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
|
||||
@@ -1,8 +1,14 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
@@ -24,21 +30,78 @@ type JetStream struct {
|
||||
|
||||
// Dial connects and returns the controller's JetStream handle.
|
||||
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))
|
||||
// **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...)
|
||||
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()
|
||||
if err != nil {
|
||||
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
|
||||
}
|
||||
|
||||
// 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,
|
||||
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
||||
// already handles (design 25 §3).
|
||||
|
||||
+29
-4
@@ -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
|
||||
// stream definitions (design 25 §3), so it alone reaches the JetStream 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
|
||||
// build is the one today: the controller asks, and reads the answer from the seat's event
|
||||
@@ -244,8 +248,15 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
case KindNode:
|
||||
// A host publishes its own node's control traffic and subscribes its own declaration —
|
||||
// and nothing of any other node's.
|
||||
pub = []string{"mesh.control." + p.Node + ".>"}
|
||||
sub = []string{"mesh.node." + p.Node + ".declare"}
|
||||
// And binding to its consumer, which asks the server about it (CONSUMER.INFO) — the one
|
||||
// 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:
|
||||
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
|
||||
@@ -278,7 +289,18 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
}
|
||||
|
||||
// 3. Seats it holds: full participation.
|
||||
// Its consumer's name, not ConsumerFor: that asks for these permissions to build the
|
||||
// consumer, and would ask forever. A subject for a consumer that turns out not to exist
|
||||
// grants nothing anybody can use.
|
||||
sub = append(sub, "_DELIVER."+consumerDurable(p))
|
||||
for _, s := range p.Holds {
|
||||
// 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 {
|
||||
sub = append(sub, seatSubject(s, "accept", a))
|
||||
}
|
||||
@@ -514,7 +536,10 @@ func ComposeAccounts(principals []Principal) (string, error) {
|
||||
// One account for the mesh: accounts in NATS isolate subject spaces entirely, and the mesh is
|
||||
// one 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.
|
||||
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 {
|
||||
perms, err := PermissionsFor(p)
|
||||
if err != nil {
|
||||
|
||||
@@ -35,10 +35,6 @@ type Readiness struct {
|
||||
Modules []string
|
||||
// ModuleCredentialled is which of those has one.
|
||||
ModuleCredentialled map[string]bool
|
||||
// StillOnTheOldBus is whether anything of the mesh's own still needs the bus it is leaving —
|
||||
// which is not a reason to stop, because that broker stays as an ordinary provider of `amqp`
|
||||
// (ADR 0119). Recorded so nobody reads the move as a retirement.
|
||||
OldBusHasOtherClients bool
|
||||
}
|
||||
|
||||
// NotReady is every reason this mesh cannot move its bus yet, in the order somebody would fix them.
|
||||
@@ -120,10 +116,11 @@ func WhatMoves(r Readiness) []string {
|
||||
out = append(out, fmt.Sprintf("move %d module runtime(s), and confirm each answers",
|
||||
len(r.Modules)))
|
||||
}
|
||||
if r.OldBusHasOtherClients {
|
||||
out = append(out, "leave the old broker running: it stays an ordinary provider of `amqp` for "+
|
||||
"whatever else uses it (ADR 0119), and this move is not its retirement")
|
||||
}
|
||||
// **The old broker goes, and it goes last** (novox/hq ADR 0131). AMQP is not a provision, so once
|
||||
// every machine reports on the new bus nothing of the mesh is left speaking to it, and its module
|
||||
// is unassigned. Said as a step so nobody reads the move as leaving a second bus behind.
|
||||
out = append(out, "then unassign the old broker's module: AMQP is not a provision (ADR 0131), and "+
|
||||
"once every machine reports on the new bus nothing of the mesh speaks to it")
|
||||
return out
|
||||
}
|
||||
|
||||
|
||||
@@ -79,9 +79,8 @@ func TestEachThingMissingNamesItsOwnRemedy(t *testing.T) {
|
||||
|
||||
// What the move would do is written out rather than summarised, because this is the one step with
|
||||
// nothing to inspect afterwards — so reading it is the last chance to disagree.
|
||||
func TestWhatMovesNamesEveryMachineAndSaysTheOldBrokerStays(t *testing.T) {
|
||||
func TestWhatMovesNamesEveryMachineAndEndsWithTheOldBrokerGoing(t *testing.T) {
|
||||
r := aMeshReadyToMove()
|
||||
r.OldBusHasOtherClients = true
|
||||
steps := strings.Join(WhatMoves(r), "\n")
|
||||
|
||||
for _, want := range []string{"anchor", "laptop", "user list", "module runtime"} {
|
||||
@@ -89,9 +88,11 @@ func TestWhatMovesNamesEveryMachineAndSaysTheOldBrokerStays(t *testing.T) {
|
||||
t.Errorf("the plan does not mention %q:\n%s", want, steps)
|
||||
}
|
||||
}
|
||||
// Said explicitly, so nobody reads the move as switching the old broker off — it stays serving
|
||||
// whatever else uses it, and that is a decision already taken.
|
||||
if !strings.Contains(steps, "not its retirement") {
|
||||
t.Errorf("the plan does not say the old broker stays:\n%s", steps)
|
||||
// Said explicitly, and last: AMQP is not a provision (novox/hq ADR 0131), so the move ends with
|
||||
// the old broker's module unassigned, not left behind as a second bus. An earlier version of this
|
||||
// test pinned the opposite, under a record 0131 superseded.
|
||||
lines := WhatMoves(r)
|
||||
if last := lines[len(lines)-1]; !strings.Contains(last, "unassign the old broker") {
|
||||
t.Errorf("the plan does not end with the old broker going:\n%s", steps)
|
||||
}
|
||||
}
|
||||
|
||||
+8
-7
@@ -21,10 +21,11 @@ jetstream {
|
||||
|
||||
accounts {
|
||||
MESH {
|
||||
jetstream: enabled
|
||||
users = [
|
||||
{ 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.>"] }
|
||||
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.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
||||
@@ -32,21 +33,21 @@ accounts {
|
||||
subscribe: { allow: ["_INBOX.enrol.one.>"] }
|
||||
} }
|
||||
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
||||
publish: { allow: ["$JS.ACK.NODES.one.>", "mesh.control.one.>"] }
|
||||
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
|
||||
subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||
} }
|
||||
{ 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"] }
|
||||
subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] }
|
||||
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
||||
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.one_telegram", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>"] }
|
||||
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
|
||||
subscribe: { allow: ["_DELIVER.two_audit", "_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
|
||||
} }
|
||||
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
|
||||
subscribe: { allow: ["_INBOX.two.shop.>"] }
|
||||
subscribe: { allow: ["_DELIVER.two_shop", "_INBOX.two.shop.>"] }
|
||||
} }
|
||||
]
|
||||
}
|
||||
|
||||
@@ -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
|
||||
// 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) {
|
||||
t.Errorf("the accounts file contains %q, which belongs to the module that raises the "+
|
||||
"server, not to the mesh", absent)
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
package catalogue
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// **The word does not come back through a manifest** (novox/hq ADR 0131). A module that wants
|
||||
// messaging wants the mesh's bus, reached through the sdk and named by the `mesh-broker` seat. Naming
|
||||
// the old wire protocol asks for the one server being retired, so both directions are refused at the
|
||||
// parser — this is judged from the manifest alone, no store needed.
|
||||
|
||||
func TestAManifestProvidingAmqpIsRefused(t *testing.T) {
|
||||
raw := []byte(`{"module":"old-broker","version":"1","provides":[{"name":"amqp","scope":"mesh"}]}`)
|
||||
_, err := ParseManifest(raw)
|
||||
if err == nil || !strings.Contains(err.Error(), `provides "amqp", which is not a provision`) {
|
||||
t.Fatalf("a module providing amqp was not refused, or not for the reason: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAManifestRequiringAmqpIsRefused(t *testing.T) {
|
||||
raw := []byte(`{"module":"forwarder","version":"1","requires":["amqp"]}`)
|
||||
_, err := ParseManifest(raw)
|
||||
if err == nil || !strings.Contains(err.Error(), `requires "amqp", which is not a provision`) {
|
||||
t.Fatalf("a module requiring amqp was not refused, or not for the reason: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// And the catalogue as checked out beside this repository names it nowhere — the three modules that
|
||||
// did are removed under design 28 task 5.4, not converted.
|
||||
func TestNoCatalogueManifestNamesAmqp(t *testing.T) {
|
||||
modules, err := filepath.Glob("../../../mesh-catalog/modules/*/module.json")
|
||||
if err != nil || len(modules) == 0 {
|
||||
t.Skip("the catalogue is not checked out beside this repository")
|
||||
}
|
||||
for _, path := range modules {
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if strings.Contains(string(raw), `"amqp"`) {
|
||||
t.Errorf("%s names amqp, which is not a provision (novox/hq ADR 0131)",
|
||||
filepath.Base(filepath.Dir(path)))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -46,23 +46,9 @@ func TestTheSeatRefusesADifferentBusToo(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The old broker no longer claims the seat: it is an ordinary provider of `amqp`
|
||||
// (novox/hq ADR 0119), so it can sit on the same mesh as the bus without contending for it.
|
||||
func TestTheAmqpBrokerDoesNotContendForTheSeat(t *testing.T) {
|
||||
lavinmq := catalogueManifest(t, "lavinmq")
|
||||
for _, c := range lavinmq.Claims {
|
||||
if c.Name == "mesh-broker" {
|
||||
t.Fatal("the amqp broker still claims mesh-broker; it is a provider, not foundation")
|
||||
}
|
||||
}
|
||||
busHeld := World{Held: []Held{{Claim: "mesh-broker", Scope: ScopeMesh,
|
||||
Node: "anchor", Module: "nats"}}}
|
||||
other := workstation()
|
||||
other.Name = "laptop"
|
||||
if _, err := Resolve(shelf(lavinmq), []string{"lavinmq"}, other, busHeld); err != nil {
|
||||
t.Fatalf("the amqp broker was refused beside the mesh bus: %v", err)
|
||||
}
|
||||
}
|
||||
// The old broker is gone from the catalogue (novox/hq ADR 0131, design 28 task 5.4), so it is no
|
||||
// longer a fixture here. That two eligible holders stand beside each other with one on record is
|
||||
// pinned in holdings_test.go against manifests this package owns.
|
||||
|
||||
// **A seat and the interface it delivers are different names, and renaming one must not rename
|
||||
// the other** (novox/hq ADR 0118). This nearly went wrong: the seats were renamed to the `mesh-*`
|
||||
|
||||
@@ -103,6 +103,11 @@ type Rendering struct {
|
||||
// **Only the users, never the server's own settings**: those are the module's, in its image and
|
||||
// its mounts (Manifest.BusUsers).
|
||||
BusUsers string
|
||||
// BusMembership is this machine's membership for the bus the mesh is moving to, sealed to it
|
||||
// (design 28, task 5.2). Empty for a machine not being moved. Written as a file the host reads
|
||||
// after the declaration has applied, so the bus it names is standing before the machine leaves
|
||||
// the one it is on.
|
||||
BusMembership string
|
||||
|
||||
// MeshRange is the private network's CIDR (the range node addresses are allocated from), for a
|
||||
// module that must name the whole mesh rather than one machine — an intrusion filter that must
|
||||
@@ -238,9 +243,23 @@ func (r Resolution) Compose(with Rendering) (Composed, error) {
|
||||
if err != nil {
|
||||
return Composed{}, err
|
||||
}
|
||||
if with.BusMembership != "" {
|
||||
// The machine's own, not any module's: how it reaches the mesh from now on. Sealed like a
|
||||
// secret and placed where the host looks for exactly this (design 28, task 5.2).
|
||||
resources = append(resources, map[string]any{
|
||||
"id": BusMembershipID(), "type": "file", "path": BusMembershipPath,
|
||||
"sealed": with.BusMembership, "mode": "0600",
|
||||
})
|
||||
}
|
||||
return Composed{Resources: resources, Owner: owner}, nil
|
||||
}
|
||||
|
||||
// BusMembershipID names the resource carrying a machine's membership for the new bus, and
|
||||
// BusMembershipPath is where the host reads it — the same constant on both sides.
|
||||
func BusMembershipID() string { return "bus-membership" }
|
||||
|
||||
const BusMembershipPath = "/var/lib/mesh/membership-next.json"
|
||||
|
||||
func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) {
|
||||
// Every manifest is placed first (novox/hq ADR 0112): the maps naming where its bindings,
|
||||
// credentials and contributions land are resolved against this node's directories, so every
|
||||
|
||||
@@ -28,8 +28,10 @@ func TestTheStoreAndTheBrokerSayWhatTheMeshGuards(t *testing.T) {
|
||||
if got := catalogueManifest(t, "postgres").Guards; !reflect.DeepEqual(got, []int{5432}) {
|
||||
t.Errorf("postgres guards %v; the store's port must be refused from outside", got)
|
||||
}
|
||||
if got := catalogueManifest(t, "lavinmq").Guards; !reflect.DeepEqual(got, []int{15672}) {
|
||||
t.Errorf("lavinmq guards %v; the management port must be refused from outside", got)
|
||||
// The bus's monitoring port, not its client port: a node reaches the bus, nobody outside
|
||||
// reads its state (novox/hq ADR 0131 — the broker that guarded 15672 has left the catalogue).
|
||||
if got := catalogueManifest(t, "nats").Guards; !reflect.DeepEqual(got, []int{8222}) {
|
||||
t.Errorf("nats guards %v; the monitoring port must be refused from outside", got)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
package catalogue
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// **A seat's holder on record settles who holds it, and lets the next holder stand beside the
|
||||
// current one** (novox/hq ADR 0131, design 28 task 5.3). Until the record existed, two assignments
|
||||
// whose modules both claimed a seat were refused outright — which left no way to hand a seat over
|
||||
// without a moment where nobody held it, and the control plane finds its own bus through one of
|
||||
// these seats. That moment was the outage of 2026-09-27.
|
||||
|
||||
func busSeatDelivering(t *testing.T, delivers string) {
|
||||
t.Helper()
|
||||
was := Seats()
|
||||
t.Cleanup(func() { UseSeats(was) })
|
||||
UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: delivers, Decision: "test"}})
|
||||
}
|
||||
|
||||
func oldBroker() Manifest {
|
||||
return Manifest{Module: "old-broker", Provides: []Offer{{Name: "mesh-bus", Scope: ScopeMesh}},
|
||||
Claims: []Claim{{Name: "mesh-broker", Scope: ScopeMesh}}}
|
||||
}
|
||||
|
||||
func newBroker() Manifest {
|
||||
return Manifest{Module: "new-broker", Provides: []Offer{{Name: "mesh-bus", Scope: ScopeMesh}},
|
||||
Claims: []Claim{{Name: "mesh-broker", Scope: ScopeMesh}}}
|
||||
}
|
||||
|
||||
// Nothing on record: exactly the old rule. One claimant holds; two are refused.
|
||||
func TestWithNoHolderOnRecordTheSoleClaimantHoldsAndTwoAreRefused(t *testing.T) {
|
||||
busSeatDelivering(t, "mesh-bus")
|
||||
node := Node{Name: "anchor"}
|
||||
|
||||
held, problems := checkClaims([]Manifest{oldBroker()}, node, nil, nil)
|
||||
if len(problems) != 0 || len(held) != 1 || held[0].Module != "old-broker" {
|
||||
t.Fatalf("a sole claimant did not hold the seat: held=%v problems=%v", held, problems)
|
||||
}
|
||||
_, problems = checkClaims([]Manifest{oldBroker(), newBroker()}, node, nil, nil)
|
||||
if len(problems) != 1 || !strings.Contains(problems[0], "both claim") {
|
||||
t.Fatalf("two claimants with nothing on record were not refused: %v", problems)
|
||||
}
|
||||
}
|
||||
|
||||
// With a holder on record, the other eligible assignment is silent: not refused, and not holding.
|
||||
func TestTheHolderOnRecordHoldsAndTheOtherClaimantStandsBesideIt(t *testing.T) {
|
||||
busSeatDelivering(t, "mesh-bus")
|
||||
node := Node{Name: "anchor"}
|
||||
record := []Held{{Claim: "mesh-broker", Scope: ScopeMesh, Node: "anchor", Module: "new-broker"}}
|
||||
|
||||
held, problems := checkClaims([]Manifest{oldBroker(), newBroker()}, node, nil, record)
|
||||
if len(problems) != 0 {
|
||||
t.Fatalf("the assignment beside the holder was refused: %v", problems)
|
||||
}
|
||||
if len(held) != 1 || held[0].Module != "new-broker" {
|
||||
t.Fatalf("the holder on record is not the one holding: %v", held)
|
||||
}
|
||||
}
|
||||
|
||||
// The record names a node too: an eligible module on another machine holds nothing, and its
|
||||
// machine's set still resolves.
|
||||
func TestAHolderOnRecordElsewhereLeavesThisMachinesClaimantSilent(t *testing.T) {
|
||||
busSeatDelivering(t, "mesh-bus")
|
||||
record := []Held{{Claim: "mesh-broker", Scope: ScopeMesh, Node: "anchor", Module: "new-broker"}}
|
||||
|
||||
held, problems := checkClaims([]Manifest{oldBroker()}, Node{Name: "laptop"}, nil, record)
|
||||
if len(problems) != 0 || len(held) != 0 {
|
||||
t.Fatalf("a claimant elsewhere than the recorded holder was not simply silent: held=%v problems=%v",
|
||||
held, problems)
|
||||
}
|
||||
}
|
||||
|
||||
// A record naming a seat's former name still applies to it after a rename (ADR 0122).
|
||||
func TestAHolderRecordedUnderAFormerNameStillHolds(t *testing.T) {
|
||||
busSeatDelivering(t, "mesh-bus")
|
||||
wasAliases := aliases
|
||||
t.Cleanup(func() { UseAliases(wasAliases) })
|
||||
UseAliases(map[string]string{"the-broker": "mesh-broker"})
|
||||
record := []Held{{Claim: "the-broker", Scope: ScopeMesh, Node: "anchor", Module: "new-broker"}}
|
||||
|
||||
held, problems := checkClaims([]Manifest{oldBroker(), newBroker()}, Node{Name: "anchor"}, nil, record)
|
||||
if len(problems) != 0 || len(held) != 1 || held[0].Module != "new-broker" {
|
||||
t.Fatalf("a record under the former name did not settle the seat: held=%v problems=%v", held, problems)
|
||||
}
|
||||
}
|
||||
|
||||
// CanHold is the one judgement registration and the handover share, against the store's row.
|
||||
func TestCanHoldJudgesClaimScopeAndWhatTheSeatDelivers(t *testing.T) {
|
||||
busSeatDelivering(t, "mesh-bus")
|
||||
seat, _ := SeatNamed("mesh-broker")
|
||||
|
||||
if err := CanHold(newBroker(), seat); err != nil {
|
||||
t.Fatalf("a module that claims the seat and provides what it delivers was refused: %v", err)
|
||||
}
|
||||
noClaim := Manifest{Module: "quiet", Provides: []Offer{{Name: "mesh-bus", Scope: ScopeMesh}}}
|
||||
if err := CanHold(noClaim, seat); err == nil || !strings.Contains(err.Error(), "does not claim") {
|
||||
t.Fatalf("a module that never claimed the seat was allowed to hold it: %v", err)
|
||||
}
|
||||
wrongScope := newBroker()
|
||||
wrongScope.Claims[0].Scope = ScopeNode
|
||||
if err := CanHold(wrongScope, seat); err == nil || !strings.Contains(err.Error(), "scope") {
|
||||
t.Fatalf("a claim at the wrong scope was allowed: %v", err)
|
||||
}
|
||||
cannotAnswer := Manifest{Module: "amqp-only", Provides: []Offer{{Name: "amqp", Scope: ScopeMesh}},
|
||||
Claims: []Claim{{Name: "mesh-broker", Scope: ScopeMesh}}}
|
||||
if err := CanHold(cannotAnswer, seat); err == nil || !strings.Contains(err.Error(), `does not provide "mesh-bus"`) {
|
||||
t.Fatalf("a holder that cannot answer for the seat was allowed: %v", err)
|
||||
}
|
||||
// And the judgement follows the store's row, not a compiled copy.
|
||||
busSeatDelivering(t, "amqp")
|
||||
seat, _ = SeatNamed("mesh-broker")
|
||||
if err := CanHold(cannotAnswer, seat); err != nil {
|
||||
t.Fatalf("with the row saying amqp, an amqp provider was refused: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A machine being moved is handed its membership for the new bus as a sealed file in its own
|
||||
// declaration — the machine's, not any module's (design 28, task 5.2).
|
||||
func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) {
|
||||
r := Resolution{Node: "anchor"}
|
||||
got, err := r.Compose(Rendering{BusMembership: "sealed-blob"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var found map[string]any
|
||||
for _, res := range got.Resources {
|
||||
if res["id"] == BusMembershipID() {
|
||||
found = res
|
||||
}
|
||||
}
|
||||
if found == nil {
|
||||
t.Fatalf("no membership resource in %v", got.Resources)
|
||||
}
|
||||
if found["path"] != BusMembershipPath || found["sealed"] != "sealed-blob" || found["mode"] != "0600" {
|
||||
t.Fatalf("the membership is not a sealed 0600 file where the host reads it: %v", found)
|
||||
}
|
||||
// And a machine not being moved is handed nothing.
|
||||
got, _ = r.Compose(Rendering{})
|
||||
for _, res := range got.Resources {
|
||||
if res["id"] == BusMembershipID() {
|
||||
t.Fatal("a machine with no membership on record was handed one")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The store's seat rows have no protocol columns yet; loading them must not drop the protocol the
|
||||
// bus is derived from, or no role's work queue is ever raised (found live, 2026-09-28).
|
||||
func TestAStoreRowWithoutAProtocolKeepsTheCompiledOne(t *testing.T) {
|
||||
was := Seats()
|
||||
t.Cleanup(func() { UseSeats(was) })
|
||||
UseSeats([]Seat{{Name: "mesh-build-machine", Scope: ScopeMesh, Decision: "row"}})
|
||||
got, ok := SeatNamed("mesh-build-machine")
|
||||
if !ok || len(got.Accepts) == 0 {
|
||||
t.Fatalf("the build machine's seat lost what it accepts when loaded from the store: %+v", got)
|
||||
}
|
||||
if got.Decision != "row" {
|
||||
t.Fatalf("the store's own columns were not kept: %+v", got)
|
||||
}
|
||||
}
|
||||
@@ -1027,6 +1027,28 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%q is not a usable slug: lower-case letters, digits, dashes and dots", m.Slug))
|
||||
}
|
||||
// **`amqp` is not a provision, and not a requirement** (novox/hq ADR 0131). A module that wants
|
||||
// messaging wants the mesh's bus — it emits and consumes through the sdk, which the mesh hands the
|
||||
// bus with the module's own credential — and the bus is whatever holds `mesh-broker`, spoken in
|
||||
// whatever that holder speaks. Naming the old wire protocol asks for a specific server, and the
|
||||
// only one that could answer is the one being retired. Refused here so the word cannot come back
|
||||
// through a manifest.
|
||||
for _, offer := range m.Provides {
|
||||
if offer.Name == "amqp" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s provides %q, which is not a provision: the mesh's bus is whatever holds "+
|
||||
"mesh-broker, and a module provides mesh-bus to be it (novox/hq ADR 0131)",
|
||||
m.Module, offer.Name))
|
||||
}
|
||||
}
|
||||
for _, r := range m.Requires {
|
||||
if r == "amqp" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s requires %q, which is not a provision: a module reaches the mesh's bus through "+
|
||||
"the sdk, and depends on the mesh-broker seat, not on a protocol (novox/hq ADR 0131)",
|
||||
m.Module, r))
|
||||
}
|
||||
}
|
||||
for _, offer := range m.Provides {
|
||||
p := offer.Name
|
||||
if !name.MatchString(p) {
|
||||
|
||||
@@ -41,6 +41,11 @@ type Node struct {
|
||||
type World struct {
|
||||
// Held is the claims already taken, for the scopes wider than one node.
|
||||
Held []Held
|
||||
// Holdings is every seat whose holder is **on record** (novox/hq ADR 0131): the one assignment
|
||||
// that holds it, chosen by a handover. A seat absent here is held by derivation — the sole
|
||||
// eligible assignment — as it always was. Present, it decides, and any other assignment whose
|
||||
// module could hold the seat is eligible and silent rather than refused.
|
||||
Holdings []Held
|
||||
// Offered is what other nodes provide at mesh scope, and everything needed to use it.
|
||||
Offered map[string][]Provider
|
||||
// Pinned is which node this machine was told to get a provision from, by name. Only consulted
|
||||
@@ -557,7 +562,7 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
|
||||
}
|
||||
|
||||
problems = append(problems, checkCapabilities(resolution.Modules, node)...)
|
||||
claims, claimProblems := checkClaims(resolution.Modules, node, elsewhere)
|
||||
claims, claimProblems := checkClaims(resolution.Modules, node, elsewhere, world.Holdings)
|
||||
problems = append(problems, claimProblems...)
|
||||
problems = append(problems, checkResources(resolution.Modules)...)
|
||||
resolution.Claims = claims
|
||||
@@ -650,14 +655,35 @@ func checkCapabilities(modules []Manifest, node Node) []string {
|
||||
//
|
||||
// Within this node's own set, and against what is already held elsewhere for the wider scopes. A
|
||||
// claim at mesh scope is the same idea as the mesh's one hub, said once instead of hard-coded.
|
||||
func checkClaims(modules []Manifest, node Node, elsewhere []Held) ([]Held, []string) {
|
||||
func checkClaims(modules []Manifest, node Node, elsewhere []Held, holdings []Held) ([]Held, []string) {
|
||||
var problems []string
|
||||
var held []Held
|
||||
|
||||
// onRecord is the recorded holder of a seat, if a handover ever named one.
|
||||
onRecord := func(claim, scope string) (Held, bool) {
|
||||
for _, h := range holdings {
|
||||
hs, ok := SeatNamed(h.Claim)
|
||||
cs, cok := SeatNamed(claim)
|
||||
if ok && cok && hs.Name == cs.Name && h.Scope == scope {
|
||||
return h, true
|
||||
}
|
||||
}
|
||||
return Held{}, false
|
||||
}
|
||||
|
||||
byScope := map[string]map[string]string{} // scope → claim → module
|
||||
for _, m := range modules {
|
||||
for _, c := range m.Claims {
|
||||
scope := c.At()
|
||||
// **A recorded holder settles it before any counting.** An assignment that could hold
|
||||
// the seat but is not the one on record is eligible, and that is all: it is not a second
|
||||
// holder, so it is not refused, and it does not hold (novox/hq ADR 0131). This is what
|
||||
// lets the next holder stand beside the current one until the seat is handed over.
|
||||
if rec, recorded := onRecord(c.Name, scope); recorded {
|
||||
if rec.Node != node.Name || rec.Module != m.Module {
|
||||
continue
|
||||
}
|
||||
}
|
||||
if byScope[scope] == nil {
|
||||
byScope[scope] = map[string]string{}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
|
||||
"type": "container", "id": "server", "name": "mesh-controller",
|
||||
"env": map[string]any{
|
||||
"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_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||
},
|
||||
@@ -27,7 +26,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
|
||||
env := control["env"].(map[string]any)
|
||||
for key, want := range map[string]string{
|
||||
"MESH_STORE_INVENTORY_PORT": "6852",
|
||||
"MESH_BROKER_AMQP_PORT": "5679",
|
||||
"MESH_BROKER_ADDRESS_PORT": "5671",
|
||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||
} {
|
||||
@@ -146,7 +144,6 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
|
||||
"MESH_STORE_INVENTORY_PORT": "6852",
|
||||
"MESH_STORE_IDENTITY_PORT": "6852",
|
||||
"MESH_STORE_LICENCES_PORT": "6852",
|
||||
"MESH_BROKER_AMQP_PORT": "5679",
|
||||
"MESH_BROKER_MANAGEMENT_PORT": "15673",
|
||||
"MESH_BROKER_ADDRESS_PORT": "5671",
|
||||
} {
|
||||
@@ -164,7 +161,7 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -175,7 +172,6 @@ var SeatPorts = map[string]string{
|
||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_STORE_IDENTITY_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_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
||||
}
|
||||
|
||||
+59
-13
@@ -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
|
||||
// a non-empty one, never erase it.
|
||||
func UseSeats(s []Seat) {
|
||||
if len(s) > 0 {
|
||||
seats = s
|
||||
if len(s) == 0 {
|
||||
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
|
||||
@@ -184,17 +205,17 @@ func claimProblems(m Manifest) []string {
|
||||
}
|
||||
|
||||
for _, c := range m.Claims {
|
||||
if seat, known := SeatNamed(c.Name); known {
|
||||
if c.At() != seat.Scope {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s claims %s at scope %q, and %s is a %s seat",
|
||||
m.Module, c.Name, c.At(), c.Name, seat.Scope))
|
||||
}
|
||||
if seat.Delivers != "" && !providesAt(m, seat.Delivers, seat.Scope) {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s claims %s, whose holder answers for %q, and %s does not provide %q at %s scope",
|
||||
m.Module, c.Name, seat.Delivers, m.Module, seat.Delivers, seat.Scope))
|
||||
}
|
||||
if _, known := SeatNamed(c.Name); known {
|
||||
// **A seat's scope and what it delivers are not judged here** (novox/hq ADR 0122).
|
||||
// This function runs wherever a manifest is parsed, and one of those places is the
|
||||
// build machine, which has no store: there, `SeatNamed` answers from the set the
|
||||
// binary shipped with, so a build would be refused for disagreeing with a compiled
|
||||
// copy of data the control plane owns. Exactly that happened — a holder of the bus
|
||||
// seat was refused for not providing what a stale compiled row said the seat
|
||||
// delivered, while the store's own row said otherwise.
|
||||
//
|
||||
// Both checks moved to CatalogueProblems, which only ever runs in the control plane,
|
||||
// after UseSeats has replaced the set with the store's.
|
||||
continue
|
||||
}
|
||||
if isSystemSeatName(c.Name) {
|
||||
@@ -222,6 +243,31 @@ func claimProblems(m Manifest) []string {
|
||||
return problems
|
||||
}
|
||||
|
||||
// CanHold is why a module could not hold a seat, or nothing: its definition must claim the seat at
|
||||
// the seat's scope, and provide what the seat delivers, if it delivers anything. The seat is the
|
||||
// store's row, so this is judged only where the store's set is loaded — at registration and in the
|
||||
// handover command (novox/hq ADR 0131), never in the parser.
|
||||
func CanHold(m Manifest, seat Seat) error {
|
||||
var claimed *Claim
|
||||
for i := range m.Claims {
|
||||
if hs, ok := SeatNamed(m.Claims[i].Name); ok && hs.Name == seat.Name {
|
||||
claimed = &m.Claims[i]
|
||||
}
|
||||
}
|
||||
if claimed == nil {
|
||||
return fmt.Errorf("%s does not claim %s", m.Module, seat.Name)
|
||||
}
|
||||
if claimed.At() != seat.Scope {
|
||||
return fmt.Errorf("%s claims %s at scope %q, and %s is a %s seat",
|
||||
m.Module, seat.Name, claimed.At(), seat.Name, seat.Scope)
|
||||
}
|
||||
if seat.Delivers != "" && !providesAt(m, seat.Delivers, seat.Scope) {
|
||||
return fmt.Errorf("%s claims %s, whose holder answers for %q, and %s does not provide %q at %s scope",
|
||||
m.Module, seat.Name, seat.Delivers, m.Module, seat.Delivers, seat.Scope)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func providesAt(m Manifest, provision, scope string) bool {
|
||||
for _, o := range m.Provides {
|
||||
if o.Name == provision && o.At() == scope {
|
||||
|
||||
@@ -190,7 +190,15 @@ func CatalogueProblems(shelf Shelf) []string {
|
||||
}
|
||||
s, isModuleSeat := declared[c.Name]
|
||||
if !isModuleSeat {
|
||||
continue // a mesh seat: already judged by claimProblems
|
||||
// **A mesh seat is judged here and nowhere else** (novox/hq ADR 0122): the set is
|
||||
// the store's, and this is the only place that runs with the store's set loaded.
|
||||
// The parser cannot do it — it also runs on the build machine, against whatever
|
||||
// set that binary was compiled with.
|
||||
seat, _ := SeatNamed(c.Name)
|
||||
if err := CanHold(m, seat); err != nil {
|
||||
problems = append(problems, err.Error())
|
||||
}
|
||||
continue
|
||||
}
|
||||
if c.At() != s.At() {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
|
||||
@@ -91,7 +91,10 @@ func TestAClaimMustMatchTheDeclaredScope(t *testing.T) {
|
||||
|
||||
// The mesh's own seats still work, and are not shadowed by the derived half.
|
||||
func TestTheMeshsOwnSeatsAreStillClaimable(t *testing.T) {
|
||||
m := Manifest{Module: "nats", Claims: []Claim{{Name: "mesh-broker", Scope: ScopeMesh}}}
|
||||
// It delivers the bus, so its holder provides the bus — the rule this check now enforces.
|
||||
m := Manifest{Module: "nats",
|
||||
Provides: []Offer{{Name: "mesh-bus", Scope: ScopeMesh}},
|
||||
Claims: []Claim{{Name: "mesh-broker", Scope: ScopeMesh}}}
|
||||
if got := problemsFor(t, Shelf{"nats": m}); got != "" {
|
||||
t.Fatalf("a mesh seat was refused by the derived check: %s", got)
|
||||
}
|
||||
@@ -106,3 +109,38 @@ func TestTheProblemsAreStable(t *testing.T) {
|
||||
t.Fatalf("unstable:\n%s\n%s", first, second)
|
||||
}
|
||||
}
|
||||
|
||||
// **A build machine has no store, so it may not judge a seat.** The set is data the control plane
|
||||
// owns (novox/hq ADR 0122), and `ParseManifest` runs on the build machine too, against whatever set
|
||||
// that binary was compiled with. When the two disagreed, a valid holder of the bus seat was refused
|
||||
// mid-rollout — the compiled row said the seat delivered one provision, the store's row said
|
||||
// another, and the build failed on the copy rather than the truth. The parser judges the manifest;
|
||||
// the seat set judges the claim, where it is loaded.
|
||||
func TestTheParserDoesNotJudgeWhatOnlyTheStoreKnows(t *testing.T) {
|
||||
was := Seats()
|
||||
t.Cleanup(func() { UseSeats(was) })
|
||||
|
||||
// A store whose bus seat delivers something this module does provide.
|
||||
UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "test"}})
|
||||
raw := []byte(`{"module":"a-bus","version":"1",` +
|
||||
`"provides":[{"name":"mesh-bus","scope":"mesh"}],` +
|
||||
`"claims":[{"name":"mesh-broker","scope":"mesh"}]}`)
|
||||
|
||||
m, err := ParseManifest(raw)
|
||||
if err != nil {
|
||||
t.Fatalf("the parser refused a claim only the seat set can judge: %v", err)
|
||||
}
|
||||
if got := CatalogueProblems(Shelf{m.Module: m}); len(got) != 0 {
|
||||
t.Fatalf("a holder that provides what the store says the seat delivers was refused: %v", got)
|
||||
}
|
||||
|
||||
// And with the store saying the seat delivers something else, registration is what refuses it.
|
||||
UseSeats([]Seat{{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "other-bus", Decision: "test"}})
|
||||
if _, err := ParseManifest(raw); err != nil {
|
||||
t.Fatalf("the parser judged it the second time: %v", err)
|
||||
}
|
||||
got := strings.Join(CatalogueProblems(Shelf{m.Module: m}), "; ")
|
||||
if !strings.Contains(got, `does not provide "other-bus"`) {
|
||||
t.Fatalf("registration did not refuse a holder that cannot answer for the seat: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -138,26 +138,32 @@ func TestAModuleDefinesAndClaimsItsOwnSeat(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// **At registration, not in the parser** (novox/hq ADR 0122): a seat's scope is a property of the
|
||||
// set, the set is the store's, and the parser also runs on a build machine that has no store.
|
||||
func TestASeatClaimedAtAnotherScopeIsRefused(t *testing.T) {
|
||||
_, err := ParseManifest(claimed(`[{"name":"npm-package-registry","scope":"node"}]`))
|
||||
if err == nil {
|
||||
t.Fatal("a mesh seat was held per node")
|
||||
m, err := ParseManifest(claimed(`[{"name":"npm-package-registry","scope":"node"}]`))
|
||||
if err != nil {
|
||||
t.Fatalf("the parser judged a scope it reads from data it may not have: %v", err)
|
||||
}
|
||||
if !strings.Contains(err.Error(), "mesh seat") {
|
||||
t.Fatalf("the refusal does not say which scope the seat is: %v", err)
|
||||
got := strings.Join(CatalogueProblems(Shelf{m.Module: m}), "; ")
|
||||
if !strings.Contains(got, "mesh seat") {
|
||||
t.Fatalf("the refusal does not say which scope the seat is: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestADeliveringSeatIsOnlyHeldByAModuleThatProvides(t *testing.T) {
|
||||
// Holding it makes the module the mesh's answer for the provision. A module that cannot answer
|
||||
// would be the answer anyway, and every consumer would be sent to it.
|
||||
// And refused at registration, where the seat set is the store's: what a seat delivers is
|
||||
// data, so a compiled copy of it may not be what refuses a build (novox/hq ADR 0122).
|
||||
raw := []byte(`{"module":"thing","version":"1","claims":[{"name":"git","scope":"mesh"}]}`)
|
||||
_, err := ParseManifest(raw)
|
||||
if err == nil {
|
||||
t.Fatal("a module holding the git seat need not provide git")
|
||||
m, err := ParseManifest(raw)
|
||||
if err != nil {
|
||||
t.Fatalf("the parser judged what a seat delivers: %v", err)
|
||||
}
|
||||
if !strings.Contains(err.Error(), `does not provide "git"`) {
|
||||
t.Fatalf("the refusal does not say what is missing: %v", err)
|
||||
got := strings.Join(CatalogueProblems(Shelf{m.Module: m}), "; ")
|
||||
if !strings.Contains(got, `does not provide "git"`) {
|
||||
t.Fatalf("the refusal does not say what is missing: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -230,3 +230,35 @@ func (i *Inventory) ForgetPerson(ctx context.Context, name string) error {
|
||||
}
|
||||
return i.ForgetBusUser(ctx, "person."+name)
|
||||
}
|
||||
|
||||
// PutBusMembership records a machine's membership for the new bus, sealed to it (design 28, 5.2).
|
||||
// Replaces any earlier one: a machine has one membership per bus, and re-minting is re-telling.
|
||||
func (i *Inventory) PutBusMembership(ctx context.Context, nodeName, sealed string) error {
|
||||
node, err := i.NodeByName(ctx, nodeName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = i.store.Pool().Exec(ctx,
|
||||
`insert into bus_membership (node, sealed) values ($1, $2)
|
||||
on conflict (node) do update set sealed = excluded.sealed, since = now()`, node.ID, sealed)
|
||||
return err
|
||||
}
|
||||
|
||||
// BusMemberships is every machine's sealed membership for the new bus, by node name.
|
||||
func (i *Inventory) BusMemberships(ctx context.Context) (map[string]string, error) {
|
||||
rows, err := i.store.Pool().Query(ctx,
|
||||
`select n.name, b.sealed from bus_membership b join node n on n.id = b.node`)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
out := map[string]string{}
|
||||
for rows.Next() {
|
||||
var name, sealed string
|
||||
if err := rows.Scan(&name, &sealed); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out[name] = sealed
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package inventory
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
)
|
||||
|
||||
// A seat's holder is a row, changed as one act (novox/hq ADR 0131). What these pin is the shape of
|
||||
// that row's life: it needs an assignment to point at, it is replaced rather than added to, and it
|
||||
// goes when the assignment does — so a seat never points at something that is not running anywhere.
|
||||
|
||||
func twoBrokersOnTwoNodes(t *testing.T) (*Inventory, context.Context) {
|
||||
t.Helper()
|
||||
old := catalogue.Manifest{Module: "old-broker", Version: "1",
|
||||
Provides: []catalogue.Offer{{Name: "mesh-bus", Scope: catalogue.ScopeMesh}},
|
||||
Claims: []catalogue.Claim{{Name: "mesh-broker", Scope: catalogue.ScopeMesh}}}
|
||||
new := catalogue.Manifest{Module: "new-broker", Version: "1",
|
||||
Provides: []catalogue.Offer{{Name: "mesh-bus", Scope: catalogue.ScopeMesh}},
|
||||
Claims: []catalogue.Claim{{Name: "mesh-broker", Scope: catalogue.ScopeMesh}}}
|
||||
inv, ctx := aMeshWith(t, old, new)
|
||||
// The holding references the seat's row, which `migrate` seeds on a real mesh.
|
||||
if _, err := inv.SeedSeats(ctx, catalogue.DefaultSeats()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, n := range []string{"anchor", "laptop"} {
|
||||
if _, err := inv.AddNode(ctx, n); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, err := inv.Assign(ctx, "anchor", "old-broker"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.Assign(ctx, "laptop", "new-broker"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return inv, ctx
|
||||
}
|
||||
|
||||
func TestAHandoverIsOneRowReplacedNotOneAdded(t *testing.T) {
|
||||
inv, ctx := twoBrokersOnTwoNodes(t)
|
||||
|
||||
if err := inv.HoldSeat(ctx, "mesh-broker", catalogue.ScopeMesh, "anchor", "old-broker"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.HoldSeat(ctx, "mesh-broker", catalogue.ScopeMesh, "laptop", "new-broker"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
held, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(held) != 1 || held[0].Node != "laptop" || held[0].Module != "new-broker" || held[0].Scope != catalogue.ScopeMesh {
|
||||
t.Fatalf("after a handover the seat is not held by exactly the new holder: %+v", held)
|
||||
}
|
||||
}
|
||||
|
||||
func TestASeatCannotBeHandedToSomethingNotAssigned(t *testing.T) {
|
||||
inv, ctx := twoBrokersOnTwoNodes(t)
|
||||
// new-broker is assigned to laptop, not anchor.
|
||||
if err := inv.HoldSeat(ctx, "mesh-broker", catalogue.ScopeMesh, "anchor", "new-broker"); err == nil {
|
||||
t.Fatal("a seat was handed to a module not assigned where it was named")
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnassigningTheHolderTakesTheHoldingWithIt(t *testing.T) {
|
||||
inv, ctx := twoBrokersOnTwoNodes(t)
|
||||
if err := inv.HoldSeat(ctx, "mesh-broker", catalogue.ScopeMesh, "laptop", "new-broker"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.Unassign(ctx, "laptop", "new-broker"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
held, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(held) != 0 {
|
||||
t.Fatalf("the holding outlived the assignment it pointed at: %+v", held)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAMachinesMembershipIsOneRowReplacedAndGoesWithTheMachine(t *testing.T) {
|
||||
inv, ctx := twoBrokersOnTwoNodes(t)
|
||||
if err := inv.PutBusMembership(ctx, "anchor", "first"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.PutBusMembership(ctx, "anchor", "second"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, err := inv.BusMemberships(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got["anchor"] != "second" || len(got) != 1 {
|
||||
t.Fatalf("a re-told membership did not replace the first: %v", got)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
-- A seat's holder is a recorded fact, not a derivation (novox/hq ADR 0131, design 26).
|
||||
--
|
||||
-- Until this, "which assignment holds the seat" was derived: the module that is assigned and
|
||||
-- claims the seat holds it, and a second eligible assignment was refused at resolution. That has no
|
||||
-- way to hand a seat from one holder to the next without a moment where nothing holds it — and the
|
||||
-- control plane finds its own bus through one of these seats, so that moment was an outage
|
||||
-- (2026-09-27). Now the holder is one row here, changed by `seat <name> --to <node>/<module>` as one
|
||||
-- act, and other assignments whose module could hold the seat are simply eligible and silent.
|
||||
--
|
||||
-- No row means what it always meant: the sole eligible assignment holds the seat, and two eligible
|
||||
-- ones are refused. So a mesh that has never handed a seat over behaves exactly as before, and the
|
||||
-- row appears the first time somebody does.
|
||||
--
|
||||
-- The seat is referenced by name because claims still are (0034); the rename cascades here so a
|
||||
-- handed-over seat survives being renamed. The holder is the assignment itself, so unassigning it
|
||||
-- takes the holding with it and the seat falls back to derivation rather than pointing at nothing.
|
||||
create table seat_holding (
|
||||
seat text primary key references seat(name) on update cascade on delete cascade,
|
||||
scope text not null,
|
||||
node uuid not null,
|
||||
module text not null,
|
||||
since timestamptz not null default now(),
|
||||
foreign key (node, module) references assignment(node, module) on delete cascade
|
||||
);
|
||||
@@ -0,0 +1,12 @@
|
||||
-- The bus seat's holder answers for the mesh's bus, not for a wire protocol (novox/hq ADR 0131).
|
||||
--
|
||||
-- The row said `amqp`, which is the protocol the old broker spoke, and so only that broker could hold
|
||||
-- the seat that names the mesh's bus — while the module that will carry the bus could not. The seat
|
||||
-- delivers `mesh-bus`; whichever module provides that may hold it, and today that is one module.
|
||||
--
|
||||
-- Safe under the current holder: the control plane composes its own bus address through the seat by
|
||||
-- name, and the overview derives holders by name. Only registration and provision-to-seat resolution
|
||||
-- read this column. So the row changes, the current holder keeps holding by derivation, the next one
|
||||
-- can register its claim, and the handover (0039) moves the seat when both are running. What must
|
||||
-- not happen in between is re-registering the current holder — registration would now refuse it.
|
||||
update seat set delivers = 'mesh-bus' where name = 'mesh-broker' and delivers = 'amqp';
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
-- A machine already enrolled is moved to the new bus by being told its membership for it
|
||||
-- (novox/hq design 28, task 5.2). Until this, a membership — bus address, fingerprint, password,
|
||||
-- transport — existed only in the enrolment reply, and nothing could hand one to a machine that
|
||||
-- had already joined. The row is the membership sealed to that machine, composed into its
|
||||
-- declaration as a file it reads after applying; the plaintext exists once, at minting, and then
|
||||
-- only on the machine. One per node: the mesh moves to one bus.
|
||||
create table bus_membership (
|
||||
node uuid primary key references node(id) on delete cascade,
|
||||
sealed text not null,
|
||||
since timestamptz not null default now()
|
||||
);
|
||||
@@ -106,3 +106,44 @@ func (i *Inventory) RenameSeat(ctx context.Context, from, to string) error {
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// HoldSeat records that one assignment holds a seat, replacing whoever held it — as one write, so
|
||||
// the seat is never without a holder in between (novox/hq ADR 0131, design 28 task 5.3). The
|
||||
// assignment must exist; the store refuses otherwise, and that refusal is the right one: a seat
|
||||
// cannot be handed to something that is not running anywhere.
|
||||
func (i *Inventory) HoldSeat(ctx context.Context, seat, scope, nodeName, module string) error {
|
||||
node, err := i.NodeByName(ctx, nodeName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = i.store.Pool().Exec(ctx,
|
||||
`insert into seat_holding (seat, scope, node, module) values ($1, $2, $3, $4)
|
||||
on conflict (seat) do update set scope = excluded.scope, node = excluded.node,
|
||||
module = excluded.module, since = now()`,
|
||||
seat, scope, node.ID, module)
|
||||
if err != nil {
|
||||
return fmt.Errorf("recording %s on %s as the holder of %s: %w", module, nodeName, seat, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Holdings is every seat whose holder is on record, as the resolver reads it. A seat with no row here
|
||||
// is held by derivation, exactly as before the table existed.
|
||||
func (i *Inventory) Holdings(ctx context.Context) ([]catalogue.Held, error) {
|
||||
rows, err := i.store.Pool().Query(ctx,
|
||||
`select h.seat, h.scope, n.name, h.module, coalesce(n.site, '')
|
||||
from seat_holding h join node n on n.id = h.node order by h.seat`)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []catalogue.Held
|
||||
for rows.Next() {
|
||||
var h catalogue.Held
|
||||
if err := rows.Scan(&h.Claim, &h.Scope, &h.Node, &h.Module, &h.Site); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, h)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// OverNats is the controller's outbound on the bus being built: the same three acts the other
|
||||
// transport has, on the subjects the permissions were derived for (design 25). A declaration is a
|
||||
// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on
|
||||
// the subject its name derives to; a tool is asked by request and reply on the module's tool subject.
|
||||
type OverNats struct{ JS *broker.JetStream }
|
||||
|
||||
// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the
|
||||
// only subject that node's consumer delivers. The host subscribes exactly this.
|
||||
func declareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
||||
|
||||
func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||
publish, cancel := context.WithTimeout(ctx, 15*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("declaring to %s: %w", node, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
|
||||
// The key is the subject: the controller's own events are named in full, and what a module
|
||||
// emits is derived before it reaches here. Headers carry the envelope the other transport put
|
||||
// in message properties (ADR 0042), so a consumer reads who and when without the payload.
|
||||
msg := nats.NewMsg(key)
|
||||
msg.Data = body
|
||||
msg.Header.Set("x-source", source)
|
||||
msg.Header.Set("x-node", node)
|
||||
msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano))
|
||||
if err := b.JS.Conn().PublishMsg(msg); err != nil {
|
||||
return fmt.Errorf("announcing %s: %w", key, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
||||
ask, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
|
||||
}
|
||||
return reply.Data, nil
|
||||
}
|
||||
|
||||
// ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one
|
||||
// JetStream connection the caller has already raised the streams on. Nothing is declared here —
|
||||
// the streams and the controller's consumers are asserted by Raise, before anything is served.
|
||||
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
|
||||
return &Server{
|
||||
inbound: Nats(js),
|
||||
bus: OverNats{JS: js},
|
||||
js: js,
|
||||
enroller: enroller,
|
||||
listener: listener,
|
||||
log: newLog(),
|
||||
}
|
||||
}
|
||||
|
||||
// Bus is the controller's outbound, whichever transport it connected over. Callers that send a
|
||||
// declaration or ask a tool use this rather than the channel, which one transport does not have.
|
||||
func (s *Server) Bus() Bus { return s.bus }
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"log"
|
||||
"os"
|
||||
"time"
|
||||
@@ -70,6 +71,7 @@ type Server struct {
|
||||
bus Bus
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
js *broker.JetStream
|
||||
|
||||
enroller Enroller
|
||||
listener Listener
|
||||
@@ -180,10 +182,12 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
||||
channel: channel,
|
||||
enroller: enroller,
|
||||
listener: listener,
|
||||
log: log.New(os.Stdout, "", log.LstdFlags),
|
||||
log: newLog(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
|
||||
|
||||
// Channel is the controller's channel, for the command line's own publishing.
|
||||
func (s *Server) Channel() *amqp.Channel { return s.channel }
|
||||
|
||||
@@ -197,6 +201,9 @@ func (s *Server) Close() {
|
||||
if s.conn != nil {
|
||||
_ = s.conn.Close()
|
||||
}
|
||||
if s.js != nil {
|
||||
s.js.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// Serve acts on what arrives until the context ends.
|
||||
|
||||
+6
-5
@@ -23,7 +23,8 @@
|
||||
"licences": "/var/lib/mesh/mesh-controller/licences",
|
||||
"broker": "/var/lib/mesh/mesh-controller/broker",
|
||||
"broker-management": "/var/lib/mesh/mesh-controller/broker-management",
|
||||
"broker-address": "/var/lib/mesh/mesh-controller/broker-address"
|
||||
"broker-address": "/var/lib/mesh/mesh-controller/broker-address",
|
||||
"bus": "/var/lib/mesh/mesh-controller/bus"
|
||||
},
|
||||
"secrets-owner": "65534:65534",
|
||||
"resources": [
|
||||
@@ -46,15 +47,14 @@
|
||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||
"MESH_STORE_IDENTITY_FILE": "/run/secrets/identity",
|
||||
"MESH_STORE_LICENCES_FILE": "/run/secrets/licences",
|
||||
"MESH_BROKER_AMQP_FILE": "/run/secrets/broker",
|
||||
"MESH_BROKER_MANAGEMENT_FILE": "/run/secrets/broker-management",
|
||||
"MESH_BROKER_ADDRESS_FILE": "/run/secrets/broker-address",
|
||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_STORE_IDENTITY_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_ADDRESS_PORT": "${seat:mesh-broker:5671}"
|
||||
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
||||
"MESH_BUS_NATS_FILE": "/run/secrets/bus"
|
||||
},
|
||||
"volumes": [
|
||||
"/var/lib/mesh-broker-tls:/broker-tls:ro",
|
||||
@@ -62,6 +62,7 @@
|
||||
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
|
||||
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences: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-address:/run/secrets/broker-address:ro"
|
||||
],
|
||||
@@ -82,7 +83,7 @@
|
||||
"on": [
|
||||
{
|
||||
"arg": "GO_BASE",
|
||||
"image": "golang@sha256:1ae0735f00daffa3aaf1363a5184c0d2dc55c78e3db4ec70241cdac97bf84b59"
|
||||
"image": "golang@sha256:8ac98ca534ac3f51e1f420a1dd2c15e74c75cfa0f23f3ad27eb5d7236c349a0c"
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user