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" // 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) 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 }