// Package broker creates and removes a site's broker login at runtime. // // Until this existed a shop was created in two places by two mechanisms: // `provision site` wrote the row and printed a password, and a person then // typed that password into Mosquitto's passwd file on the host and reloaded // the broker. Nothing else in the product needed a shell, so this one step // was what stopped a tenant opening a second branch on their own - and it was // fragile even for us: the file was mounted read-only in the container, the // first attempt failed, and the password had to be re-rolled. // // Mosquitto 2.0's dynamic-security plugin takes the same operations as // commands on a control topic, from a client that holds the `admin` role. // The server already holds a broker login; this gives it that role and uses // it. No file, no reload, no docker socket, and the per-site credential // model is unchanged - one username per shop, topics only under its own // prefix. // // Isolation is a ROLE PER SITE with literal topics, not one role with `%u`: // the 2.0 plugin does not substitute `%u` in ACL topics (measured - the // publish was denied). The role is created and deleted with the client, so // there is still exactly one thing to get right and it is done in one place. package broker import ( "context" "crypto/rand" "encoding/hex" "encoding/json" "errors" "fmt" "log" "strings" "sync" "time" paho "github.com/eclipse/paho.mqtt.golang" ) const ( controlTopic = "$CONTROL/dynamic-security/v1" responseTopic = "$CONTROL/dynamic-security/v1/response" // A control round trip on a healthy broker is milliseconds; this is a // bound on a broker that is up but not answering control commands, which // is what a broker WITHOUT the plugin looks like. commandTimeout = 8 * time.Second connectTimeout = 10 * time.Second ) // ErrUnavailable is returned when the broker cannot be reached or does not // answer control commands. Callers turn it into a 502/503, never a 500: the // operator's next step is to look at the broker, not at the server. var ErrUnavailable = errors.New("broker: dynamic security not available") // SiteBroker is what the API and the provisioner depend on. A fake satisfies // it in tests; Dynsec satisfies it in production. type SiteBroker interface { // EnsureSite creates the broker login for one site, or resets its password // if it already exists. Idempotent: safe to run again after any failure. EnsureSite(ctx context.Context, username, password string) error // DeleteSite removes the login and its role. Missing is not an error. DeleteSite(ctx context.Context, username string) error } // Dynsec drives the plugin over its own connection - not the ingest client's. // That one has SetOrderMatters and blocking handlers, and a provisioning call // must neither wait behind a slow visit nor delay one. type Dynsec struct { url, user, pass string log *log.Logger mu sync.Mutex client paho.Client pending map[string]chan []response // keyed by batch id } func New(url, user, pass string, logger *log.Logger) *Dynsec { if logger == nil { logger = log.Default() } return &Dynsec{url: url, user: user, pass: pass, log: logger, pending: map[string]chan []response{}} } type response struct { Command string `json:"command"` Error string `json:"error,omitempty"` CorrelationData string `json:"correlationData,omitempty"` Data json.RawMessage `json:"data,omitempty"` } // SiteTopics are exactly what a shop PC may do: publish its own visits, // heartbeats and status, and read its own commands. This mirrors the acl // file the broker used to be configured with, rule for rule. func siteACLs(username string) []map[string]any { prefix := "bv/" + username acl := func(kind, topic string) map[string]any { return map[string]any{"acltype": kind, "topic": topic, "allow": true, "priority": 0} } return []map[string]any{ acl("publishClientSend", prefix+"/visit"), acl("publishClientSend", prefix+"/heartbeat"), acl("publishClientSend", prefix+"/status"), acl("publishClientReceive", prefix+"/cmd/#"), acl("subscribePattern", prefix+"/cmd/#"), } } // RoleName is the per-site role. Exported so the migration can name it the // same way. func RoleName(username string) string { return "site." + username } func (d *Dynsec) EnsureSite(ctx context.Context, username, password string) error { if strings.TrimSpace(username) == "" || password == "" { return errors.New("broker: username and password are required") } role := RoleName(username) cmds := []map[string]any{{"command": "createRole", "rolename": role}} for _, a := range siteACLs(username) { c := map[string]any{"command": "addRoleACL", "rolename": role} for k, v := range a { c[k] = v } cmds = append(cmds, c) } cmds = append(cmds, map[string]any{ "command": "createClient", "username": username, "password": password, "roles": []map[string]any{{"rolename": role}}, }) res, err := d.run(ctx, cmds) if err != nil { return err } clientExisted := false for _, r := range res { switch { case r.Error == "": case strings.Contains(r.Error, "already exists") && r.Command != "createClient": // A re-run after a partial failure. Fine. case r.Command == "createClient" && strings.Contains(r.Error, "already exists"): clientExisted = true default: return fmt.Errorf("broker: %s: %s", r.Command, r.Error) } } if !clientExisted { return nil } // The login exists from an earlier run: make its password THIS one - the // database holds this one and the enrolment hands it out - and make sure // it carries the role. addClientRole on a client that already has it // answers "Internal error", so check first rather than guess from prose. res, err = d.run(ctx, []map[string]any{ {"command": "setClientPassword", "username": username, "password": password}, {"command": "getClient", "username": username}, }) if err != nil { return err } hasRole := false for _, r := range res { if r.Error != "" { return fmt.Errorf("broker: %s: %s", r.Command, r.Error) } if r.Command == "getClient" { var got struct { Client struct { Roles []struct{ Rolename string } `json:"roles"` } `json:"client"` } _ = json.Unmarshal(r.Data, &got) for _, rr := range got.Client.Roles { if rr.Rolename == role { hasRole = true } } } } if hasRole { return nil } res, err = d.run(ctx, []map[string]any{{"command": "addClientRole", "username": username, "rolename": role}}) if err != nil { return err } if res[0].Error != "" { return fmt.Errorf("broker: addClientRole: %s", res[0].Error) } return nil } func (d *Dynsec) DeleteSite(ctx context.Context, username string) error { res, err := d.run(ctx, []map[string]any{ {"command": "deleteClient", "username": username}, {"command": "deleteRole", "rolename": RoleName(username)}, }) if err != nil { return err } for _, r := range res { if r.Error != "" && !strings.Contains(r.Error, "not found") { return fmt.Errorf("broker: %s: %s", r.Command, r.Error) } } return nil } // Close drops the control connection. Safe when never connected. func (d *Dynsec) Close() { d.mu.Lock() c := d.client d.client = nil d.mu.Unlock() if c != nil && c.IsConnected() { c.Disconnect(250) } } // run sends one batch and waits for its one response message. Every command // carries the batch id as correlationData, and the plugin echoes it, so a // response for somebody else's batch - two servers, or a retry - is never // mistaken for ours. func (d *Dynsec) run(ctx context.Context, cmds []map[string]any) ([]response, error) { if err := d.connect(ctx); err != nil { return nil, err } id := newID() for _, c := range cmds { c["correlationData"] = id } body, err := json.Marshal(map[string]any{"commands": cmds}) if err != nil { return nil, err } ch := make(chan []response, 1) d.mu.Lock() d.pending[id] = ch client := d.client d.mu.Unlock() defer func() { d.mu.Lock() delete(d.pending, id) d.mu.Unlock() }() tok := client.Publish(controlTopic, 1, false, body) if !tok.WaitTimeout(commandTimeout) { return nil, fmt.Errorf("%w: publish timed out", ErrUnavailable) } if tok.Error() != nil { return nil, fmt.Errorf("%w: %v", ErrUnavailable, tok.Error()) } select { case res := <-ch: return res, nil case <-time.After(commandTimeout): return nil, fmt.Errorf("%w: no answer on %s - is the dynamic-security plugin enabled and does %q hold the admin role?", ErrUnavailable, responseTopic, d.user) case <-ctx.Done(): return nil, ctx.Err() } } func (d *Dynsec) onResponse(_ paho.Client, m paho.Message) { var env struct { Responses []response `json:"responses"` } if err := json.Unmarshal(m.Payload(), &env); err != nil || len(env.Responses) == 0 { return } id := env.Responses[0].CorrelationData d.mu.Lock() ch, ok := d.pending[id] d.mu.Unlock() if ok { select { case ch <- env.Responses: default: } } } func (d *Dynsec) connect(ctx context.Context) error { d.mu.Lock() defer d.mu.Unlock() if d.client != nil && d.client.IsConnectionOpen() { return nil } opts := paho.NewClientOptions(). AddBroker(d.url). SetClientID(fmt.Sprintf("behavision-dynsec-%s", newID()[:8])). SetUsername(d.user). SetPassword(d.pass). SetAutoReconnect(true). SetCleanSession(true). SetKeepAlive(30 * time.Second). SetConnectTimeout(connectTimeout) opts.OnConnect = func(c paho.Client) { // Re-subscribed on every (re)connect: clean session keeps nothing. c.Subscribe(responseTopic, 1, d.onResponse) } c := paho.NewClient(opts) tok := c.Connect() if !tok.WaitTimeout(connectTimeout) { return fmt.Errorf("%w: connect to %s timed out", ErrUnavailable, d.url) } if tok.Error() != nil { return fmt.Errorf("%w: %v", ErrUnavailable, tok.Error()) } // The subscription must be in place before the first command is sent, or // its answer is published to nobody. st := c.Subscribe(responseTopic, 1, d.onResponse) if !st.WaitTimeout(connectTimeout) || st.Error() != nil { c.Disconnect(100) return fmt.Errorf("%w: cannot subscribe to %s (does %q hold the admin role?)", ErrUnavailable, responseTopic, d.user) } d.client = c return nil } func newID() string { var b [12]byte _, _ = rand.Read(b[:]) return hex.EncodeToString(b[:]) }