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
267 lines
9.0 KiB
Go
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
|