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
208 lines
5.8 KiB
Go
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...)
|
|
}
|
|
}
|