Files
Behavision/agent/pkg/mqtt/pump.go
Suriyakumarvijayanayagam 50d122e5f0 An unused constant was a Stop() that could hang forever
Ran staticcheck across all three Go modules for the first time. server
(23k lines) and desktop came back clean. agent had seven findings, and
one of them was not tidiness.

`stopGrace = 10 * time.Second` was declared and wired to nothing.
Stop() cancels the context, cmd.Cancel kills the process tree, and then
Stop() blocks on cmd.Wait() - which, with no WaitDelay set, waits not
just for the process but for every writer of its stdout pipe to close.
One grandchild still holding that pipe hangs Wait, hangs Stop, and on the
desktop app that is the tray's Quit never returning. The constant named
the intent and nothing read it. cmd.WaitDelay = stopGrace is the line
that was missing.

The rest were real but small: an unused field in the live relay, an
unused sleep helper in the pump, and "net/url" imported twice under two
names - both genuinely used, in two functions doing the same job for the
same reason, so they are unified rather than one deleted. My first pass
deleted the wrong one on a bad grep and the build caught it immediately.

Three findings are suppressed rather than fixed, with the reason stated:

- Two "error strings should not end with punctuation". Both are
  multi-line messages a shop operator reads at a counter, not errors
  anything wraps. ST1005 exists because wrapped errors concatenate
  mid-sentence; stripping the full stops would run three sentences
  together to satisfy a rule that does not apply.
- A deliberately nil context in a pump test - the point of the test is
  that an unconnected client does not panic. It already carried
  //nolint:staticcheck, which is golangci-lint's directive and
  staticcheck ignores, which is why it kept being reported.

Also tidied agent/go.mod, which had paho and x/sys marked indirect while
being imported directly.

All three modules clean, all suites pass: 21 Go packages, 226 engine
tests.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-30 15:14:09 +05:30

208 lines
5.8 KiB
Go

