// Broker client: the thin adapter behind the Publisher interface. // // Everything that decides *what to send and when* is in pump.go and is tested // without a broker. This file only knows how to put bytes on a topic, which is // why it is the one part that needs a real connection to exercise. // // Targets Mosquitto. No clustering, no shared subscriptions, no broker-side // rules — a store publishes its own events under its own prefix and that is // the whole interaction. package mqtt import ( "context" "crypto/tls" "crypto/x509" "errors" "fmt" "log" neturl "net/url" "os" "strings" "time" paho "github.com/eclipse/paho.mqtt.golang" ) // ClientOptions configures a broker connection. type ClientOptions struct { // BrokerURL is tls://host:8883 in production, tcp://host:1883 for local // testing only. Credentials and footfall must never cross the internet in // the clear, so Connect refuses tcp:// to a non-loopback host. BrokerURL string ClientID string Username string Password string // CAFile pins a private CA. Empty uses the system roots, which is what a // Let's Encrypt certificate on the broker needs. CAFile string // InsecureSkipVerify disables certificate checking. Only ever for a // self-signed staging box, and it is logged loudly when set, because a // forgotten one silently removes the protection TLS was added for. InsecureSkipVerify bool // PublishTimeout bounds a single publish. Without it a half-open // connection blocks the pump indefinitely and the queue grows behind it. PublishTimeout time.Duration Log *log.Logger } // Client implements Publisher. type Client struct { opts ClientOptions client paho.Client } // NewClient dials the broker. It returns as soon as the connection is // established; reconnection afterwards is automatic and the pump reads // Connected() to decide whether to try. func NewClient(opts ClientOptions) (*Client, error) { if opts.BrokerURL == "" { return nil, errors.New("mqtt: no broker url") } if opts.PublishTimeout <= 0 { opts.PublishTimeout = 10 * time.Second } if err := checkTransport(opts.BrokerURL); err != nil { return nil, err } po := paho.NewClientOptions(). AddBroker(opts.BrokerURL). SetClientID(opts.ClientID). SetUsername(opts.Username). SetPassword(opts.Password). // The broker holds no state for us: every event is already durable on // our own disk, so a clean session avoids the broker queueing a // second copy we would then have to de-duplicate. SetCleanSession(true). SetAutoReconnect(true). SetConnectRetry(true). SetConnectRetryInterval(5 * time.Second). SetMaxReconnectInterval(2 * time.Minute). SetKeepAlive(30 * time.Second). SetConnectTimeout(15 * time.Second). // Publishes must fail fast rather than pile up in memory while the // link is down; the spool is what holds them, not the client. SetMessageChannelDepth(1). SetOrderMatters(true) if strings.HasPrefix(opts.BrokerURL, "tls://") || strings.HasPrefix(opts.BrokerURL, "ssl://") { cfg, err := tlsConfig(opts) if err != nil { return nil, err } po.SetTLSConfig(cfg) } c := &Client{opts: opts} po.OnConnect = func(paho.Client) { c.logf("broker connected: %s", opts.BrokerURL) } po.OnConnectionLost = func(_ paho.Client, err error) { c.logf("broker connection lost: %v", err) } c.client = paho.NewClient(po) tok := c.client.Connect() if !tok.WaitTimeout(20 * time.Second) { return c, fmt.Errorf("mqtt: connect to %s timed out", opts.BrokerURL) } if err := tok.Error(); err != nil { return c, fmt.Errorf("mqtt: connect to %s: %w", opts.BrokerURL, err) } return c, nil } // Publish sends one message at QoS 1 and waits for the broker's PUBACK. // // QoS 1, not 0 or 2. At QoS 0 the broker never confirms, so the pump would ack // and delete an event that was dropped on the wire. QoS 2 costs two extra // round trips to remove a duplicate the server can drop itself from the event // id — at-least-once with idempotent consumers is the cheaper contract. func (c *Client) Publish(ctx context.Context, topic string, payload []byte) error { if c.client == nil { return errors.New("mqtt: no client") } if !c.client.IsConnected() { return errors.New("mqtt: not connected") } tok := c.client.Publish(topic, 1, false, payload) // Honour both the caller's context and a hard timeout: a half-open TCP // connection can leave a token that never completes, which would stall the // pump forever with the queue growing behind it. done := make(chan struct{}) go func() { tok.Wait(); close(done) }() select { case <-ctx.Done(): return ctx.Err() case <-time.After(c.opts.PublishTimeout): return fmt.Errorf("mqtt: publish to %s timed out", topic) case <-done: return tok.Error() } } // Connected reports whether the broker link is up. func (c *Client) Connected() bool { return c.client != nil && c.client.IsConnected() } // Close disconnects cleanly, giving in-flight publishes a moment to land. func (c *Client) Close() { if c.client != nil && c.client.IsConnected() { c.client.Disconnect(1000) } } func (c *Client) logf(format string, args ...any) { if c.opts.Log != nil { c.opts.Log.Printf(format, args...) } } // checkTransport refuses plaintext MQTT to anywhere but the local machine. // // The payloads carry customer visit records and the connection carries the // tenant's broker password. A tcp:// URL to a public host is not a // configuration choice, it is a mistake, and it is one that works — which is // exactly why it has to be rejected here rather than noticed later. func checkTransport(raw string) error { if !strings.HasPrefix(raw, "tcp://") && !strings.HasPrefix(raw, "mqtt://") { return nil } if os.Getenv("BEHAVISION_ALLOW_PLAINTEXT_MQTT") == "1" { return nil } // url.Parse, not hand-rolled splitting: an IPv6 literal is bracketed and // full of colons, so scanning for the first ":" turns "[::1]:1883" into // "[" and refuses a perfectly good loopback address. u, err := neturl.Parse(raw) if err != nil { return fmt.Errorf("mqtt: cannot parse broker url %q: %w", raw, err) } switch u.Hostname() { case "localhost", "127.0.0.1", "::1", "": return nil } return fmt.Errorf("mqtt: refusing plaintext connection to %q - use tls:// "+ "(set BEHAVISION_ALLOW_PLAINTEXT_MQTT=1 only for local testing)", u.Hostname()) } func tlsConfig(opts ClientOptions) (*tls.Config, error) { cfg := &tls.Config{MinVersion: tls.VersionTLS12} if opts.InsecureSkipVerify { cfg.InsecureSkipVerify = true if opts.Log != nil { opts.Log.Print("WARNING: MQTT certificate verification is DISABLED") } return cfg, nil } if opts.CAFile == "" { return cfg, nil // system roots } pem, err := os.ReadFile(opts.CAFile) if err != nil { return nil, fmt.Errorf("mqtt: ca file: %w", err) } pool := x509.NewCertPool() if !pool.AppendCertsFromPEM(pem) { return nil, fmt.Errorf("mqtt: no certificates found in %s", opts.CAFile) } cfg.RootCAs = pool return cfg, nil } // osWriteFile is indirected so tests can build without importing os twice. var osWriteFile = os.WriteFile