diff --git a/cmd/mesh-control/main.go b/cmd/mesh-control/main.go index 43a58d6..613f709 100644 --- a/cmd/mesh-control/main.go +++ b/cmd/mesh-control/main.go @@ -135,6 +135,7 @@ func usage() { module list what modules this mesh knows about module moved the source has a newer commit than the mesh built module forget remove one, unless a node is running it + module issue --node a broker account for a module, scoped to its emits and consumes status [--json] what is wrong, what is quiet, and what is out of date board [--listen ADDR] the same three questions, as a page that holds nothing api --issuer URL [--listen A] assign and unassign over http, for a surface that is not here diff --git a/cmd/mesh-control/modules.go b/cmd/mesh-control/modules.go index cfa81cd..3d10bd7 100644 --- a/cmd/mesh-control/modules.go +++ b/cmd/mesh-control/modules.go @@ -2,6 +2,8 @@ package main import ( "context" + "crypto/rand" + "encoding/base64" "encoding/json" "errors" "flag" @@ -10,6 +12,7 @@ import ( "sort" "strings" + "github.com/novox/mesh-control/internal/broker" "github.com/novox/mesh-control/internal/catalogue" "github.com/novox/mesh-control/internal/inventory" "github.com/novox/mesh-control/internal/overlay" @@ -194,8 +197,93 @@ func moduleCommand(ctx context.Context, args []string) error { fmt.Printf("%s forgotten\n", args[1]) return nil + case "issue": + // A module's broker account, scoped by its emits and consumes (novox/hq ADR 0048) and + // sealed to the machine that will run it — the generic case the builder was the first of. + set := flag.NewFlagSet("module issue", flag.ContinueOnError) + forNode := set.String("node", "", + "the machine that will run it, so the credential is delivered instead of printed") + positionals, err := parseAround(set, args[1:]) + if err != nil { + return err + } + if len(positionals) != 1 { + return errors.New("module issue --node ") + } + module := positionals[0] + if *forNode == "" { + return errors.New("module issue needs --node: a module's account is sealed to the " + + "machine that runs it, and the mesh cannot read it back to print") + } + + shelf, err := inv.Catalogue(ctx) + if err != nil { + return err + } + m, ok := shelf[module] + if !ok { + return fmt.Errorf("this mesh knows no module %q; `module add` it first", module) + } + + management, err := broker.ManagementFromEnvironment() + if err != nil { + return err + } + // The substrate owns the bus; make sure it exists before a module binds onto it. + if err := management.EnsureEventExchanges(ctx); err != nil { + return err + } + + secret := make([]byte, 32) + if _, err := rand.Read(secret); err != nil { + return err + } + password := base64.RawURLEncoding.EncodeToString(secret) + account, err := management.CreateModuleAccount(ctx, *forNode, module, password, m.Emits, m.Consumes) + if err != nil { + return err + } + // A consumer's queue, with its dead-letter, is the substrate's to declare — its own account + // may not (ADR 0048). Made now, so it exists before the module binds onto it. + if len(m.Consumes) > 0 { + if err := management.EnsureModuleQueue(ctx, *forNode, module); err != nil { + return err + } + } + + known, err := broker.FromEnvironment() + if err != nil { + return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err) + } + // The URL and what verifies the broker, together — a mesh's broker presents its own + // certificate, in no public trust store, so a URL alone fails at TLS (as `builder issue`). + held, err := json.Marshal(struct { + URL string `json:"url"` + Fingerprint string `json:"fingerprint,omitempty"` + Node string `json:"node"` + Module string `json:"module"` + }{ + URL: fmt.Sprintf("amqps://%s:%s@%s/", account, password, known.Address), + Fingerprint: known.Fingerprint, + // The node and module the account is for, so the runtime names its queue as the mesh + // scoped it (..events) without a manifest having to interpolate a node. + Node: *forNode, + Module: module, + }) + if err != nil { + return err + } + if err := inv.AcceptSecretForModule(ctx, *forNode, module, "broker", string(held)); err != nil { + return err + } + fmt.Printf("broker account %s created for %s, scoped to what it emits and consumes\n", + account, module) + fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n", + *forNode, *forNode) + return nil + default: - return fmt.Errorf("module has no %q; it has add, list, moved and forget", args[0]) + return fmt.Errorf("module has no %q; it has add, list, moved, forget and issue", args[0]) } } diff --git a/internal/broker/management.go b/internal/broker/management.go index 2873408..3270826 100644 --- a/internal/broker/management.go +++ b/internal/broker/management.go @@ -67,6 +67,119 @@ 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 +} + +// EnsureModuleQueue declares a consuming module's queue with its dead-letter exchange, idempotently. +// The substrate declares it because a scoped module account may not: the broker refuses a queue with +// a dead-letter exchange to a non-administrator (novox/hq ADR 0048), so a consumer passively checks +// the queue the mesh made rather than declaring its own. +func (m *Management) EnsureModuleQueue(ctx context.Context, node, module string) error { + queue := ModuleQueueFor(node, module) + if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(queue), map[string]any{ + "durable": true, + "arguments": map[string]any{"x-dead-letter-exchange": DeadExchangeName}, + }); err != nil { + return fmt.Errorf("cannot declare the queue for %s on %s: %w", module, node, err) + } + return 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 +280,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 +}