From 48581af35cd8d1614d2a18bb5ab2facfe30dd579 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 18:20:57 +0200 Subject: [PATCH] Grant the bus's own module the snapshot API and nothing else (hq ADR 0235) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The night's backup of the bus takes each stream through JetStream's snapshot API, run by the nats module under its own account. The module holding mesh-broker is composed that account: stream names and info, the snapshot request, its flow-control acks, its own inbox — no write, which the writers table checks. A bus module declaring anything else to say on the bus is refused by module check rather than silently granted nothing. The genesis user list is unchanged: the controller's grants are. --- internal/broker/nats.go | 41 ++++++ internal/broker/nats_golden_test.go | 3 + internal/broker/snapshot_live_test.go | 183 ++++++++++++++++++++++++ internal/broker/snapshot_test.go | 75 ++++++++++ internal/broker/testdata/composed.conf | 4 + internal/broker/users.go | 5 +- internal/catalogue/bus_snapshot_test.go | 67 +++++++++ internal/catalogue/manifest.go | 23 +++ internal/catalogue/seats.go | 6 +- internal/inventory/busrecords.go | 3 + internal/inventory/busrecords_test.go | 55 +++++++ 11 files changed, 463 insertions(+), 2 deletions(-) create mode 100644 internal/broker/snapshot_live_test.go create mode 100644 internal/broker/snapshot_test.go create mode 100644 internal/catalogue/bus_snapshot_test.go diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 7b0cd29..173b0c5 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -110,6 +110,11 @@ type Principal struct { State []string Reads []string + // SnapshotsTheBus is the bus's own module, the one holding mesh-broker (novox/hq ADR 0235). Its + // whole authority is BusSnapshotGrants: it copies the streams for the night's backup and nothing + // else — not its declarations, which the node's tool runtime serves for it. + SnapshotsTheBus bool + // PasswordHash is the bcrypt hash the mesh minted. The plaintext is sealed to the principal // and never appears here: this file is written to a node's disk and read by a server, and a // secret that can be read from a configuration file is a secret with a wider blast radius @@ -215,6 +220,31 @@ type Permissions struct { // refused, which is why a holder answers before it does what can reload it. const ResponseTTL = time.Minute +// BusSnapshotGrants is the whole authority of the bus's own module (novox/hq ADR 0235): the +// snapshot API and what it needs, and nothing that changes a stream. +// +// **Read-only, and the writers table proves it**: none of these overlaps a subject a write of the +// mesh's state is a publish to — not a stream's definition, not a bucket, not a message. A snapshot is +// the server reading its own blocks out to the asker, while it goes on taking writes; it neither +// pauses nor reconfigures the stream. Named one by one rather than `$JS.API.>`, which would be the +// controller's authority over every stream: +// +// - the stream names, and one stream's information (whether it is on disk at all — a memory stream +// cannot be snapshotted and is said as skipped); +// - the snapshot request itself, for any stream: the archive comes to the asker's own inbox; +// - the acknowledgement the server waits for past its window: a publish to the subject the server +// put on each chunk, `$JS.SNAPSHOT.ACK....`, which only the server's +// own internal subscription hears. +// +// Restoring is not here: a restore creates a stream, which is the controller's to define, and the +// mesh restores into a new store beside the live one, swapped in by a person (to-be 43). +var BusSnapshotGrants = struct{ Publish []string }{Publish: []string{ + "$JS.API.STREAM.NAMES", + "$JS.API.STREAM.INFO.*", + "$JS.API.STREAM.SNAPSHOT.*", + "$JS.SNAPSHOT.ACK.>", +}} + // PermissionsFor derives a principal's authority. Pure, and the only place authority is decided: // a permission that cannot be derived from a declaration is a permission nobody can explain. func PermissionsFor(p Principal) (Permissions, error) { @@ -230,6 +260,17 @@ func PermissionsFor(p Principal) (Permissions, error) { } } + if p.Kind == KindModule && p.SnapshotsTheBus { + pub := append([]string(nil), BusSnapshotGrants.Publish...) + sort.Strings(pub) + if err := CheckWriters(p, pub); err != nil { + return Permissions{}, err + } + // Its own inbox for the answers and the archive's chunks, and nothing else: nothing is asked + // of it, so it answers nothing. + return Permissions{Publish: pub, Subscribe: []string{p.inbox()}}, nil + } + var pub, sub []string switch p.Kind { case KindController: diff --git a/internal/broker/nats_golden_test.go b/internal/broker/nats_golden_test.go index 6c3b2da..0157174 100644 --- a/internal/broker/nats_golden_test.go +++ b/internal/broker/nats_golden_test.go @@ -28,6 +28,9 @@ func TestTheComposedConfigMatchesTheGolden(t *testing.T) { Emits: []string{"order.placed"}, PasswordHash: "$2a$11$ssssssssssssssssssssss"}, {Kind: KindModule, Node: "two", Module: "audit", Consumes: []string{"shop.order.placed"}, PasswordHash: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa"}, + // The bus's own module: the snapshot API and its inbox, nothing else (novox/hq ADR 0235). + {Kind: KindModule, Node: "one", Module: "nats", SnapshotsTheBus: true, + Serves: []string{"nats_streams"}, PasswordHash: "$2a$11$bbbbbbbbbbbbbbbbbbbbbb"}, }) if err != nil { t.Fatal(err) diff --git a/internal/broker/snapshot_live_test.go b/internal/broker/snapshot_live_test.go new file mode 100644 index 0000000..af9d33c --- /dev/null +++ b/internal/broker/snapshot_live_test.go @@ -0,0 +1,183 @@ +package broker + +import ( + "encoding/json" + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" + "golang.org/x/crypto/bcrypt" +) + +// The bus's own module's composed authority, against the bus's release in a throwaway container +// (novox/hq ADR 0235): it can take a whole snapshot of a stream through the server's chunked +// protocol, and every write — to a stream's definition, its messages, a bucket — is refused. +// +// MESH_TEST_DOCKER=1 go test ./internal/broker/ -run TestTheComposedSnapshotUser +// +// The program that takes the night's snapshot lives in the nats module (mesh-catalog +// modules/nats/snapshot) and is tested there against these same grants; this proves the grants are +// what the controller composes, by composing them. +func TestTheComposedSnapshotUserCanSnapshotAndCannotWrite(t *testing.T) { + if os.Getenv("MESH_TEST_DOCKER") != "1" { + t.Skip("MESH_TEST_DOCKER is not 1: this starts a throwaway nats container") + } + hash := func(pw string) string { + h, err := bcrypt.GenerateFromPassword([]byte(pw), bcrypt.MinCost) + if err != nil { + t.Fatal(err) + } + return string(h) + } + accounts, err := ComposeAccounts([]Principal{ + {Kind: KindController, PasswordHash: hash("controller")}, + {Kind: KindModule, Node: "anchor", Module: "nats", SnapshotsTheBus: true, PasswordHash: hash("snap")}, + }) + if err != nil { + t.Fatal(err) + } + conf, data := t.TempDir(), t.TempDir() + _ = os.WriteFile(filepath.Join(conf, "accounts.conf"), []byte(accounts), 0o644) + _ = os.WriteFile(filepath.Join(conf, "nats.conf"), []byte("port: 4222\njetstream { store_dir: \"/data\" }\ninclude accounts.conf\n"), 0o644) + name := fmt.Sprintf("mesh-controller-snapshot-test-%d", time.Now().UnixNano()) + run := exec.Command("docker", "run", "-d", "--rm", "--name", name, "--user", fmt.Sprintf("%d:%d", os.Getuid(), os.Getgid()), + "-p", "127.0.0.1::4222", "-v", conf+":/etc/nats:ro", "-v", data+":/data", "nats:2.11-alpine", "-c", "/etc/nats/nats.conf") + if out, err := run.CombinedOutput(); err != nil { + t.Fatalf("%v\n%s", err, out) + } + t.Cleanup(func() { _ = exec.Command("docker", "rm", "-f", name).Run() }) + out, err := exec.Command("docker", "port", name, "4222/tcp").Output() + if err != nil { + t.Fatal(err) + } + _, port, _ := strings.Cut(strings.TrimSpace(strings.Split(string(out), "\n")[0]), ":") + url := "nats://127.0.0.1:" + port + + dial := func(user, pw string, errs chan error) *nats.Conn { + var nc *nats.Conn + var err error + for deadline := time.Now().Add(15 * time.Second); ; time.Sleep(200 * time.Millisecond) { + nc, err = nats.Connect(url, nats.UserInfo(user, pw), nats.CustomInboxPrefix("_INBOX."+user), + nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) { + if errs != nil { + select { + case errs <- err: + default: + } + } + })) + if err == nil || time.Now().After(deadline) { + break + } + } + if err != nil { + t.Fatal(err) + } + t.Cleanup(nc.Close) + return nc + } + + // The controller defines the streams and fills them, as it does. + controller := dial("controller", "controller", nil) + js, _ := controller.JetStream() + if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}}); err != nil { + t.Fatal(err) + } + kv, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: HandActsBucket}) + if err != nil { + t.Fatal(err) + } + if _, err := kv.Put("act", []byte("done by hand")); err != nil { + t.Fatal(err) + } + + errs := make(chan error, 32) + snap := dial("anchor.nats", "snap", errs) + + // What it may do: list, read a stream's information, and take a whole snapshot. + names, err := snap.Request("$JS.API.STREAM.NAMES", nil, 5*time.Second) + if err != nil || !strings.Contains(string(names.Data), "KV_"+HandActsBucket) { + t.Fatalf("it cannot list the streams: %v", err) + } + deliver := snap.NewRespInbox() + done := make(chan string, 1) + var bytes int + sub, err := snap.Subscribe(deliver, func(m *nats.Msg) { + if len(m.Data) == 0 { + done <- m.Header.Get("Status") + return + } + bytes += len(m.Data) + if m.Reply != "" { + _ = snap.Publish(m.Reply, nil) + } + }) + if err != nil { + t.Fatal(err) + } + defer sub.Unsubscribe() + req, _ := json.Marshal(map[string]any{"deliver_subject": deliver}) + answer, err := snap.Request("$JS.API.STREAM.SNAPSHOT.KV_"+HandActsBucket, req, 5*time.Second) + if err != nil || strings.Contains(string(answer.Data), `"error"`) { + t.Fatalf("the snapshot was not granted: %v", err) + } + select { + case status := <-done: + if status != "" && status != "204" { + t.Fatalf("the snapshot ended with %s", status) + } + case <-time.After(10 * time.Second): + t.Fatal("the snapshot never finished") + } + if bytes == 0 { + t.Fatal("the snapshot delivered nothing") + } + select { + case e := <-errs: + t.Fatalf("taking a snapshot met a refusal: %v", e) + default: + } + + // What it may not: every write, and every subscription beyond its own inbox. + for _, subject := range []string{"$JS.API.STREAM.CREATE.NEW", "$JS.API.STREAM.UPDATE.EVENTS", + "$JS.API.STREAM.DELETE.EVENTS", "$JS.API.STREAM.PURGE.EVENTS", "$JS.API.STREAM.MSG.DELETE.EVENTS", + "$JS.API.STREAM.RESTORE.NEW", "$JS.API.CONSUMER.CREATE.EVENTS", "$KV." + HandActsBucket + ".act", + "mesh.mod.nats.event.anything", "mesh.node.anchor.declare"} { + if _, err := snap.Request(subject, []byte(`{}`), 300*time.Millisecond); err == nil { + t.Errorf("%s was answered", subject) + } + select { + case e := <-errs: + if !strings.Contains(strings.ToLower(e.Error()), "permissions violation") { + t.Errorf("%s: %v", subject, e) + } + case <-time.After(2 * time.Second): + t.Errorf("the bus did not refuse %s", subject) + } + } + for _, subject := range []string{"mesh.>", "_INBOX.controller.>", "$KV.>"} { + if _, err := snap.SubscribeSync(subject); err != nil { + t.Fatal(err) + } + _ = snap.Flush() + select { + case e := <-errs: + if !strings.Contains(strings.ToLower(e.Error()), "permissions violation") { + t.Errorf("subscribing %s: %v", subject, e) + } + case <-time.After(2 * time.Second): + t.Errorf("the bus let it subscribe %s", subject) + } + } + if info, err := js.StreamInfo("EVENTS"); err != nil || info == nil { + t.Fatalf("the stream is gone: %v", err) + } + if e, err := kv.Get("act"); err != nil || string(e.Value()) != "done by hand" { + t.Fatalf("the bucket changed: %v", err) + } +} diff --git a/internal/broker/snapshot_test.go b/internal/broker/snapshot_test.go new file mode 100644 index 0000000..dd9d34f --- /dev/null +++ b/internal/broker/snapshot_test.go @@ -0,0 +1,75 @@ +package broker + +import ( + "slices" + "testing" +) + +// The bus's own module copies the bus and does nothing else on it (novox/hq ADR 0235): its whole +// authority is the snapshot API and its own inbox, whatever else it declared — its tools are the +// node's runtime's to serve. +func TestTheBussOwnModuleMaySnapshotAndNothingElse(t *testing.T) { + p := Principal{Kind: KindModule, Node: "anchor", Module: "nats", SnapshotsTheBus: true, + Serves: []string{"nats_streams"}, PasswordHash: "x"} + perms, err := PermissionsFor(p) + if err != nil { + t.Fatal(err) + } + want := []string{"$JS.API.STREAM.INFO.*", "$JS.API.STREAM.NAMES", "$JS.API.STREAM.SNAPSHOT.*", "$JS.SNAPSHOT.ACK.>"} + if !slices.Equal(perms.Publish, want) { + t.Errorf("it may publish %v, want exactly %v", perms.Publish, want) + } + if !slices.Equal(perms.Subscribe, []string{"_INBOX.anchor.nats.>"}) { + t.Errorf("it may subscribe %v, want its own inbox alone", perms.Subscribe) + } + if perms.AllowResponses { + t.Error("nothing is asked of it, and it may answer") + } +} + +// Read-only, by the writers table: no grant of it overlaps a subject a write of the mesh's state is a +// publish to — a stream's definition, a bucket, a message. +func TestTheSnapshotGrantsWriteNothing(t *testing.T) { + p := Principal{Kind: KindModule, Node: "anchor", Module: "nats"} + if err := CheckWriters(p, BusSnapshotGrants.Publish); err != nil { + t.Fatal(err) + } + for _, grant := range BusSnapshotGrants.Publish { + for _, write := range []string{"$JS.API.STREAM.CREATE.X", "$JS.API.STREAM.UPDATE.X", "$JS.API.STREAM.DELETE.X", + "$JS.API.STREAM.PURGE.X", "$JS.API.STREAM.MSG.DELETE.X", "$JS.API.STREAM.RESTORE.X", + "$JS.API.CONSUMER.CREATE.X", "$JS.API.CONSUMER.DURABLE.CREATE.X.y", "$KV.b.k", "mesh.mod.m.event.e"} { + if SubjectsOverlap(grant, write) { + t.Errorf("%s would let the bus's own module publish %s", grant, write) + } + } + } +} + +// A module that is not the bus gets nothing of it, and the flag travels from the records to the user. +func TestOnlyTheBussOwnModuleIsGrantedTheSnapshot(t *testing.T) { + users, err := Users(Records{Nodes: []string{"anchor"}, Assigned: map[string][]Declared{"anchor": { + {Module: "nats", SnapshotsTheBus: true}, + {Module: "shop", Emits: []string{"order.placed"}}, + }}}) + if err != nil { + t.Fatal(err) + } + for _, u := range users { + perms, err := PermissionsFor(u) + if err != nil { + t.Fatal(err) + } + snapshots := slices.Contains(perms.Publish, "$JS.API.STREAM.SNAPSHOT.*") + switch u.Username() { + case "anchor.nats": + if !snapshots { + t.Error("the bus's own module is not granted the snapshot") + } + case "controller": + default: + if snapshots { + t.Errorf("%s may snapshot the bus", u.Username()) + } + } + } +} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index d8934f5..e7ab6c2 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -36,6 +36,10 @@ accounts { publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] } subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.ask.report", "mesh.node.one.declare"] } } } + { user: "one.nats", password: "$2a$11$bbbbbbbbbbbbbbbbbbbbbb", permissions: { + publish: { allow: ["$JS.API.STREAM.INFO.*", "$JS.API.STREAM.NAMES", "$JS.API.STREAM.SNAPSHOT.*", "$JS.SNAPSHOT.ACK.>"] } + subscribe: { allow: ["_INBOX.one.nats.>"] } + } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.telegram", "$SRV.INFO.telegram.>", "$SRV.PING", "$SRV.PING.telegram", "$SRV.PING.telegram.>", "$SRV.STATS", "$SRV.STATS.telegram", "$SRV.STATS.telegram.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } diff --git a/internal/broker/users.go b/internal/broker/users.go index 4c1589f..0bcc8c6 100644 --- a/internal/broker/users.go +++ b/internal/broker/users.go @@ -41,6 +41,9 @@ type Declared struct { // delivered to it and nothing can connect as it (novox/hq issue 195). Said in the negative so a // record that does not say is composed as it always was. NoAccount bool + // SnapshotsTheBus says the module holds mesh-broker — it is the bus — and so is the one module + // granted the snapshot API, to copy the bus's streams for the night's backup (novox/hq ADR 0235). + SnapshotsTheBus bool } // Records is what composing a user list needs to know about the mesh, and nothing more. @@ -98,7 +101,7 @@ func Users(r Records) ([]Principal, error) { Kind: KindModule, Node: node, Module: d.Module, Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves, Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, Invokes: d.Invokes, - State: stateNames(d.State), Reads: d.Reads, + State: stateNames(d.State), Reads: d.Reads, SnapshotsTheBus: d.SnapshotsTheBus, }) } if runtimeHere { diff --git a/internal/catalogue/bus_snapshot_test.go b/internal/catalogue/bus_snapshot_test.go new file mode 100644 index 0000000..7ba2571 --- /dev/null +++ b/internal/catalogue/bus_snapshot_test.go @@ -0,0 +1,67 @@ +package catalogue + +import ( + "os" + "strings" + "testing" +) + +// The module holding mesh-broker is the bus, and its account is granted the bus's snapshot API and +// nothing else (novox/hq ADR 0235). Anything it declared to say or hear on the bus would be granted +// nothing, so it is refused at the parser rather than silently dropped. +func TestTheBussAccountSaysNothingOnTheBus(t *testing.T) { + raw := []byte(`{"module":"bus","version":"1","provides":[{"name":"mesh-bus","scope":"mesh"}], + "claims":[{"name":"mesh-broker","scope":"mesh"}],"own-secrets":{"broker":"/run/broker"}, + "emits":["something.happened"]}`) + _, err := ParseManifest(raw) + if err == nil || !strings.Contains(err.Error(), "granted the bus's snapshot API and nothing else") { + t.Fatalf("a bus module declaring what it emits was not refused, or not for the reason: %v", err) + } + + quiet := []byte(`{"module":"bus","version":"1","provides":[{"name":"mesh-bus","scope":"mesh"}], + "claims":[{"name":"mesh-broker","scope":"mesh"}],"own-secrets":{"broker":"/run/broker"}, + "tools":["bus_streams"]}`) + m, err := ParseManifest(quiet) + if err != nil { + t.Fatalf("a bus module with an account and tools was refused: %v", err) + } + if !m.ClaimsSeat(BrokerSeat) { + t.Fatal("it does not read as holding the broker seat") + } +} + +// The catalogue's bus protects its streams by the snapshot, not as live files: its streams' item is a +// dump into its snapshots, it has the account the dump runs as, and the dump runs the snapshot +// program in the bus's own container. +func TestTheCataloguesBusIsBackedUpBySnapshot(t *testing.T) { + raw, err := os.ReadFile("../../../mesh-catalog/modules/nats/module.json") + if err != nil { + t.Skip("the catalogue is not checked out beside this repository") + } + m, err := ParseManifest(raw) + if err != nil { + t.Fatal(err) + } + if _, account := m.OwnSecrets["broker"]; !account { + t.Fatal("the bus declares no account, so its snapshot could not reach it") + } + var found bool + if m.Data == nil { + t.Fatal("the bus declares no data") + } + for _, it := range m.Data.Own { + if it.ID != "jetstream" { + continue + } + found = true + if !it.Backup.IsDump() || it.Backup.Into != "snapshots" { + t.Fatalf("the bus's streams are protected by %+v, not a snapshot into its snapshots", it.Backup) + } + if !strings.Contains(it.Backup.Dump, "mesh-nats-snapshot snapshot") { + t.Fatalf("the dump does not run the snapshot program: %s", it.Backup.Dump) + } + } + if !found { + t.Fatal("the bus declares no item for its streams") + } +} diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index b2ca58c..630e70b 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -1402,6 +1402,29 @@ func ParseManifest(raw []byte) (Manifest, error) { m.Module, offer.Name)) } } + // **The bus's own account copies the bus and says nothing on it** (novox/hq ADR 0235). The module + // holding mesh-broker is granted the snapshot API and nothing else, so anything it declared to say + // or hear on the bus would be granted nothing — refused here rather than silently dropped. + if _, account := m.OwnSecrets["broker"]; account && m.ClaimsSeat(BrokerSeat) { + var said []string + if len(m.EmitsAll()) > 0 { + said = append(said, "emits") + } + for _, f := range []struct { + name string + n int + }{{"consumes", len(m.Consumes)}, {"uses", len(m.Uses)}, {"invokes", len(m.Invokes)}, + {"state", len(m.State)}, {"reads", len(m.Reads)}} { + if f.n > 0 { + said = append(said, f.name) + } + } + if len(said) > 0 { + problems = append(problems, fmt.Sprintf( + "%s holds %s and declares a bus account, which is granted the bus's snapshot API and nothing "+ + "else (novox/hq ADR 0235); its %s would be granted nothing", m.Module, BrokerSeat, strings.Join(said, ", "))) + } + } for _, r := range m.Requires { if r == "amqp" { problems = append(problems, fmt.Sprintf( diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index 4ebd6b8..e5928ab 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -106,7 +106,7 @@ var defaultSeats = append([]Seat{ // the mesh's own transport. ADR 0128 then made that connection something a module requires // rather than receives ambiently — 23 of the catalogue's modules never speak, and an ambient // connection would mint a credential for each. - {Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "novox/hq ADR 0079"}, + {Name: BrokerSeat, Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "novox/hq ADR 0079"}, // The vault: the controller seals every minted credential with what it provides, which is the // test for a seat of the mesh's own (novox/hq ADR 0161) — a second provider of `secret` is a // second claimant, refused by name, rather than a candidate for a pin. @@ -635,3 +635,7 @@ func buildAgentVerbs() []Verb { Input: schema(map[string]string{}, nil)}, } } + +// BrokerSeat is the mesh's bus. Its holder is the bus, and its account — when it declares one — may +// snapshot the bus's streams and do nothing else (novox/hq ADR 0235). +const BrokerSeat = "mesh-broker" diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index 8de4341..af7b5fc 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -141,6 +141,9 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio if _, reads := m.OwnSecrets["broker"]; !reads { d.NoAccount = true } + // The bus's own module — the one holding mesh-broker — copies the bus's streams for the night's + // backup, and its account is granted that and nothing else (novox/hq ADR 0235). + d.SnapshotsTheBus = m.ClaimsSeat(catalogue.BrokerSeat) for _, c := range m.Claims { // Every seat with a protocol, the mesh's own included. One that says only who does a job is // not here and grants nothing, which is most of them. diff --git a/internal/inventory/busrecords_test.go b/internal/inventory/busrecords_test.go index fa9dfcb..3cafa96 100644 --- a/internal/inventory/busrecords_test.go +++ b/internal/inventory/busrecords_test.go @@ -241,3 +241,58 @@ func TestIssuingATokenRecordsTheAccountItIsThePasswordOf(t *testing.T) { } } } + +// The module holding mesh-broker is the bus, and its user is the bus's own: granted the snapshot API +// for the night's backup and nothing else, whatever tools it declares (novox/hq ADR 0235). Another +// module on the same machine is granted none of it. +func TestTheBussOwnModuleBecomesTheSnapshotUser(t *testing.T) { + bus := catalogue.Manifest{ + Module: "nats", Version: "1", + Provides: []catalogue.Offer{{Name: "mesh-bus", Scope: catalogue.ScopeMesh}}, + Claims: []catalogue.Claim{{Name: catalogue.BrokerSeat, Scope: catalogue.ScopeMesh}}, + Tools: []string{"nats_streams"}, + OwnSecrets: catalogue.OwnSecrets{"broker": {Path: "/run/broker"}}, + } + shop := catalogue.Manifest{Module: "shop", Version: "1", Emits: []string{"order.placed"}, + OwnSecrets: catalogue.OwnSecrets{"broker": {Path: "/run/broker"}}} + inv, ctx := aMeshWith(t, bus, shop) + if _, err := inv.AddNode(ctx, "one"); err != nil { + t.Fatal(err) + } + for _, module := range []string{"nats", "shop"} { + if _, err := inv.Assign(ctx, "one", module); err != nil { + t.Fatal(err) + } + } + records, err := inv.BusRecords(ctx) + if err != nil { + t.Fatal(err) + } + users, err := broker.Users(records) + if err != nil { + t.Fatal(err) + } + seen := map[string]bool{} + for _, u := range users { + perms, err := broker.PermissionsFor(u) + if err != nil { + t.Fatal(err) + } + snapshots := granted(perms.Publish, "$JS.API.STREAM.SNAPSHOT.*") + switch u.Username() { + case "one.nats": + seen[u.Username()] = true + if !snapshots || len(perms.Publish) != len(broker.BusSnapshotGrants.Publish) { + t.Fatalf("the bus's own user is granted %v, not the snapshot API alone", perms.Publish) + } + case "one.shop": + seen[u.Username()] = true + if snapshots { + t.Fatal("a module that is not the bus may snapshot it") + } + } + } + if !seen["one.nats"] || !seen["one.shop"] { + t.Fatalf("users derived: %v", seen) + } +}