Files
Behavision/agent/pkg/mqtt/client.go
Suriyakumarvijayanayagam 2cd7a78ddc A headless PC can be claimed, and a refused broker says so
Both found by operating the stack rather than writing it: the local
processes were OOM-killed and bringing them back hit two gaps.

The headless agent had no way to be claimed at all. Bootstrap lived only
in desktop/internal/cloud, so the one configuration the agent binary
exists for - a back-office PC with no window - could only be onboarded by
hand-editing agent.json, which is the state the desktop's Setup screen was
built to end. `behavision-agent claim <code>` closes it; the CLI joins its
arguments because the code is printed in groups for reading aloud and an
operator pasting it will paste the spaces too.

Second: after the site's broker password was re-rolled, mosquitto logged
"not authorised" while the agent logged "timed out". Those need opposite
actions - re-link this PC, or go and look at the network - and paho's
SetConnectRetry collapses them, because it retries internally and the
connect token never completes. describeStall asks whether a TCP socket
opens at all, and says what is known rather than guessing at a reason the
broker never gives.

Verified end to end: minted a code from the platform as the owner,
claimed with the new command, broker connected, and the shop went to
online: true with 1/1 cameras on w600k_r50.

Also corrects this machine's memory in CLAUDE.md from 16 GB to 8 GB. It
feeds the model-fallback reasoning, and the local gallery already holds
17 embeddings tagged w600k_mbf beside 19 tagged w600k_r50 - the fallback
has silently fired before.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
2026-09-04 13:32:38 +05:30

267 lines
9.0 KiB
Go

// 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"
"net"
"net/url"
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) {
// SetConnectRetry means paho retries internally and this token never
// completes, so a REFUSED connection and an UNREACHABLE broker both
// arrive here as a timeout. They need opposite actions - re-link this
// PC, or go and look at the network - and reporting both as "timed
// out" sent the diagnosis to the wrong place. Measured: mosquitto
// logged "not authorised" while the agent logged a timeout.
return c, fmt.Errorf("mqtt: %s", describeStall(opts.BrokerURL))
}
if err := tok.Error(); err != nil {
return c, fmt.Errorf("mqtt: connect to %s: %w", opts.BrokerURL, err)
}
return c, nil
}
// describeStall says which of the two failures this is, by asking the one
// question that separates them: can we open a socket to the broker at all?
//
// It cannot name the exact reason - the broker does not tell a rejected client
// why, and a TLS failure looks the same from here - so it says what is known
// and what to check, rather than guessing. Being reachable but not accepted is
// overwhelmingly a credential this PC no longer has, which is what happens when
// a site is re-provisioned.
func describeStall(brokerURL string) string {
host := brokerHostPort(brokerURL)
if host == "" {
return fmt.Sprintf("connect to %s timed out", brokerURL)
}
conn, err := net.DialTimeout("tcp", host, 5*time.Second)
if err != nil {
return fmt.Sprintf("cannot reach the broker at %s: %v - check the "+
"network and that the broker is running", host, err)
}
_ = conn.Close()
return fmt.Sprintf("the broker at %s is reachable but did not accept this "+
"PC - usually its credentials are no longer valid; re-link it with "+
"`behavision-agent claim <code>`", host)
}
// brokerHostPort extracts host:port for the reachability probe. Parsed with
// net/url, never by scanning for the first ":" - an IPv6 literal is bracketed
// and full of them.
func brokerHostPort(brokerURL string) string {
u, err := url.Parse(brokerURL)
if err != nil || u.Host == "" {
return ""
}
if u.Port() != "" {
return u.Host
}
if strings.HasPrefix(brokerURL, "tls://") || strings.HasPrefix(brokerURL, "ssl://") {
return net.JoinHostPort(u.Hostname(), "8883")
}
return net.JoinHostPort(u.Hostname(), "1883")
}
// 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