Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6005a8471f | ||
|
|
e2ee0dfe98 | ||
|
|
70341cfbc7 | ||
|
|
77643aa3f4 | ||
|
|
1fd6194ff8 | ||
|
|
c37018fdd2 |
@@ -60,7 +60,7 @@ func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.En
|
||||
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)
|
||||
broker.BareAddress(busAddress), err)
|
||||
}
|
||||
return link.ConnectNats(js, enroller, listener), nil
|
||||
}
|
||||
@@ -726,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()
|
||||
|
||||
@@ -765,6 +765,6 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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,64 @@ 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 &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,
|
||||
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
||||
// already handles (design 25 §3).
|
||||
|
||||
+15
-3
@@ -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
|
||||
@@ -252,7 +256,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
"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", "_DELIVER." + p.Node}
|
||||
|
||||
case KindModule:
|
||||
// 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.
|
||||
// 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 {
|
||||
sub = append(sub, "_DELIVER.SEAT_"+upperSnake(s.Name)+"_worker")
|
||||
for _, a := range s.Accepts {
|
||||
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 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 {
|
||||
|
||||
+6
-5
@@ -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: {
|
||||
@@ -33,20 +34,20 @@ accounts {
|
||||
} }
|
||||
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
||||
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: {
|
||||
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" }
|
||||
} }
|
||||
{ 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)
|
||||
|
||||
Reference in New Issue
Block a user