Compare commits
11
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d1e488efaf | ||
|
|
c5dc7e732a | ||
|
|
7efcccd013 | ||
|
|
964285f08c | ||
|
|
5698dda11f | ||
|
|
6005a8471f | ||
|
|
e2ee0dfe98 | ||
|
|
70341cfbc7 | ||
|
|
77643aa3f4 | ||
|
|
1fd6194ff8 | ||
|
|
c37018fdd2 |
@@ -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
|
// machine told about both would take work from one and answer on the other, and every log line would
|
||||||
// say it was fine.
|
// say it was fine.
|
||||||
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
||||||
|
// **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()
|
address, onNATS, err := broker.OnNATS()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
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
|
// **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels
|
||||||
// out of band — here, sealed with the credential — and the endpoint is verified once at connect.
|
// out of band — here, sealed with the credential — and the endpoint is verified once at connect.
|
||||||
type Credential struct {
|
type Credential struct {
|
||||||
URL string `json:"url"`
|
URL string `json:"url"`
|
||||||
// Fingerprint is SHA-256 over the broker certificate's DER bytes, or empty to verify the
|
|
||||||
// ordinary way.
|
|
||||||
Fingerprint string `json:"fingerprint,omitempty"`
|
Fingerprint string `json:"fingerprint,omitempty"`
|
||||||
|
// User and Password ride beside the address on the bus being built (design 25): a credential
|
||||||
|
// embedded in a URL leaks into every log line that prints a connection, so the mesh seals them
|
||||||
|
// as two fields and this machine joins them once, here, to dial.
|
||||||
|
User string `json:"user,omitempty"`
|
||||||
|
Password string `json:"password,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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.
|
// 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{
|
allowed := map[string]bool{
|
||||||
"string(body)": true, // once.go: the result JSON, which IS stdout
|
"string(body)": true, // once.go: the result JSON, which IS stdout
|
||||||
"version)": true, // --version
|
"version)": true, // --version
|
||||||
`"stopping")`: true, // the loop.s shutdown line
|
`"stopping")`: true, // the loop.s shutdown line
|
||||||
"usage)": true, // --help text, for a human
|
"usage)": true, // --help text, for a human
|
||||||
}
|
}
|
||||||
for _, file := range []string{"once.go", "main.go"} {
|
for _, file := range []string{"once.go", "main.go"} {
|
||||||
src, err := os.ReadFile(file)
|
src, err := os.ReadFile(file)
|
||||||
|
|||||||
@@ -60,7 +60,7 @@ func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.En
|
|||||||
js, err := broker.Dial(busAddress)
|
js, err := broker.Dial(busAddress)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
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
|
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)
|
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()
|
||||||
|
|
||||||
@@ -761,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,
|
// The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder,
|
||||||
// deliberately: work queues until somebody arrives to do it, so assigning a build machine a week
|
// deliberately: work queues until somebody arrives to do it, so assigning a build machine a week
|
||||||
// after something started asking for builds flushes the backlog instead of having lost it.
|
// after something started asking for builds flushes the backlog instead of having lost it.
|
||||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), nil); err != nil {
|
// With the seats' holders, so each role's work queue gets the consumer its holder takes
|
||||||
|
// work from. Passed as nil until the first live raise, which left the build machine bound to a
|
||||||
|
// consumer nothing had created (2026-09-28).
|
||||||
|
holders, err := seatHolders(ctx, inv)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
||||||
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
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,8 +1,14 @@
|
|||||||
package broker
|
package broker
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"crypto/sha256"
|
||||||
|
"crypto/tls"
|
||||||
|
"crypto/x509"
|
||||||
|
"encoding/hex"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/nats-io/nats.go"
|
"github.com/nats-io/nats.go"
|
||||||
@@ -24,21 +30,78 @@ type JetStream struct {
|
|||||||
|
|
||||||
// Dial connects and returns the controller's JetStream handle.
|
// Dial connects and returns the controller's JetStream handle.
|
||||||
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
||||||
// A name, because a connection nobody can identify in the server's own monitoring is one
|
|
||||||
// nobody can attribute a problem to.
|
|
||||||
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
||||||
|
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
|
||||||
|
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
|
||||||
|
// checks nothing else, and so does this). Without this, the first connection failed with
|
||||||
|
// "certificate is not valid for any names" against a bus that was answering (2026-09-28).
|
||||||
|
if path := strings.TrimSpace(os.Getenv(CertificateVar)); path != "" {
|
||||||
|
pinned, err := pinnedTo(path)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
opts = append(opts, nats.Secure(pinned))
|
||||||
|
}
|
||||||
|
// **Its own inbox, and nothing wider.** Every principal is granted `_INBOX.<its user>.>` and
|
||||||
|
// no other inbox; the client's default prefix is random, and the server refused the first
|
||||||
|
// subscription to it (2026-09-28). The user is in the URL, so the prefix follows from it.
|
||||||
|
if user, _, _ := CredentialIn(url); user != "" {
|
||||||
|
opts = append(opts, nats.CustomInboxPrefix("_INBOX."+user))
|
||||||
|
}
|
||||||
|
// The address in an error is the address alone. The URL carries this controller's password,
|
||||||
|
// and an error here is written on the assumption it will be logged.
|
||||||
|
where := BareAddress(url)
|
||||||
conn, err := nats.Connect(url, opts...)
|
conn, err := nats.Connect(url, opts...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err)
|
return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
|
||||||
}
|
}
|
||||||
js, err := conn.JetStream()
|
js, err := conn.JetStream()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
conn.Close()
|
conn.Close()
|
||||||
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err)
|
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
|
||||||
}
|
}
|
||||||
return &JetStream{conn: conn, js: js}, nil
|
return &JetStream{conn: conn, js: js}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// pinnedTo is a TLS configuration that accepts exactly the certificate in the file and no other:
|
||||||
|
// the leaf's SHA-256, compared on every handshake, with the name and the chain deliberately not
|
||||||
|
// consulted — a self-signed certificate with no names is the ordinary case for a mesh's bus.
|
||||||
|
func pinnedTo(path string) (*tls.Config, error) {
|
||||||
|
want, err := FingerprintOf(path)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return PinnedToFingerprint(want), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// DialPinned is Dial with the server's certificate pinned by a fingerprint the caller already holds
|
||||||
|
// — a module or a build machine that was handed one beside its credential, and has no file.
|
||||||
|
func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) {
|
||||||
|
if strings.TrimSpace(fingerprint) != "" {
|
||||||
|
opts = append(opts, nats.Secure(PinnedToFingerprint(fingerprint)))
|
||||||
|
}
|
||||||
|
return Dial(url, opts...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// PinnedToFingerprint accepts exactly the certificate with this SHA-256 and no other.
|
||||||
|
func PinnedToFingerprint(want string) *tls.Config {
|
||||||
|
return &tls.Config{
|
||||||
|
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
|
||||||
|
MinVersion: tls.VersionTLS12,
|
||||||
|
VerifyPeerCertificate: func(rawCerts [][]byte, _ [][]*x509.Certificate) error {
|
||||||
|
if len(rawCerts) == 0 {
|
||||||
|
return errors.New("the bus presented no certificate")
|
||||||
|
}
|
||||||
|
sum := sha256.Sum256(rawCerts[0])
|
||||||
|
got := "sha256:" + hex.EncodeToString(sum[:])
|
||||||
|
if got != want {
|
||||||
|
return fmt.Errorf("the bus presented a certificate this mesh does not know (%s…), expected %s…", got[:23], want[:23])
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
||||||
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
||||||
// already handles (design 25 §3).
|
// already handles (design 25 §3).
|
||||||
|
|||||||
+21
-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
|
// 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,18 @@ 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 {
|
||||||
|
// Taking work from the role's queue: the worker consumer it binds (asked about,
|
||||||
|
// delivered on, acknowledged), each on the seat's own stream. The first machine to
|
||||||
|
// take work over the new bus was refused the asking (2026-09-28).
|
||||||
|
worker := "SEAT_" + upperSnake(s.Name) + "_worker"
|
||||||
|
stream := seatStreamName(s.Name)
|
||||||
|
sub = append(sub, "_DELIVER."+worker)
|
||||||
|
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>")
|
||||||
for _, a := range s.Accepts {
|
for _, a := range s.Accepts {
|
||||||
sub = append(sub, seatSubject(s, "accept", a))
|
sub = append(sub, seatSubject(s, "accept", a))
|
||||||
}
|
}
|
||||||
@@ -521,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 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 {
|
||||||
|
|||||||
+7
-6
@@ -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.>", "$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: ["_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.>"] }
|
||||||
} }
|
} }
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -143,3 +143,18 @@ func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The store's seat rows have no protocol columns yet; loading them must not drop the protocol the
|
||||||
|
// bus is derived from, or no role's work queue is ever raised (found live, 2026-09-28).
|
||||||
|
func TestAStoreRowWithoutAProtocolKeepsTheCompiledOne(t *testing.T) {
|
||||||
|
was := Seats()
|
||||||
|
t.Cleanup(func() { UseSeats(was) })
|
||||||
|
UseSeats([]Seat{{Name: "mesh-build-machine", Scope: ScopeMesh, Decision: "row"}})
|
||||||
|
got, ok := SeatNamed("mesh-build-machine")
|
||||||
|
if !ok || len(got.Accepts) == 0 {
|
||||||
|
t.Fatalf("the build machine's seat lost what it accepts when loaded from the store: %+v", got)
|
||||||
|
}
|
||||||
|
if got.Decision != "row" {
|
||||||
|
t.Fatalf("the store's own columns were not kept: %+v", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -13,7 +13,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
|
|||||||
"type": "container", "id": "server", "name": "mesh-controller",
|
"type": "container", "id": "server", "name": "mesh-controller",
|
||||||
"env": map[string]any{
|
"env": map[string]any{
|
||||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||||
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
|
|
||||||
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
||||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||||
},
|
},
|
||||||
@@ -27,7 +26,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
|
|||||||
env := control["env"].(map[string]any)
|
env := control["env"].(map[string]any)
|
||||||
for key, want := range map[string]string{
|
for key, want := range map[string]string{
|
||||||
"MESH_STORE_INVENTORY_PORT": "6852",
|
"MESH_STORE_INVENTORY_PORT": "6852",
|
||||||
"MESH_BROKER_AMQP_PORT": "5679",
|
|
||||||
"MESH_BROKER_ADDRESS_PORT": "5671",
|
"MESH_BROKER_ADDRESS_PORT": "5671",
|
||||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||||
} {
|
} {
|
||||||
@@ -146,7 +144,6 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
|
|||||||
"MESH_STORE_INVENTORY_PORT": "6852",
|
"MESH_STORE_INVENTORY_PORT": "6852",
|
||||||
"MESH_STORE_IDENTITY_PORT": "6852",
|
"MESH_STORE_IDENTITY_PORT": "6852",
|
||||||
"MESH_STORE_LICENCES_PORT": "6852",
|
"MESH_STORE_LICENCES_PORT": "6852",
|
||||||
"MESH_BROKER_AMQP_PORT": "5679",
|
|
||||||
"MESH_BROKER_MANAGEMENT_PORT": "15673",
|
"MESH_BROKER_MANAGEMENT_PORT": "15673",
|
||||||
"MESH_BROKER_ADDRESS_PORT": "5671",
|
"MESH_BROKER_ADDRESS_PORT": "5671",
|
||||||
} {
|
} {
|
||||||
@@ -164,7 +161,7 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any)
|
env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any)
|
||||||
if env["MESH_STORE_INVENTORY_PORT"] != "" || env["MESH_BROKER_AMQP_PORT"] != "" {
|
if env["MESH_STORE_INVENTORY_PORT"] != "" {
|
||||||
t.Errorf("with no settings, the control plane is told %v", env)
|
t.Errorf("with no settings, the control plane is told %v", env)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -175,7 +172,6 @@ var SeatPorts = map[string]string{
|
|||||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||||
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
|
||||||
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
|
||||||
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
|
|
||||||
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
|
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
|
||||||
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -109,9 +109,30 @@ func DefaultSeats() []Seat { return append([]Seat(nil), defaultSeats...) }
|
|||||||
// than running on the set the binary shipped with. So the store can only ever *replace* the set with
|
// than running on the set the binary shipped with. So the store can only ever *replace* the set with
|
||||||
// a non-empty one, never erase it.
|
// a non-empty one, never erase it.
|
||||||
func UseSeats(s []Seat) {
|
func UseSeats(s []Seat) {
|
||||||
if len(s) > 0 {
|
if len(s) == 0 {
|
||||||
seats = s
|
return
|
||||||
}
|
}
|
||||||
|
// **The store's rows carry no protocol yet, and the protocol is what the bus is derived
|
||||||
|
// from.** ADR 0129 gives a seat what it accepts, emits and serves; ADR 0122 moved the set into
|
||||||
|
// a table that has name, scope, delivers and decision and nothing else, and the columns for
|
||||||
|
// the rest are not there yet. So a row replacing a compiled entry would silently drop the
|
||||||
|
// protocol, and the roles' work queues would never be raised — found live as "no response
|
||||||
|
// from stream" the first time a build was submitted over the new bus (2026-09-28). Until the
|
||||||
|
// table gains the columns, a row without a protocol keeps the compiled one of the same name.
|
||||||
|
byName := map[string]Seat{}
|
||||||
|
for _, d := range defaultSeats {
|
||||||
|
byName[d.Name] = d
|
||||||
|
}
|
||||||
|
merged := make([]Seat, 0, len(s))
|
||||||
|
for _, row := range s {
|
||||||
|
if len(row.Accepts)+len(row.Emits)+len(row.Serves) == 0 {
|
||||||
|
if d, known := byName[row.Name]; known {
|
||||||
|
row.Accepts, row.Emits, row.Serves = d.Accepts, d.Emits, d.Serves
|
||||||
|
}
|
||||||
|
}
|
||||||
|
merged = append(merged, row)
|
||||||
|
}
|
||||||
|
seats = merged
|
||||||
}
|
}
|
||||||
|
|
||||||
// aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from
|
// aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from
|
||||||
|
|||||||
Reference in New Issue
Block a user