The broker settings take a _FILE twin like the store connections; the catalogue engine refuses a secret placeholder in a container's env and a secret-carrying env-file unless the container says why with secrets-in-environment, which stays in the catalogue and never reaches the machine.
364 lines
16 KiB
Go
364 lines
16 KiB
Go
package broker
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"github.com/novox/mesh-controller/internal/envfile"
|
|
"io"
|
|
"net/http"
|
|
"net/url"
|
|
"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, err := envfile.Value(ManagementVar)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
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 0042): one topic exchange every event rides, a second for tool
|
|
// RPC kept apart, and a dead-letter home for a poison event. The foundation 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 0042), 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 0043): 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 0042's
|
|
// origin reservation — a module publishes only under `module.<self>.*` — 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 0047) — serve.<module>.<tool> — 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 0047: 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 0043). 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 foundation 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 0043), 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
|
|
// foundation 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 + "$",
|
|
// Two exchanges, 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 node exchange; announcing what was built goes through
|
|
// the events exchange, which is a different act with a different audience (novox/hq
|
|
// ADR 0072). A builder that could answer and not announce would leave the module graph
|
|
// knowing less than the registry does.
|
|
"write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + regexp.QuoteMeta(EventsExchangeName) + ")$",
|
|
// 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
|
|
}
|