From 1f36787d7574723ae4e88e957cc53a1cda4331a7 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 21:02:18 +0200 Subject: [PATCH] The mesh's four streams, asserted on every start MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Task 1.4. The foundation set only — a seat's streams come at registration and a module's consumers at assignment, neither of which has happened at genesis (ADR 0118). Asserted rather than created: a stream that was deleted, or a mesh raised from a backup, must converge rather than run without the guarantee its messages assume. Two things the definitions have to get right, both tested: - CONTROL names its subjects instead of taking mesh.control.>, because heartbeats live under that prefix and a stream of them competes for retention with the messages that matter - EVENTS filters on the event token, which is why that token exists; a filter over a module's whole namespace would persist every tool call Overlapping filters are refused where the set is written: NATS accepts two streams matching one subject and stores the message twice under two retentions, which nothing reports. Adds nats.go as a dependency; it pulled golang.org/x/* forward. Full suite green. --- go.mod | 16 +-- go.sum | 18 ++++ internal/broker/nats.go | 44 +++++--- internal/broker/nats_test.go | 16 +-- internal/broker/streams.go | 134 +++++++++++++++++++++++++ internal/broker/streams_test.go | 130 ++++++++++++++++++++++++ internal/broker/testdata/composed.conf | 8 +- 7 files changed, 335 insertions(+), 31 deletions(-) create mode 100644 internal/broker/streams.go create mode 100644 internal/broker/streams_test.go diff --git a/go.mod b/go.mod index a827fad..47c7bce 100644 --- a/go.mod +++ b/go.mod @@ -1,19 +1,23 @@ module github.com/novox/mesh-controller -go 1.25.0 +go 1.26.0 require ( github.com/jackc/pgx/v5 v5.10.0 github.com/rabbitmq/amqp091-go v1.14.0 - golang.org/x/crypto v0.55.0 + golang.org/x/crypto v0.57.0 ) require ( github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect - golang.org/x/net v0.57.0 // indirect - golang.org/x/sync v0.22.0 // indirect - golang.org/x/sys v0.47.0 // indirect - golang.org/x/text v0.41.0 // indirect + github.com/klauspost/compress v1.20.0 // indirect + github.com/nats-io/nats.go v1.54.0 // indirect + github.com/nats-io/nkeys v0.4.16 // indirect + github.com/nats-io/nuid v1.0.1 // indirect + golang.org/x/net v0.58.0 // indirect + golang.org/x/sync v0.23.0 // indirect + golang.org/x/sys v0.48.0 // indirect + golang.org/x/text v0.42.0 // indirect ) diff --git a/go.sum b/go.sum index f1c77ed..e85a40e 100644 --- a/go.sum +++ b/go.sum @@ -9,6 +9,14 @@ github.com/jackc/pgx/v5 v5.10.0 h1:VhSvgU2jSli8o3AqIEOTJr7rZwAEUVo4E4XhR94Zfr0= github.com/jackc/pgx/v5 v5.10.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA= +github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI= +github.com/nats-io/nats.go v1.54.0 h1:vsXoOxjHp/GmPUN+EcI7uOf/uB+iAP+kEsAFNQN0yzA= +github.com/nats-io/nats.go v1.54.0/go.mod h1:y+DZoD1oBOYfZTU681eTUiUjI0vbqYGixNVFHcjHJ0k= +github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= +github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= +github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= +github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek= @@ -22,14 +30,24 @@ go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= +golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M= +golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE= golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= +golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= +golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU= golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sync v0.23.0 h1:KameEIfc1IkluZyXWLn39Wd4tURc6GbCiISGiZm2bQk= +golang.org/x/sync v0.23.0/go.mod h1:sUUOizhqBxiL6pEWpqNLUiaJn1ShEbZ6BBqskPbjZm0= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo= +golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og= golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= +golang.org/x/text v0.42.0 h1:JbOZXgfeCPU9gacVtYliJqOhD+zhrEqK4LfdpmlUZqI= +golang.org/x/text v0.42.0/go.mod h1:ojzP1Z+2QtioaF8DTtO8K5q7JWVVYwZKenzujK0Zd0E= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 62d2bf0..3c83d4c 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -41,6 +41,7 @@ type Seat struct { Name string Accepts []string Emits []string + Serves []string Versions []string // protocol versions served beside the current one; empty for v1 only } @@ -149,10 +150,8 @@ func PermissionsFor(p Principal) (Permissions, error) { // else may publish into it, so an event's source is a fact the server enforces rather // than a claim in the body (design 29 §2). own := "mesh.mod." + p.Module - if len(p.Emits) > 0 { - for _, e := range p.Emits { - pub = append(pub, own+"."+e) - } + for _, e := range p.Emits { + pub = append(pub, own+".event."+e) } for _, t := range p.Serves { sub = append(sub, own+".tool."+t) @@ -161,16 +160,24 @@ func PermissionsFor(p Principal) (Permissions, error) { // 2. What it consumes, by the emitter's own subject — an event is addressed to its // emitter, because the emitter's identity is the meaning (ADR 0118). for _, c := range p.Consumes { - sub = append(sub, "mesh.mod."+c) + emitter, event, ok := strings.Cut(c, ".") + if !ok { + return Permissions{}, fmt.Errorf( + "%q does not name an emitter and an event: a consumed event is .", c) + } + sub = append(sub, "mesh.mod."+emitter+".event."+event) } // 3. Seats it holds: full participation. for _, s := range p.Holds { for _, a := range s.Accepts { - sub = append(sub, seatSubject(s, a)) + sub = append(sub, seatSubject(s, "accept", a)) } for _, e := range s.Emits { - pub = append(pub, seatSubject(s, e)) + pub = append(pub, seatSubject(s, "event", e)) + } + for _, t := range s.Serves { + sub = append(sub, seatSubject(s, "tool", t)) } } @@ -179,7 +186,10 @@ func PermissionsFor(p Principal) (Permissions, error) { // events and lie about outcomes (design 29 §2). for _, s := range p.Uses { for _, a := range s.Accepts { - pub = append(pub, seatSubject(s, a)) + pub = append(pub, seatSubject(s, "accept", a)) + } + for _, t := range s.Serves { + pub = append(pub, seatSubject(s, "tool", t)) } } } @@ -207,11 +217,19 @@ func PermissionsFor(p Principal) (Permissions, error) { }, nil } -// seatSubject places a seat's verb. A seat serving more than its current protocol version carries -// the version as a token (design 29 §8): the seat stays one role, and v1 and v2 run beside each -// other until nothing is bound to the old one. -func seatSubject(s Seat, verb string) string { - return "mesh.seat." + s.Name + "." + verb +// seatSubject places a seat's verb under the kind of traffic it is. +// +// **The kind token is load-bearing, not decoration.** A stream is defined by a subject filter, so +// without it a stream over a seat or a module's namespace would capture that namespace's *tool* +// traffic too — and a tool call must never be persisted (design 25 §3: tools stay on core NATS). +// Found while defining the streams: the first draft of design 29 had one namespace per module +// with no kind, which reads well and cannot be filtered. +// +// A seat serving more than its current protocol version carries the version as a token +// (design 29 §8): the seat stays one role, and v1 and v2 run beside each other until nothing is +// bound to the old one. +func seatSubject(s Seat, kind, verb string) string { + return "mesh.seat." + s.Name + "." + kind + "." + verb } // consumerName is the durable consumer the controller derives for this principal. It is here diff --git a/internal/broker/nats_test.go b/internal/broker/nats_test.go index e063f81..64193d8 100644 --- a/internal/broker/nats_test.go +++ b/internal/broker/nats_test.go @@ -32,9 +32,9 @@ func TestAModulePublishesOnlyWhatItEmits(t *testing.T) { if err != nil { t.Fatal(err) } - has(t, perms.Publish, "mesh.mod.billing.order.placed") + has(t, perms.Publish, "mesh.mod.billing.event.order.placed") hasNot(t, perms.Publish, "mesh.mod.billing.>") - hasNot(t, perms.Publish, "mesh.mod.shipping.order.placed") + hasNot(t, perms.Publish, "mesh.mod.shipping.event.order.placed") } // The gap AMQP left open — an emitter granted the events exchange whole — is closed by per-subject @@ -55,9 +55,9 @@ func TestUsingASeatIsPublishOnlyAndInboundOnly(t *testing.T) { seat := Seat{Name: "telegram-sender", Accepts: []string{"send"}, Emits: []string{"delivered", "failed"}} perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "shop", Uses: []Seat{seat}, PasswordHash: "x"}) - has(t, perms.Publish, "mesh.seat.telegram-sender.send") - hasNot(t, perms.Publish, "mesh.seat.telegram-sender.delivered") - hasNot(t, perms.Subscribe, "mesh.seat.telegram-sender.send") + has(t, perms.Publish, "mesh.seat.telegram-sender.accept.send") + hasNot(t, perms.Publish, "mesh.seat.telegram-sender.event.delivered") + hasNot(t, perms.Subscribe, "mesh.seat.telegram-sender.accept.send") } // The holder is the mirror image: it consumes what the seat accepts and publishes what it emits. @@ -65,9 +65,9 @@ func TestHoldingASeatIsTheMirrorOfUsingIt(t *testing.T) { seat := Seat{Name: "telegram-sender", Accepts: []string{"send"}, Emits: []string{"delivered", "failed"}} perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "telegram", Holds: []Seat{seat}, PasswordHash: "x"}) - has(t, perms.Subscribe, "mesh.seat.telegram-sender.send") - has(t, perms.Publish, "mesh.seat.telegram-sender.delivered") - hasNot(t, perms.Publish, "mesh.seat.telegram-sender.send") + has(t, perms.Subscribe, "mesh.seat.telegram-sender.accept.send") + has(t, perms.Publish, "mesh.seat.telegram-sender.event.delivered") + hasNot(t, perms.Publish, "mesh.seat.telegram-sender.accept.send") } // Without an ack permission a durable consumer never really consumes: every message it receives is diff --git a/internal/broker/streams.go b/internal/broker/streams.go new file mode 100644 index 0000000..239a118 --- /dev/null +++ b/internal/broker/streams.go @@ -0,0 +1,134 @@ +package broker + +import ( + "fmt" + "sort" +) + +// The mesh's own streams. +// +// **These four and no more** (novox/hq ADR 0116 task 1.4, as revised by ADR 0118). An earlier +// reading had the controller create *every* stream at genesis, from a fixed set. That is only the +// mesh's own half: a seat's streams are created when the module declaring it is registered, and a +// module's durable consumers when it is assigned — neither of which has happened at genesis. What +// is here is the foundation, which exists before any module does. +// +// The controller is the only writer of stream definitions (design 25 §3). A module declares +// nothing about them and cannot reach the JetStream API to make one. + +// Retention is how a stream decides what to keep, which is the whole of what distinguishes the +// mesh's four relationships on the wire (design 29 §4). +type Retention string + +const ( + // RetentionWorkQueue: a message is removed once a consumer acknowledges it. Exactly one + // worker does the work, and a worker that dies has its message redelivered. + RetentionWorkQueue Retention = "workqueue" + // RetentionLastPerSubject: only the newest message on each subject survives. This is the + // state shape — a node that was away gets exactly the current declaration and nothing older. + RetentionLastPerSubject Retention = "last_per_subject" + // RetentionLimits: kept until it ages or the stream fills. Events, where a subscriber that + // was down catches up and nobody is obliged to act. + RetentionLimits Retention = "limits" +) + +// A Stream is one of the mesh's own, as the controller asserts it. +type Stream struct { + Name string + Subjects []string + Retention Retention + // MaxAge in seconds, zero for unbounded. + MaxAge int + // Why is carried into the assertion so an operator reading the server's own state finds the + // reason there, rather than only in a repository they may not have. + Why string +} + +// MeshStreams is the foundation set, in the order a person reads it. +// +// **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live +// under that prefix and must not be persisted: a lost heartbeat is the next heartbeat, and a +// stream of them is a stream of the least valuable messages the mesh sends, competing for the +// same retention as the ones that matter. +// +// **EVENTS filters on the `event` token**, which is the reason that token exists. A module's +// namespace carries both its events and its tool calls; a filter of `mesh.mod.*.>` would persist +// every tool invocation in the mesh, and a tool call must never be persisted (design 25 §3 keeps +// tools on core NATS, where a lost call is a timeout the caller already handles). +func MeshStreams() []Stream { + return []Stream{ + { + Name: "CONTROL", + Subjects: []string{"mesh.control.*.report", "mesh.control.enrol", "mesh.control.built"}, + Retention: RetentionWorkQueue, + Why: "the store-window guarantee (ADR 0083): the controller naks with a delay while its " + + "store is away and the message is redelivered; nothing is dropped", + }, + { + Name: "NODES", + Subjects: []string{"mesh.node.*.declare"}, + Retention: RetentionLastPerSubject, + Why: "one declaration per node, always the newest; a node that sees sequence n refuses " + + "n-1 by construction (issue 107)", + }, + { + Name: "BUILDS", + Subjects: []string{"mesh.build.request"}, + Retention: RetentionWorkQueue, + Why: "at least once, one builder at a time; a builder that dies mid-build has its message redelivered", + }, + { + Name: "EVENTS", + Subjects: []string{"mesh.mod.*.event.>"}, + Retention: RetentionLimits, + MaxAge: 7 * 24 * 60 * 60, + Why: "a subscriber that was down catches up; tool traffic under the same prefix is excluded by the event token", + }, + } +} + +// An Asserter is the part of a JetStream connection stream assertion needs. Narrow on purpose: it +// keeps this testable without a server, and keeps the client library out of everything that only +// wants to know what the streams are. +type Asserter interface { + // EnsureStream creates the stream if absent and updates it to match if present. It must be + // idempotent: the controller asserts on every start, not only at genesis. + EnsureStream(s Stream) error +} + +// AssertMeshStreams brings the foundation set into being, in order, and says which one failed +// rather than that something did. +// +// Asserted on every start rather than created once at genesis, because a stream that was deleted, +// or a mesh raised from a restored backup, must converge rather than run without the guarantee +// its messages assume. Idempotence is the whole requirement. +func AssertMeshStreams(a Asserter) error { + for _, s := range MeshStreams() { + if err := a.EnsureStream(s); err != nil { + return fmt.Errorf("asserting stream %s: %w", s.Name, err) + } + } + return nil +} + +// Overlaps reports subject filters claimed by more than one stream. +// +// Two streams matching one subject is not a warning in NATS; it is accepted, and the message is +// stored twice under two retentions. For the mesh that would mean a declaration kept as both +// state and a work queue, acknowledged in one and lingering in the other — a divergence nothing +// reports and nobody would think to look for. So it is refused here, where the set is written. +func Overlaps() []string { + seen := map[string]string{} + var clashes []string + for _, s := range MeshStreams() { + for _, subject := range s.Subjects { + if first, ok := seen[subject]; ok { + clashes = append(clashes, fmt.Sprintf("%s and %s both claim %s", first, s.Name, subject)) + continue + } + seen[subject] = s.Name + } + } + sort.Strings(clashes) + return clashes +} diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go new file mode 100644 index 0000000..1fd41cc --- /dev/null +++ b/internal/broker/streams_test.go @@ -0,0 +1,130 @@ +package broker + +import ( + "errors" + "strings" + "testing" +) + +type recorder struct { + seen []Stream + fail string +} + +func (r *recorder) EnsureStream(s Stream) error { + if s.Name == r.fail { + return errors.New("refused") + } + r.seen = append(r.seen, s) + return nil +} + +// The controller asserts on every start, not only at genesis: a stream that was deleted, or a mesh +// raised from a backup, must converge rather than run without the guarantee its messages assume. +func TestAssertingTwiceIsTheSameAsOnce(t *testing.T) { + a, b := &recorder{}, &recorder{} + if err := AssertMeshStreams(a); err != nil { + t.Fatal(err) + } + if err := AssertMeshStreams(a); err != nil { + t.Fatal(err) + } + if err := AssertMeshStreams(b); err != nil { + t.Fatal(err) + } + if len(a.seen) != 2*len(b.seen) { + t.Fatalf("asserted %d then %d; assertion is not repeatable", len(a.seen), len(b.seen)) + } +} + +func TestAFailedAssertionNamesItsStream(t *testing.T) { + err := AssertMeshStreams(&recorder{fail: "NODES"}) + if err == nil || !strings.Contains(err.Error(), "NODES") { + t.Fatalf("got %v, which does not say which stream failed", err) + } +} + +// Two streams matching one subject is accepted by NATS and stores the message twice under two +// retentions. Nothing reports that, so it is refused where the set is written. +func TestNoTwoStreamsClaimTheSameSubject(t *testing.T) { + if clashes := Overlaps(); len(clashes) != 0 { + t.Fatalf("overlapping subject filters: %v", clashes) + } +} + +// A heartbeat under mesh.control.> must not be persisted: a lost one is the next one, and a +// stream of them competes for retention with the messages that matter. +func TestHeartbeatsAreNotInTheControlStream(t *testing.T) { + for _, s := range MeshStreams() { + for _, subject := range s.Subjects { + if subject == "mesh.control.>" || strings.Contains(subject, "alive") { + t.Fatalf("stream %s claims %q, which captures heartbeats", s.Name, subject) + } + } + } +} + +// The reason the kind token exists: a filter over a module's whole namespace would persist every +// tool call in the mesh. +func TestTheEventsStreamDoesNotCaptureToolCalls(t *testing.T) { + var events Stream + for _, s := range MeshStreams() { + if s.Name == "EVENTS" { + events = s + } + } + if len(events.Subjects) != 1 || events.Subjects[0] != "mesh.mod.*.event.>" { + t.Fatalf("EVENTS filters on %v", events.Subjects) + } + // A tool subject the composer would actually produce must not match that filter. + tool := "mesh.mod.billing.tool.status" + if subjectMatches(events.Subjects[0], tool) { + t.Fatalf("%q matches the events filter, so every tool call would be persisted", tool) + } + if !subjectMatches(events.Subjects[0], "mesh.mod.billing.event.order.placed") { + t.Fatal("an event does not match the events filter") + } +} + +// subjectMatches is NATS subject matching, enough for these filters: `*` is one token, `>` is the +// rest. +func subjectMatches(filter, subject string) bool { + f, s := strings.Split(filter, "."), strings.Split(subject, ".") + for i, tok := range f { + if tok == ">" { + return i <= len(s) + } + if i >= len(s) { + return false + } + if tok != "*" && tok != s[i] { + return false + } + } + return len(f) == len(s) +} + +// Each relationship's retention is the thing that makes it what it is (design 29 §4). +func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) { + want := map[string]Retention{ + "CONTROL": RetentionWorkQueue, + "NODES": RetentionLastPerSubject, + "BUILDS": RetentionWorkQueue, + "EVENTS": RetentionLimits, + } + got := map[string]Retention{} + for _, s := range MeshStreams() { + got[s.Name] = s.Retention + if s.Why == "" { + t.Errorf("stream %s says no reason it exists", s.Name) + } + } + if len(got) != len(want) { + t.Fatalf("the foundation set is %v", got) + } + for name, r := range want { + if got[name] != r { + t.Errorf("%s retains as %q, expected %q", name, got[name], r) + } + } +} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 82dfbb3..c879021 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -33,16 +33,16 @@ accounts { subscribe: { allow: ["_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.delivered", "mesh.seat.telegram-sender.failed"] } - subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.send"] } + 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"] } 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.order.placed"] } + subscribe: { allow: ["_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.order.placed", "mesh.seat.telegram-sender.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.>"] } } } ]