Merge pull request 'Grant the bus's own module the snapshot API and nothing else (hq ADR 0235)' (#90) from feat/bus-snapshot into main
mesh/delivery delivered

This commit was merged in pull request #90.
This commit is contained in:
2026-10-06 16:24:53 +00:00
11 changed files with 463 additions and 2 deletions
+41
View File
@@ -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.<stream>.<id>.<size>.<index>`, 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:
+3
View File
@@ -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)
+183
View File
@@ -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)
}
}
+75
View File
@@ -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())
}
}
}
}
+4
View File
@@ -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"] }
+4 -1
View File
@@ -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 {
+67
View File
@@ -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")
}
}
+23
View File
@@ -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(
+5 -1
View File
@@ -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"
+3
View File
@@ -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.
+55
View File
@@ -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)
}
}