package broker import ( "bytes" "context" "encoding/json" "fmt" "io" "net/http" "net/url" "os" "regexp" "strings" "time" ) // The mesh runs the broker, so there is no chicken-and-egg in a node needing an account before it // can connect: the account is created when the token is issued, and the one-time secret in that // token IS the password. A node's first connection is already authenticated, and enrolment is // what happens over it. // // novox/hq ADR 0004's *a node holds its own identity and nothing else* is why the account is per // node rather than shared. A shared enrolment account would let any node consume another's queue, // which is the shared-credential fault that record exists to remove, reappearing at the transport. // ManagementVar holds the broker's management API, credentials included. const ManagementVar = "MESH_BROKER_MANAGEMENT" // safeName is what a node may be called at the broker. // // The name goes into a URL path and into permission patterns, which are regular expressions. A // name carrying a `.` or a `*` would silently widen what that node may reach — so it is // constrained here rather than escaped later, because an escape that is forgotten once is a node // reading everybody's queues. var safeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,62}$`) // Management is the broker's administrative interface. type Management struct { base *url.URL client *http.Client } // ManagementFromEnvironment reads where the management API is, if it is configured. func ManagementFromEnvironment() (*Management, error) { raw := strings.TrimSpace(os.Getenv(ManagementVar)) if raw == "" { return nil, ErrNotConfigured } base, err := url.Parse(raw) if err != nil || base.Host == "" { // The value carries a password, so it is not quoted back. return nil, fmt.Errorf("%s is not a usable URL", ManagementVar) } return &Management{base: base, client: &http.Client{Timeout: 15 * time.Second}}, nil } // QueueFor is the queue a node consumes from. One per node, named after it. func QueueFor(node string) string { return "node." + node } // ExchangeName is where nodes publish what they have to say. One exchange, and the control plane // is the only consumer behind it (novox/hq ADR 0006 — one consumer, so two cannot silently split // the traffic between them). const ExchangeName = "mesh" // BuildQueueName is where build work waits. Duplicated from `link` rather than imported, for the // same reason QueueFor above is: this package must not depend on the one that uses it, and a // 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) rpc := regexp.QuoteMeta(RPCExchangeName) // A module serves each of its tools on its own queue, namespaced by the module (novox/hq // ADR 0052) — serve.. — so the account may declare, bind and read exactly its own, // and no other module's. serve := "serve\\." + regexp.QuoteMeta(module) + "\\..*" // Declare its own events queue and its own tool serve queues. configure = "^(" + queue + "|" + serve + ")$" // Write to bind its queue and serve queues (binding is a write on the queue), and to the RPC // exchange to publish replies (ADR 0052: replies ride mesh.rpc, never the default exchange, which // would let it publish into any queue). To the events exchange only if it emits. writes := []string{queue, serve, rpc} if len(emits) > 0 { writes = append(writes, events) } write = "^(" + strings.Join(writes, "|") + ")$" // Read its own queue and serve queues to consume them, and the RPC exchange to bind its serve // queues onto. The events exchange to bind onto only if it consumes. reads := []string{queue, serve, rpc} 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 // are anchored: a node called `laptop` must not be able to read `laptop-of-somebody-else`. func (m *Management) CreateNodeAccount(ctx context.Context, node, password string) error { if !safeName.MatchString(node) { return fmt.Errorf( "%q cannot be a broker account name: it becomes part of a permission pattern, so it "+ "is lower-case letters, digits and dashes", node) } if err := m.put(ctx, "/api/users/"+url.PathEscape(node), map[string]string{"password": password, "tags": ""}); err != nil { return fmt.Errorf("cannot create the broker account for %s: %w", node, err) } queue := regexp.QuoteMeta(QueueFor(node)) if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(node), map[string]string{ "configure": "^" + queue + "$", "write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + queue + ")$", "read": "^" + queue + "$", }); err != nil { return fmt.Errorf("cannot scope the broker account for %s: %w", node, err) } return nil } // CreateBuilderAccount scopes an account to taking build work and answering it. // // **A build machine is not a node**, and giving it a node's account would let it read another // machine's declarations. What it needs is narrower and different: read the build queue, and // write to the exchange and to whatever temporary queue an asker is waiting on. // // The reply queues are the reason `write` is not simply the exchange. `RequestBuild` declares an // exclusive queue with a generated name and waits on it, so a builder that could not write to it // could take work and never answer — which is the failure that looks like a builder that is not // running. func (m *Management) CreateBuilderAccount(ctx context.Context, name, password string) error { if !safeName.MatchString(name) { return fmt.Errorf( "%q cannot be a broker account name: it becomes part of a permission pattern, so it "+ "is lower-case letters, digits and dashes", name) } if err := m.put(ctx, "/api/users/"+url.PathEscape(name), map[string]string{"password": password, "tags": ""}); err != nil { return fmt.Errorf("cannot create the broker account for %s: %w", name, err) } builds := regexp.QuoteMeta(BuildQueueName) if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(name), map[string]string{ // It declares the build queue, because whichever builder starts first must be able to — // and a queue nobody may declare is a queue that exists only if the control plane has // already run, which makes the order they start in matter. "configure": "^" + builds + "$", // The exchange, and nothing else. **Not the default exchange**: permission there is // granted per exchange rather than per queue, so a builder allowed to use it could // publish into any node's queue — the privilege a build machine most obviously should // not have. Answers go through the exchange, which is why they can. "write": "^" + regexp.QuoteMeta(ExchangeName) + "$", // The build queue and nothing else. Not another machine's declarations. "read": "^" + builds + "$", }); err != nil { return fmt.Errorf("cannot scope the broker account for %s: %w", name, err) } return nil } // RemoveNodeAccount withdraws a node's access. func (m *Management) RemoveNodeAccount(ctx context.Context, node string) error { if !safeName.MatchString(node) { return fmt.Errorf("%q is not a broker account name", node) } return m.do(ctx, http.MethodDelete, "/api/users/"+url.PathEscape(node), nil) } // Accounts lists the broker's users, so a picture can be read from the system rather than assumed // (novox/hq ADR 0018). func (m *Management) Accounts(ctx context.Context) ([]string, error) { body, err := m.get(ctx, "/api/users") if err != nil { return nil, err } var users []struct { Name string `json:"name"` } if err := json.Unmarshal(body, &users); err != nil { return nil, err } names := make([]string, 0, len(users)) for _, u := range users { names = append(names, u.Name) } return names, nil } 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 { return nil, err } response, err := m.client.Do(request) if err != nil { return nil, err } defer response.Body.Close() if response.StatusCode >= 300 { return nil, fmt.Errorf("the broker's management API answered %s to GET %s", response.Status, path) } return io.ReadAll(io.LimitReader(response.Body, 1<<20)) } func (m *Management) do(ctx context.Context, method, path string, body any) error { request, err := m.request(ctx, method, path, body) if err != nil { return err } response, err := m.client.Do(request) if err != nil { return err } defer response.Body.Close() if response.StatusCode >= 300 { detail, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) return fmt.Errorf("the broker's management API answered %s to %s %s: %s", response.Status, method, path, strings.TrimSpace(string(detail))) } return nil } func (m *Management) request(ctx context.Context, method, path string, body any) (*http.Request, error) { var payload io.Reader if body != nil { raw, err := json.Marshal(body) if err != nil { return nil, err } payload = bytes.NewReader(raw) } // Path joined by hand rather than through url.Parse: %2f is the default vhost and must reach // the broker still encoded. Parsing would decode it to a slash and address a different route. target := strings.TrimSuffix(m.base.String(), "/") if user := m.base.User; user != nil { target = strings.TrimSuffix(m.base.Scheme+"://"+m.base.Host, "/") } request, err := http.NewRequestWithContext(ctx, method, target+path, payload) if err != nil { return nil, err } if user := m.base.User; user != nil { password, _ := user.Password() request.SetBasicAuth(user.Username(), password) } if body != nil { request.Header.Set("Content-Type", "application/json") } return request, nil }