Compare commits

...
Author SHA1 Message Date
jschoubben 3907ea0db0 The control plane serves and pushes on the bus it is told to
The seams were there and nothing chose a side: serve, push, ask and build all
opened the old bus's connection and declared over its channel, whatever
MESH_BUS_NATS said. So the switch moved every host and left the control plane
unable to follow — "this control plane has no MESH_BROKER_AMQP" with the new bus
named and standing (2026-09-28). That was task 4.3 of design 28, still open.

One place now decides: connectLink reads the switch, refuses both buses named at
once, raises the new bus's streams and this controller's consumers when it is
handed the inventory, and opens the link over whichever bus it is on. Every
caller that sent a declaration or asked a tool through the old channel goes
through the server's bus instead, which the new transport has and the channel is
not. OverNats is that outbound: a declaration is a JetStream publish into the
node's own subject, an event is announced on the subject its name derives to, a
tool is request and reply on the module's tool subject.
2026-09-28 01:29:44 +02:00
jschoubben 40f5e9a41c Merge pull request 'A machine may bind its consumer' (#102) from fix/a-node-may-bind-its-consumer into main 2026-09-27 23:18:20 +00:00
jschoubben 2c2eb51878 Only CONSUMER.INFO was missing from a machine's grants; the rest was already there 2026-09-28 01:17:40 +02:00
jschoubben aa2d0b51ea Golden: a machine's user may bind its consumer, ack, and hear its inbox 2026-09-28 01:17:15 +02:00
jschoubben 64d154d9d7 A machine may bind its consumer and hear the answer
Binding to a consumer asks the server about it and hears the answer on the
client's inbox; hearing a declaration acknowledges it. A machine's user was granted
none of that and was refused the first time one dialled a permissioned server:
"this node cannot read its declarations". Its inbox is its own prefix, which the
host now sets.
2026-09-28 01:16:46 +02:00
jschoubben ffa390f916 Merge pull request 'The mint leaves the control plane's old-bus secret alone' (#101) from fix/mint-leaves-the-control-planes-old-secret-alone into main 2026-09-27 23:08:35 +00:00
jschoubben 4d62e6caf1 The mint leaves the control plane's old-bus secret alone
The control plane is a module too, and its broker secret is the old bus's
credential it is still using while the mint runs. Writing the new bus's blob there
cut the mesh off from its own old bus mid-move. Its new-bus credential is the
controller principal's bus secret; the module principal is skipped.
2026-09-28 01:08:11 +02:00
jschoubben 386ae676ca Merge pull request 'The control plane mounts the bus secret it reads' (#100) from fix/the-controller-mounts-its-bus-secret into main 2026-09-27 23:03:42 +00:00
jschoubben f8a9c3d6bc The control plane mounts the bus secret it reads
MESH_BUS_NATS_FILE named /run/secrets/bus and nothing put a file there: the
manifest binds each secret explicitly, and the switch added the secret and the
variable but not the bind. Found live — the control plane came up on the new bus
and could not read its own credential.
2026-09-28 01:03:18 +02:00
jschoubben 83671fae5f Merge pull request 'The network map resolves each machine with the seat holders on record' (#99) from fix/the-network-map-knows-the-holders into main 2026-09-27 22:56:55 +00:00
9 changed files with 137 additions and 19 deletions
+1 -1
View File
@@ -37,7 +37,7 @@ func askCommand(ctx context.Context, args []string) error {
arguments = json.RawMessage(positionals[2]) arguments = json.RawMessage(positionals[2])
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
+2 -2
View File
@@ -368,7 +368,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
} }
defer ident.Close() defer ident.Close()
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -472,7 +472,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
return err return err
} }
defer ident.Close() defer ident.Close()
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
+35 -13
View File
@@ -38,6 +38,33 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
// nothing in it was wrong, and no one edit was the one that should have been a new file. // nothing in it was wrong, and no one edit was the one that should have been a new file.
// serve is the control plane running: one connection to the broker, one queue, one consumer. // serve is the control plane running: one connection to the broker, one queue, one consumer.
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
// so nothing served here finds them missing.
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
busAddress, onNATS, err := broker.OnNATS()
if err != nil {
return nil, err
}
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
return nil, err
}
if !onNATS {
return link.Connect(enroller, listener)
}
if inv != nil {
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
return nil, err
}
}
js, err := broker.Dial(busAddress)
if err != nil {
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
busAddress, err)
}
return link.ConnectNats(js, enroller, listener), nil
}
func serve(ctx context.Context) error { func serve(ctx context.Context) error {
open, err := openStores(ctx) open, err := openStores(ctx)
if err != nil { if err != nil {
@@ -90,7 +117,7 @@ func serve(ctx context.Context) error {
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known, work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
OnNATS: onNATS} OnNATS: onNATS}
server, err := link.Connect(work, work) server, err := connectLink(ctx, inv, work, work)
if err != nil { if err != nil {
return err return err
} }
@@ -100,11 +127,6 @@ func serve(ctx context.Context) error {
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was // somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
// replaced all have records and no objects — and a node whose consumer is missing hears nothing // replaced all have records and no objects — and a node whose consumer is missing hears nothing
// while everything else about it looks correct. // while everything else about it looks correct.
if onNATS {
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
return err
}
}
// And build results nobody was waiting for. A build triggered any other way than `build` // And build results nobody was waiting for. A build triggered any other way than `build`
// would otherwise be reported into the void, which is the same as not reporting it. // would otherwise be reported into the void, which is the same as not reporting it.
server.Records(builds{inv}) server.Records(builds{inv})
@@ -158,13 +180,13 @@ func declare(ctx context.Context, args []string) error {
return err return err
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
defer server.Close() defer server.Close()
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
return err return err
} }
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
@@ -271,7 +293,7 @@ func pushCommand(ctx context.Context, args []string) error {
return err return err
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -332,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err return err
} }
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
// After it is away, not before. A digest recorded for something that failed to send would // After it is away, not before. A digest recorded for something that failed to send would
@@ -415,7 +437,7 @@ func pushCommand(ctx context.Context, args []string) error {
return declarationWith(held, open, node, plan, settings, gens, Allocating) return declarationWith(held, open, node, plan, settings, gens, Allocating)
}, },
func(s readyNode, body []byte) error { func(s readyNode, body []byte) error {
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, if err := link.Declare(ctx, server.Bus(), ident, s.node, body,
15*time.Second); err != nil { 15*time.Second); err != nil {
return err return err
} }
@@ -628,7 +650,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
len(refusals), strings.Join(refusals, "\n\n")) len(refusals), strings.Join(refusals, "\n\n"))
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -639,7 +661,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
if err != nil { if err != nil {
return err return err
} }
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
record, err := inv.NodeByName(ctx, s.node) record, err := inv.NodeByName(ctx, s.node)
+8
View File
@@ -340,6 +340,14 @@ func rolloutMint(ctx context.Context, again bool) error {
machines++ machines++
case broker.KindModule: case broker.KindModule:
if p.Module == "mesh-controller" {
// The control plane is a module too, and its `broker` secret is the old bus's
// credential it is still using while this runs. Writing the new bus's blob there
// cut the mesh off from its own old bus mid-move (2026-09-28). Its new-bus credential
// is the controller principal's `bus` secret above; nothing else is needed here.
skipped++
continue
}
m, inShelf := shelf[p.Module] m, inShelf := shelf[p.Module]
if !inShelf { if !inShelf {
skipped++ skipped++
+8 -1
View File
@@ -244,7 +244,14 @@ func PermissionsFor(p Principal) (Permissions, error) {
case KindNode: case KindNode:
// A host publishes its own node's control traffic and subscribes its own declaration — // A host publishes its own node's control traffic and subscribes its own declaration —
// and nothing of any other node's. // and nothing of any other node's.
pub = []string{"mesh.control." + p.Node + ".>"} // And binding to its consumer, which asks the server about it (CONSUMER.INFO) — the one
// 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"} sub = []string{"mesh.node." + p.Node + ".declare"}
case KindModule: case KindModule:
+1 -1
View File
@@ -32,7 +32,7 @@ accounts {
subscribe: { allow: ["_INBOX.enrol.one.>"] } subscribe: { allow: ["_INBOX.enrol.one.>"] }
} } } }
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: { { user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
publish: { allow: ["$JS.ACK.NODES.one.>", "mesh.control.one.>"] } publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] } subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] }
} } } }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
+73
View File
@@ -0,0 +1,73 @@
package link
import (
"context"
"fmt"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
)
// OverNats is the controller's outbound on the bus being built: the same three acts the other
// transport has, on the subjects the permissions were derived for (design 25). A declaration is a
// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on
// the subject its name derives to; a tool is asked by request and reply on the module's tool subject.
type OverNats struct{ JS *broker.JetStream }
// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the
// only subject that node's consumer delivers. The host subscribes exactly this.
func declareSubject(node string) string { return "mesh.node." + node + ".declare" }
func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error {
publish, cancel := context.WithTimeout(ctx, 15*time.Second)
defer cancel()
if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil {
return fmt.Errorf("declaring to %s: %w", node, err)
}
return nil
}
func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
// The key is the subject: the controller's own events are named in full, and what a module
// emits is derived before it reaches here. Headers carry the envelope the other transport put
// in message properties (ADR 0042), so a consumer reads who and when without the payload.
msg := nats.NewMsg(key)
msg.Data = body
msg.Header.Set("x-source", source)
msg.Header.Set("x-node", node)
msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano))
if err := b.JS.Conn().PublishMsg(msg); err != nil {
return fmt.Errorf("announcing %s: %w", key, err)
}
return nil
}
func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
ask, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args)
if err != nil {
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
}
return reply.Data, nil
}
// ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one
// JetStream connection the caller has already raised the streams on. Nothing is declared here —
// the streams and the controller's consumers are asserted by Raise, before anything is served.
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
return &Server{
inbound: Nats(js),
bus: OverNats{JS: js},
js: js,
enroller: enroller,
listener: listener,
log: newLog(),
}
}
// Bus is the controller's outbound, whichever transport it connected over. Callers that send a
// declaration or ask a tool use this rather than the channel, which one transport does not have.
func (s *Server) Bus() Bus { return s.bus }
+8 -1
View File
@@ -7,6 +7,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"github.com/novox/mesh-controller/internal/broker"
"log" "log"
"os" "os"
"time" "time"
@@ -70,6 +71,7 @@ type Server struct {
bus Bus bus Bus
conn *amqp.Connection conn *amqp.Connection
channel *amqp.Channel channel *amqp.Channel
js *broker.JetStream
enroller Enroller enroller Enroller
listener Listener listener Listener
@@ -180,10 +182,12 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
channel: channel, channel: channel,
enroller: enroller, enroller: enroller,
listener: listener, listener: listener,
log: log.New(os.Stdout, "", log.LstdFlags), log: newLog(),
}, nil }, nil
} }
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
// Channel is the controller's channel, for the command line's own publishing. // Channel is the controller's channel, for the command line's own publishing.
func (s *Server) Channel() *amqp.Channel { return s.channel } func (s *Server) Channel() *amqp.Channel { return s.channel }
@@ -197,6 +201,9 @@ func (s *Server) Close() {
if s.conn != nil { if s.conn != nil {
_ = s.conn.Close() _ = s.conn.Close()
} }
if s.js != nil {
s.js.Close()
}
} }
// Serve acts on what arrives until the context ends. // Serve acts on what arrives until the context ends.
+1
View File
@@ -62,6 +62,7 @@
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro", "/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro", "/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro",
"/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro", "/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro",
"/var/lib/mesh/mesh-controller/bus:/run/secrets/bus:ro",
"/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro", "/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro",
"/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro" "/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro"
], ],