// 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...) } }