diff --git a/internal/broker/management.go b/internal/broker/management.go index 2873408..caf8739 100644 --- a/internal/broker/management.go +++ b/internal/broker/management.go @@ -67,6 +67,104 @@ const ExchangeName = "mesh" // constant that differed would be caught by the test that asserts they agree. const BuildQueueName = "builds" +// The events bus (novox/hq ADR 0047): one topic exchange every event rides, a second for tool +// RPC kept apart, and a dead-letter home for a poison event. The substrate owns these — a module's +// account cannot declare them, only bind its own queue to the events one. +const ( + EventsExchangeName = "mesh.events" + RPCExchangeName = "mesh.rpc" + DeadExchangeName = "mesh.events.dead" +) + +// ModuleQueueFor is the durable queue a module consumes its events from — one per module per node +// (novox/hq ADR 0047), named so the account that may read it is exactly this module's. +func ModuleQueueFor(node, module string) string { return node + "." + module + ".events" } + +// modulePermissions is a module's authority on the bus, derived from its manifest (novox/hq +// ADR 0048): what it consumes and what it emits, and nothing else. Pure, so the scope is tested as +// patterns without a broker — the way a builder's is. +// +// A note on the limit: the broker's write permission is per exchange, not per routing key (LavinMQ +// has no topic permissions), so an emitting module is granted the events exchange whole. ADR 0047's +// origin reservation — a module publishes only under `module..*` — is stamped by the sdk, not +// enforced here; that gap is the broker's, and is recorded rather than hidden. A pure consumer like +// the audit logger is unaffected: it is granted no write to the exchange at all. +func modulePermissions(node, module string, emits, consumes []string) (configure, write, read string) { + queue := regexp.QuoteMeta(ModuleQueueFor(node, module)) + events := regexp.QuoteMeta(EventsExchangeName) + + // Declare only its own queue. + configure = "^" + queue + "$" + + // Write to its own queue — binding a queue to an exchange is a write on the queue — and to the + // events exchange only if it emits. + writes := []string{queue} + if len(emits) > 0 { + writes = append(writes, events) + } + write = "^(" + strings.Join(writes, "|") + ")$" + + // Read its own queue to consume it, and the events exchange to bind onto, only if it consumes. + reads := []string{queue} + if len(consumes) > 0 { + reads = append(reads, events) + } + read = "^(" + strings.Join(reads, "|") + ")$" + return configure, write, read +} + +// CreateModuleAccount gives an assigned module its own broker account, scoped by what it emits and +// consumes (novox/hq ADR 0048). The account name carries the node so the same module on two machines +// holds two accounts, each sealed to its own; the permissions carry the module so one module cannot +// read another's queue. Generic — the builder is one instance of this rule, not a separate kind. +func (m *Management) CreateModuleAccount(ctx context.Context, node, module, password string, emits, consumes []string) (string, error) { + if !safeName.MatchString(node) { + return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", node) + } + if !safeName.MatchString(module) { + return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", module) + } + account := node + "-" + module + if !safeName.MatchString(account) { + return "", fmt.Errorf("%q is not a usable broker account name", account) + } + + if err := m.put(ctx, "/api/users/"+url.PathEscape(account), + map[string]string{"password": password, "tags": ""}); err != nil { + return "", fmt.Errorf("cannot create the broker account for %s on %s: %w", module, node, err) + } + + configure, write, read := modulePermissions(node, module, emits, consumes) + if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(account), map[string]string{ + "configure": configure, "write": write, "read": read, + }); err != nil { + return "", fmt.Errorf("cannot scope the broker account for %s on %s: %w", module, node, err) + } + return account, nil +} + +// EnsureEventExchanges declares the bus's exchanges and the dead-letter home, idempotently. The +// substrate owns them (a module's account may not declare an exchange), and a dead-letter exchange +// with no queue behind it drops what it receives — so a durable queue bound to `#` retains a poison +// event for inspection, which is the whole reason the trail exists. +func (m *Management) EnsureEventExchanges(ctx context.Context) error { + for _, exchange := range []string{EventsExchangeName, RPCExchangeName, DeadExchangeName} { + if err := m.put(ctx, "/api/exchanges/%2f/"+url.PathEscape(exchange), + map[string]any{"type": "topic", "durable": true}); err != nil { + return fmt.Errorf("cannot declare the %s exchange: %w", exchange, err) + } + } + if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(DeadExchangeName), + map[string]any{"durable": true}); err != nil { + return fmt.Errorf("cannot declare the dead-letter queue: %w", err) + } + if err := m.post(ctx, "/api/bindings/%2f/e/"+url.PathEscape(DeadExchangeName)+ + "/q/"+url.PathEscape(DeadExchangeName), map[string]string{"routing_key": "#"}); err != nil { + return fmt.Errorf("cannot bind the dead-letter queue: %w", err) + } + return nil +} + // CreateNodeAccount gives a node its own broker account, with the token's secret as the password. // // Scoped so a node can reach its own queue and the one exchange, and nothing else. The patterns @@ -167,6 +265,10 @@ func (m *Management) put(ctx context.Context, path string, body any) error { return m.do(ctx, http.MethodPut, path, body) } +func (m *Management) post(ctx context.Context, path string, body any) error { + return m.do(ctx, http.MethodPost, path, body) +} + func (m *Management) get(ctx context.Context, path string) ([]byte, error) { request, err := m.request(ctx, http.MethodGet, path, nil) if err != nil { diff --git a/internal/broker/module_account_test.go b/internal/broker/module_account_test.go new file mode 100644 index 0000000..d14bce8 --- /dev/null +++ b/internal/broker/module_account_test.go @@ -0,0 +1,96 @@ +package broker_test + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "regexp" + "strings" + "testing" + + "github.com/novox/mesh-control/internal/broker" +) + +// The audit logger consumes everything and emits nothing. Its account must let it declare and read +// its own queue and read the events exchange to bind onto — and must not reach another module's +// queue, nor grant any write to the events exchange (novox/hq ADR 0048). +func TestAConsumerReadsItsOwnQueueAndTheEventsExchangeAndNoOthers(t *testing.T) { + // Rebuilt through CreateModuleAccount's own path by asking for the same scope it would apply. + // The queue this module reads: + mine := broker.ModuleQueueFor("anchor", "audit-logger") + other := broker.ModuleQueueFor("anchor", "plex") + + read := scope(t, "read", "anchor", "audit-logger", nil, []string{"#"}) + if !read.MatchString(mine) { + t.Error("the audit logger may not read its own queue, so it consumes nothing") + } + if !read.MatchString(broker.EventsExchangeName) { + t.Error("the audit logger may not read the events exchange, so it cannot bind onto it") + } + if read.MatchString(other) { + t.Error("the audit logger may read another module's queue") + } + + // It emits nothing, so it is granted no write to the events exchange — only its own queue, to bind. + write := scope(t, "write", "anchor", "audit-logger", nil, []string{"#"}) + if write.MatchString(broker.EventsExchangeName) { + t.Error("a pure consumer was granted write to the events exchange") + } + if !write.MatchString(mine) { + t.Error("the audit logger may not write to its own queue, so it cannot bind it") + } +} + +// An emitter is granted the events exchange to write; a consumer is not. +func TestAnEmitterMayWriteTheEventsExchangeAndAConsumerMayNot(t *testing.T) { + emitter := scope(t, "write", "anchor", "umami", []string{"module.umami.site.created"}, nil) + if !emitter.MatchString(broker.EventsExchangeName) { + t.Error("an emitting module may not write the events exchange, so it cannot emit") + } +} + +// scope reconstructs one of the three permission patterns CreateModuleAccount would apply, by +// reading it back from a captured request against a stub management API. +func scope(t *testing.T, which, node, module string, emits, consumes []string) *regexp.Regexp { + t.Helper() + pat := capturePermission(t, which, node, module, emits, consumes) + re, err := regexp.Compile(pat) + if err != nil { + t.Fatalf("the %s pattern does not compile: %v", which, err) + } + return re +} + +// capturePermission runs CreateModuleAccount against a stub management API and returns the pattern +// it set for `which` ("configure"/"write"/"read"). The scope is tested where it is applied, not +// reconstructed by the test — so a change to the mapping cannot pass a test that hard-codes the old +// one. +func capturePermission(t *testing.T, which, node, module string, emits, consumes []string) string { + t.Helper() + var captured map[string]string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasPrefix(r.URL.Path, "/api/permissions/") { + _ = json.NewDecoder(r.Body).Decode(&captured) + } + w.WriteHeader(http.StatusNoContent) + })) + defer server.Close() + + t.Setenv(broker.ManagementVar, server.URL) + m, err := broker.ManagementFromEnvironment() + if err != nil { + t.Fatalf("stub management not usable: %v", err) + } + if _, err := m.CreateModuleAccount(context.Background(), node, module, "pw", emits, consumes); err != nil { + t.Fatalf("CreateModuleAccount: %v", err) + } + if captured == nil { + t.Fatal("no permissions were set") + } + pattern, ok := captured[which] + if !ok { + t.Fatalf("no %s permission was set; got %v", which, captured) + } + return pattern +}