// Package mqtt moves queued events to the broker.
//
// Split from the broker client on purpose: everything that decides *what to
// send and when* lives here and is testable without a broker, while the paho
// binding is a thin adapter that only knows how to put bytes on a topic.
package mqtt
import (
"context"
"errors"
"log"
"time"
"github.com/loyaly/behavision-agent/pkg/spool"
)
// Publisher is the broker, reduced to what the pump needs.
type Publisher interface {
// Publish must return nil only once the broker has confirmed receipt.
// Returning early would let the pump ack an event that never arrived.
Publish(ctx context.Context, topic string, payload []byte) error
Connected() bool
}
// Queue is the durable side, reduced likewise.
type Queue interface {
Peek(n int) ([]spool.Entry, error)
Ack(seqs ...uint64) error
Len() int
Dropped() uint64
}
const (
batchSize = 32
idleInterval = 2 * time.Second
minRetry = 1 * time.Second
maxRetry = 30 * time.Second
defaultHeartbe = 60 * time.Second
)
// Pump drains the queue into the broker and emits a heartbeat.
type Pump struct {
Queue Queue
Publisher Publisher
// Heartbeat topic. Without it "the site is offline" and "nobody visited"
// are indistinguishable on the server, which for a footfall product is a
// silent hole in the customer's report.
HeartbeatTopic string
HeartbeatPayload func() []byte
HeartbeatInterval time.Duration
Log *log.Logger
// Wake, when set, makes the pump drain immediately instead of waiting out
// idleInterval. Without it a visit that lands one millisecond after a drain
// sits on disk for two seconds before anyone is told - and that delay is on
// the path a shop screen or a mobile app sees as "how long after someone
// walks in does their face appear".
//
// A doorbell, not a queue: it carries nothing, because the pump re-reads
// the spool either way. Buffered by one and written non-blockingly, so a
// burst of arrivals cannot stall the recognition pipeline behind a pump
// that is mid-publish.
Wake <-chan struct{}
}
// Waker is the writing end of the Wake channel, held by whatever appends to the
// queue. NewWaker returns both halves so a caller cannot accidentally build one
// that blocks its own producer.
type Waker struct{ ch chan struct{} }
func NewWaker() *Waker { return &Waker{ch: make(chan struct{}, 1)} }
// Wake rings the pump. Never blocks: a full slot already means "there is work",
// which is the entire message, so a second ring adds nothing.
func (w *Waker) Wake() {
if w == nil {
return
}
select {
case w.ch <- struct{}{}:
default:
}
}
// C is the channel to hand the pump.
func (w *Waker) C() <-chan struct{} {
if w == nil {
return nil
}
return w.ch
}
// Run drains until ctx is cancelled.
func (p *Pump) Run(ctx context.Context) {
interval := p.HeartbeatInterval
if interval <= 0 {
interval = defaultHeartbe
}
beat := time.NewTicker(interval)
defer beat.Stop()
retry := minRetry
for {
if ctx.Err() != nil {
return
}
// Non-blocking, for the case where there is a backlog and the loop
// never reaches the waiting select below.
select {
case <-beat.C:
p.heartbeat(ctx)
default:
}
sent, err := p.drainOnce(ctx)
if ctx.Err() != nil {
return
}
var wait time.Duration
switch {
case err != nil:
// The broker is down or refusing. Back off rather than spinning:
// a store with no internet would otherwise burn a core all night.
p.logf("publish failed, retrying in %s: %v", retry, err)
wait = retry
if retry < maxRetry {
retry *= 2
if retry > maxRetry {
retry = maxRetry
}
}
case sent == 0:
retry = minRetry
wait = idleInterval
default:
// Something went through; there may be more waiting, so loop
// immediately rather than sleeping through a backlog.
retry = minRetry
}
if wait == 0 {
continue
}
// The heartbeat must be able to interrupt this wait. Sleeping through
// it would delay every beat by the idle interval, and on a quiet site
// the pump is idle essentially always.
// A nil Wake channel blocks forever in a select, which is exactly the
// right behaviour: an agent with no waker falls back to the timer.
select {
case <-ctx.Done():
return
case <-beat.C:
p.heartbeat(ctx)
case <-p.Wake:
// Something was queued. Loop straight round and drain it rather
// than sleeping out the rest of the idle interval.
case <-time.After(wait):
}
}
}
// drainOnce sends at most one batch and returns how many were acked.
func (p *Pump) drainOnce(ctx context.Context) (int, error) {
if !p.Publisher.Connected() {
return 0, errors.New("broker not connected")
}
entries, err := p.Queue.Peek(batchSize)
if err != nil || len(entries) == 0 {
return 0, err
}
sent := 0
for _, e := range entries {
if err := p.Publisher.Publish(ctx, e.Topic, e.Payload); err != nil {
// Stop at the first failure instead of skipping past it. Events
// are a per-visitor timeline and the server reads them in order;
// publishing around a stuck one would reorder a customer's visits.
return sent, err
}
// Acked one at a time, immediately after its own confirmation. A batch
// ack would re-send everything before a mid-batch failure on restart.
if err := p.Queue.Ack(e.Seq); err != nil {
return sent, err
}
sent++
}
return sent, nil
}
func (p *Pump) heartbeat(ctx context.Context) {
if p.HeartbeatTopic == "" || p.HeartbeatPayload == nil {
return
}
if !p.Publisher.Connected() {
return
}
// Not queued: a heartbeat is only meaningful now. Spooling them would
// replay a week of "I am alive" the moment a site reconnects.
if err := p.Publisher.Publish(ctx, p.HeartbeatTopic, p.HeartbeatPayload()); err != nil {
p.logf("heartbeat failed: %v", err)
}
}
func (p *Pump) logf(format string, args ...any) {
if p.Log != nil {
p.Log.Printf(format, args...)
}
}