Compare commits

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

One place now decides: connectLink reads the switch, refuses both buses named at
once, raises the new bus's streams and this controller's consumers when it is
handed the inventory, and opens the link over whichever bus it is on. Every
caller that sent a declaration or asked a tool through the old channel goes
through the server's bus instead, which the new transport has and the channel is
not. OverNats is that outbound: a declaration is a JetStream publish into the
node's own subject, an event is announced on the subject its name derives to, a
tool is request and reply on the module's tool subject.
2026-09-28 01:29:44 +02:00
jschoubben 40f5e9a41c Merge pull request 'A machine may bind its consumer' (#102) from fix/a-node-may-bind-its-consumer into main 2026-09-27 23:18:20 +00:00
9 changed files with 202 additions and 32 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
} }
+37 -15
View File
@@ -38,6 +38,33 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
// nothing in it was wrong, and no one edit was the one that should have been a new file. // nothing in it was wrong, and no one edit was the one that should have been a new file.
// serve is the control plane running: one connection to the broker, one queue, one consumer. // serve is the control plane running: one connection to the broker, one queue, one consumer.
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
// so nothing served here finds them missing.
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
busAddress, onNATS, err := broker.OnNATS()
if err != nil {
return nil, err
}
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
return nil, err
}
if !onNATS {
return link.Connect(enroller, listener)
}
if inv != nil {
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
return nil, err
}
}
js, err := broker.Dial(busAddress)
if err != nil {
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
broker.BareAddress(busAddress), err)
}
return link.ConnectNats(js, enroller, listener), nil
}
func serve(ctx context.Context) error { func serve(ctx context.Context) error {
open, err := openStores(ctx) open, err := openStores(ctx)
if err != nil { if err != nil {
@@ -90,7 +117,7 @@ func serve(ctx context.Context) error {
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known, work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
OnNATS: onNATS} OnNATS: onNATS}
server, err := link.Connect(work, work) server, err := connectLink(ctx, inv, work, work)
if err != nil { if err != nil {
return err return err
} }
@@ -100,11 +127,6 @@ func serve(ctx context.Context) error {
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was // somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
// replaced all have records and no objects — and a node whose consumer is missing hears nothing // replaced all have records and no objects — and a node whose consumer is missing hears nothing
// while everything else about it looks correct. // while everything else about it looks correct.
if onNATS {
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
return err
}
}
// And build results nobody was waiting for. A build triggered any other way than `build` // And build results nobody was waiting for. A build triggered any other way than `build`
// would otherwise be reported into the void, which is the same as not reporting it. // would otherwise be reported into the void, which is the same as not reporting it.
server.Records(builds{inv}) server.Records(builds{inv})
@@ -158,13 +180,13 @@ func declare(ctx context.Context, args []string) error {
return err return err
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
defer server.Close() defer server.Close()
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
return err return err
} }
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
@@ -271,7 +293,7 @@ func pushCommand(ctx context.Context, args []string) error {
return err return err
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -332,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err return err
} }
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
// After it is away, not before. A digest recorded for something that failed to send would // After it is away, not before. A digest recorded for something that failed to send would
@@ -415,7 +437,7 @@ func pushCommand(ctx context.Context, args []string) error {
return declarationWith(held, open, node, plan, settings, gens, Allocating) return declarationWith(held, open, node, plan, settings, gens, Allocating)
}, },
func(s readyNode, body []byte) error { func(s readyNode, body []byte) error {
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, if err := link.Declare(ctx, server.Bus(), ident, s.node, body,
15*time.Second); err != nil { 15*time.Second); err != nil {
return err return err
} }
@@ -628,7 +650,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
len(refusals), strings.Join(refusals, "\n\n")) len(refusals), strings.Join(refusals, "\n\n"))
} }
server, err := link.Connect(nil, nil) server, err := connectLink(ctx, nil, nil, nil)
if err != nil { if err != nil {
return err return err
} }
@@ -639,7 +661,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
if err != nil { if err != nil {
return err return err
} }
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
record, err := inv.NodeByName(ctx, s.node) record, err := inv.NodeByName(ctx, s.node)
@@ -704,7 +726,7 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
js, err := broker.Dial(address) js, err := broker.Dial(address)
if err != nil { if err != nil {
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w", return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
address, err) broker.BareAddress(address), err)
} }
defer js.Close() defer js.Close()
@@ -743,6 +765,6 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
return err return err
} }
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n", fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
address, len(names)) broker.BareAddress(address), len(names))
return nil return nil
} }
+53 -4
View File
@@ -1,8 +1,14 @@
package broker package broker
import ( import (
"crypto/sha256"
"crypto/tls"
"crypto/x509"
"encoding/hex"
"errors" "errors"
"fmt" "fmt"
"os"
"strings"
"time" "time"
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
@@ -24,21 +30,64 @@ type JetStream struct {
// Dial connects and returns the controller's JetStream handle. // Dial connects and returns the controller's JetStream handle.
func Dial(url string, opts ...nats.Option) (*JetStream, error) { func Dial(url string, opts ...nats.Option) (*JetStream, error) {
// A name, because a connection nobody can identify in the server's own monitoring is one
// nobody can attribute a problem to.
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second)) opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
// checks nothing else, and so does this). Without this, the first connection failed with
// "certificate is not valid for any names" against a bus that was answering (2026-09-28).
if path := strings.TrimSpace(os.Getenv(CertificateVar)); path != "" {
pinned, err := pinnedTo(path)
if err != nil {
return nil, err
}
opts = append(opts, nats.Secure(pinned))
}
// **Its own inbox, and nothing wider.** Every principal is granted `_INBOX.<its user>.>` and
// no other inbox; the client's default prefix is random, and the server refused the first
// subscription to it (2026-09-28). The user is in the URL, so the prefix follows from it.
if user, _, _ := CredentialIn(url); user != "" {
opts = append(opts, nats.CustomInboxPrefix("_INBOX."+user))
}
// The address in an error is the address alone. The URL carries this controller's password,
// and an error here is written on the assumption it will be logged.
where := BareAddress(url)
conn, err := nats.Connect(url, opts...) conn, err := nats.Connect(url, opts...)
if err != nil { if err != nil {
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err) return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
} }
js, err := conn.JetStream() js, err := conn.JetStream()
if err != nil { if err != nil {
conn.Close() conn.Close()
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err) return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
} }
return &JetStream{conn: conn, js: js}, nil return &JetStream{conn: conn, js: js}, nil
} }
// pinnedTo is a TLS configuration that accepts exactly the certificate in the file and no other:
// the leaf's SHA-256, compared on every handshake, with the name and the chain deliberately not
// consulted — a self-signed certificate with no names is the ordinary case for a mesh's bus.
func pinnedTo(path string) (*tls.Config, error) {
want, err := FingerprintOf(path)
if err != nil {
return nil, err
}
return &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
},
}, nil
}
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat, // Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
// a tool call — where a lost message is answered by the next one or by a timeout the caller // a tool call — where a lost message is answered by the next one or by a timeout the caller
// already handles (design 25 §3). // already handles (design 25 §3).
+15 -3
View File
@@ -169,7 +169,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// The controller owns the mesh's own traffic and the streams. It is the only writer of // The controller owns the mesh's own traffic and the streams. It is the only writer of
// stream definitions (design 25 §3), so it alone reaches the JetStream API. // stream definitions (design 25 §3), so it alone reaches the JetStream API.
pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"} pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"}
sub = []string{"mesh.control.>", "$JS.API.>"} // **And where its consumers deliver.** A push consumer delivers on `_DELIVER.<its name>`,
// and a client bound to it subscribes exactly that; the server refused it for every
// principal the first time one bound a consumer (2026-09-28). Each kind below is granted
// its own consumers' delivery subjects and no other's.
sub = []string{"mesh.control.>", "$JS.API.>", "_DELIVER." + ControllerName, "_DELIVER." + ControllerName + ".>"}
// Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A // Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A
// build is the one today: the controller asks, and reads the answer from the seat's event // build is the one today: the controller asks, and reads the answer from the seat's event
@@ -252,7 +256,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
"mesh.control." + p.Node + ".>", "mesh.control." + p.Node + ".>",
"$JS.API.CONSUMER.INFO.NODES." + p.Node, "$JS.API.CONSUMER.INFO.NODES." + p.Node,
} }
sub = []string{"mesh.node." + p.Node + ".declare"} sub = []string{"mesh.node." + p.Node + ".declare", "_DELIVER." + p.Node}
case KindModule: case KindModule:
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing // 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
@@ -285,7 +289,12 @@ func PermissionsFor(p Principal) (Permissions, error) {
} }
// 3. Seats it holds: full participation. // 3. Seats it holds: full participation.
// Its consumer's name, not ConsumerFor: that asks for these permissions to build the
// consumer, and would ask forever. A subject for a consumer that turns out not to exist
// grants nothing anybody can use.
sub = append(sub, "_DELIVER."+consumerDurable(p))
for _, s := range p.Holds { for _, s := range p.Holds {
sub = append(sub, "_DELIVER.SEAT_"+upperSnake(s.Name)+"_worker")
for _, a := range s.Accepts { for _, a := range s.Accepts {
sub = append(sub, seatSubject(s, "accept", a)) sub = append(sub, seatSubject(s, "accept", a))
} }
@@ -521,7 +530,10 @@ func ComposeAccounts(principals []Principal) (string, error) {
// One account for the mesh: accounts in NATS isolate subject spaces entirely, and the mesh is // One account for the mesh: accounts in NATS isolate subject spaces entirely, and the mesh is
// one space (design 25 §4). The cost of that — that permissions are the only isolation — is // one space (design 25 §4). The cost of that — that permissions are the only isolation — is
// paid in the scoping of every inbox and every ack subject. // paid in the scoping of every inbox and every ack subject.
b.WriteString("accounts {\n MESH {\n users = [\n") // JetStream is enabled per account once accounts exist at all: with only the global block set,
// a user in MESH is told "JetStream not enabled for account" the first time it binds a
// consumer, which is the first thing every host does (2026-09-28).
b.WriteString("accounts {\n MESH {\n jetstream: enabled\n users = [\n")
for _, p := range sorted { for _, p := range sorted {
perms, err := PermissionsFor(p) perms, err := PermissionsFor(p)
if err != nil { if err != nil {
+6 -5
View File
@@ -21,10 +21,11 @@ jetstream {
accounts { accounts {
MESH { MESH {
jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] } publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] } subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
@@ -33,20 +34,20 @@ accounts {
} } } }
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: { { user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] } publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] } subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} } } }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] } subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.one_telegram", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>"] } publish: { allow: ["$JS.ACK.EVENTS.two_audit.>"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] } subscribe: { allow: ["_DELIVER.two_audit", "_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
} } } }
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: { { user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] } publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>"] } subscribe: { allow: ["_DELIVER.two_shop", "_INBOX.two.shop.>"] }
} } } }
] ]
} }
+7 -1
View File
@@ -186,7 +186,13 @@ func TestWhatTheMeshWritesIsUsersAndNothingAboutTheServer(t *testing.T) {
} }
// None of the server's own settings. Each of these in the mesh's file is a value the controller // None of the server's own settings. Each of these in the mesh's file is a value the controller
// would then own, and the module could no longer change its own image without the mesh agreeing. // would then own, and the module could no longer change its own image without the mesh agreeing.
for _, absent := range []string{"port:", "http:", "jetstream", "tls {", "store_dir", "cert_file"} { // `jetstream {` is the server's block (its store, its limits); `jetstream: enabled` inside the
// account is the account's, and the mesh owns the account — a user in it is told "JetStream
// not enabled for account" without it (2026-09-28).
if !strings.Contains(got, "jetstream: enabled") {
t.Errorf("the account does not enable JetStream, so no user in it can bind a consumer")
}
for _, absent := range []string{"port:", "http:", "jetstream {", "tls {", "store_dir", "cert_file"} {
if strings.Contains(got, absent) { if strings.Contains(got, absent) {
t.Errorf("the accounts file contains %q, which belongs to the module that raises the "+ t.Errorf("the accounts file contains %q, which belongs to the module that raises the "+
"server, not to the mesh", absent) "server, not to the mesh", absent)
+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.