Files
backend_fiesta/messaging/livehub.go
2026-08-27 11:17:17 +05:30

349 lines
11 KiB
Go

package messaging
import (
"encoding/json"
"log"
"os"
"strconv"
"strings"
"sync"
"time"
"nearle/models"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
// Live console events.
//
// The console polls every 30 seconds. That is a safe floor, and it means a bill
// rung at 10:00:01 shows up on the owner's screen at 10:00:30. This closes that
// gap: the tills already publish every sale to the broker, so a second consumer
// listens to the same stream and pushes a nudge to whichever consoles are
// watching that outlet.
//
// Two things it deliberately is NOT:
//
// - **It does not carry data.** An event says "stock at outlet 1185 changed,
// these product ids are involved" and nothing more. The console then re-reads
// through the normal API, which is the only thing that knows the authoritative
// number after locks, dedup and rejections. Pushing figures from here would
// mean a screen showing a total the database never agreed to.
// - **It does not replace the poll.** The console keeps its 30-second refetch.
// A dropped connection, a full buffer or a replica restart then costs latency
// rather than correctness, which is the trade an operations screen wants.
//
// ── Why a second consumer and not a hook in the ingest path ──────────────────
//
// `StartPosMqttConsumer` runs on ONE replica (see `posConsumerElected`) because
// two writers ingesting the same bills would fight over the same rows. Hanging
// this off that consumer would mean only consoles that happened to land on the
// elected replica ever received anything, which is a bug that looks like flaky
// network. This client only reads and touches no database, so every replica can
// safely run one and serve its own connected consoles.
//
// Enabled by MQTT_URL, the same switch the ingest consumer uses. Unset, the hub
// still runs and simply never emits — the console degrades to polling, which is
// exactly what it does today.
const (
// How many events a single console connection may fall behind before it
// starts losing them. Small on purpose: these are nudges, and a client that
// cannot keep up with a handful is better served by its next poll than by a
// backlog of stale ones.
liveBuffer = 16
// Belt and braces around a burst. A terminal replaying a day of offline
// bills would otherwise emit one event per batch as fast as the broker can
// deliver them; the console cannot use more than a couple a second.
liveMinInterval = 400 * time.Millisecond
)
// LiveEvent is one nudge. Kept small — it crosses the wire on every sale.
type LiveEvent struct {
// "sale", "customer" or "terminal".
Type string `json:"type"`
// tenantlocations.locationid. The topic's store segment carries it as a
// string (see models.PosOrderBatch).
Locationid int `json:"locationid"`
// Which products moved, when the event knows. Empty means "something did".
Productids []int `json:"productids,omitempty"`
// Which till, for the presence board.
Terminalid string `json:"terminalid,omitempty"`
At string `json:"at"`
}
type liveSubscriber struct {
ch chan LiveEvent
last time.Time
}
// LiveHub fans events out to the consoles currently watching an outlet.
type LiveHub struct {
mu sync.RWMutex
subscribers map[int]map[*liveSubscriber]struct{}
client mqtt.Client
dropped uint64
}
// Hub is the process-wide hub. Nil until StartLiveHub runs; every method
// tolerates a nil receiver so callers never have to check.
var Hub *LiveHub
// StartLiveHub creates the hub and, if a broker is configured, connects a
// read-only consumer to it.
//
// Never returns nil: a hub with no broker is a working hub with no events, and
// the SSE endpoint must still accept connections so the console's reconnect
// logic has something to talk to.
func StartLiveHub() *LiveHub {
hub := &LiveHub{subscribers: make(map[int]map[*liveSubscriber]struct{})}
Hub = hub
url := strings.TrimSpace(os.Getenv("MQTT_URL"))
if url == "" {
log.Println("live: MQTT_URL not set, console event stream will stay quiet")
return hub
}
// A distinct client id per replica. Sharing one with the ingest consumer
// would make the broker disconnect whichever connected first.
clientID := "nearle-console-live-" + hostSuffix()
opts := mqtt.NewClientOptions().
AddBroker(url).
SetClientID(clientID).
SetAutoReconnect(true).
SetConnectRetry(true).
SetConnectRetryInterval(5 * time.Second).
SetCleanSession(true). // No queued backlog on reconnect; stale nudges are noise.
SetOrderMatters(false)
if user := strings.TrimSpace(os.Getenv("MQTT_USERNAME")); user != "" {
opts.SetUsername(user)
opts.SetPassword(os.Getenv("MQTT_PASSWORD"))
}
opts.OnConnect = func(client mqtt.Client) {
for topic, handler := range map[string]mqtt.MessageHandler{
topicOrders: hub.onOrders,
topicCustomers: hub.onCustomers,
topicHealth: hub.onHealth,
} {
// QoS 0. A missed nudge costs at most one poll interval, and QoS 1
// would have the broker retaining state for a listener that does
// not need it.
if token := client.Subscribe(topic, 0, handler); token.Wait() && token.Error() != nil {
log.Printf("live: could not subscribe to %s: %v", topic, token.Error())
continue
}
log.Printf("live: watching %s", topic)
}
}
client := mqtt.NewClient(opts)
hub.client = client
// Connect in the background: a broker that is slow to answer must not hold
// up the HTTP server coming online.
go func() {
if token := client.Connect(); token.Wait() && token.Error() != nil {
log.Printf("live: broker unreachable, console falls back to polling: %v", token.Error())
}
}()
return hub
}
/* ── Subscription ─────────────────────────────────────────────────────────── */
// Subscribe registers a console connection watching one outlet.
//
// Returns the channel to read and the function that releases it. The caller
// MUST call release — an SSE handler that returns without it leaks a channel
// the hub goes on writing to for the life of the process.
func (h *LiveHub) Subscribe(locationid int) (<-chan LiveEvent, func()) {
if h == nil {
// A closed channel reads immediately and forever, which would spin the
// handler. An open one that never delivers is the honest no-op.
return make(chan LiveEvent), func() {}
}
sub := &liveSubscriber{ch: make(chan LiveEvent, liveBuffer)}
h.mu.Lock()
if h.subscribers[locationid] == nil {
h.subscribers[locationid] = make(map[*liveSubscriber]struct{})
}
h.subscribers[locationid][sub] = struct{}{}
h.mu.Unlock()
var once sync.Once
release := func() {
once.Do(func() {
h.mu.Lock()
if set := h.subscribers[locationid]; set != nil {
delete(set, sub)
if len(set) == 0 {
delete(h.subscribers, locationid)
}
}
h.mu.Unlock()
close(sub.ch)
})
}
return sub.ch, release
}
// Broadcast delivers an event to everyone watching that outlet.
//
// Never blocks. A subscriber whose buffer is full loses the event and is
// counted rather than waited for: one wedged console must not stall the fan-out
// to every other console, and the poll will collect what was missed.
func (h *LiveHub) Broadcast(event LiveEvent) {
if h == nil || event.Locationid == 0 {
return
}
if event.At == "" {
event.At = time.Now().UTC().Format(time.RFC3339)
}
// A write lock, not a read lock, even though this only walks the map: the
// loop below writes `sub.last` and `h.dropped`. Under RLock those are
// concurrent writes from every broker goroutine at once — a data race, and
// one `go test -race` would fail on. Fan-out is a handful of non-blocking
// sends, so holding the write lock costs nothing worth reclaiming.
h.mu.Lock()
defer h.mu.Unlock()
now := time.Now()
for sub := range h.subscribers[event.Locationid] {
// Coalesce a burst per subscriber, not globally: a busy outlet must not
// throttle a quiet one.
if now.Sub(sub.last) < liveMinInterval {
continue
}
select {
case sub.ch <- event:
sub.last = now
default:
h.dropped++
}
}
}
/* ── Broker handlers ──────────────────────────────────────────────────────── */
func (h *LiveHub) onOrders(_ mqtt.Client, msg mqtt.Message) {
var batch models.PosOrderBatch
if err := json.Unmarshal(msg.Payload(), &batch); err != nil {
return // The ingest consumer logs the bad payload; one complaint is enough.
}
store, terminal := topicIdentity(msg.Topic())
if batch.Storeid == "" {
batch.Storeid = store
}
if batch.Terminalid == "" {
batch.Terminalid = terminal
}
// Deduplicated: a bill with four lines of the same SKU is one product.
seen := make(map[int]struct{})
productids := make([]int, 0, 8)
for _, order := range batch.Orders {
for _, item := range order.Items {
id, err := strconv.Atoi(strings.TrimSpace(item.Productid))
if err != nil || id == 0 {
continue
}
if _, done := seen[id]; done {
continue
}
seen[id] = struct{}{}
productids = append(productids, id)
}
}
h.Broadcast(LiveEvent{
Type: "sale",
Locationid: locationFromStore(batch.Storeid),
Productids: productids,
Terminalid: batch.Terminalid,
})
}
func (h *LiveHub) onCustomers(_ mqtt.Client, msg mqtt.Message) {
var batch models.PosCustomerBatch
if err := json.Unmarshal(msg.Payload(), &batch); err != nil {
return
}
store, terminal := topicIdentity(msg.Topic())
if batch.Storeid == "" {
batch.Storeid = store
}
h.Broadcast(LiveEvent{
Type: "customer",
Locationid: locationFromStore(batch.Storeid),
Terminalid: firstNonEmpty(batch.Terminalid, terminal),
})
}
// onHealth drives the presence board.
//
// Read from the topic rather than the body, the same rule bills follow: a till
// that could name its own store could appear in another tenant's console.
func (h *LiveHub) onHealth(_ mqtt.Client, msg mqtt.Message) {
store, terminal := topicIdentity(msg.Topic())
h.Broadcast(LiveEvent{
Type: "terminal",
Locationid: locationFromStore(store),
Terminalid: terminal,
})
}
/* ── Helpers ──────────────────────────────────────────────────────────────── */
// locationFromStore parses the topic's store segment.
//
// `models.PosOrderBatch` documents it: "Storeid carries the numeric
// tenantlocations.locationid as a string." Anything unparseable yields 0, which
// Broadcast drops — an event with no outlet has nowhere to go.
func locationFromStore(store string) int {
id, err := strconv.Atoi(strings.TrimSpace(store))
if err != nil {
return 0
}
return id
}
func firstNonEmpty(values ...string) string {
for _, value := range values {
if strings.TrimSpace(value) != "" {
return value
}
}
return ""
}
// hostSuffix keeps each replica's client id distinct without needing config.
func hostSuffix() string {
if host := strings.TrimSpace(os.Getenv("HOSTNAME")); host != "" {
return host
}
return strconv.FormatInt(time.Now().UnixNano(), 36)
}
// Stats reports what the hub is doing, for the health endpoint.
func (h *LiveHub) Stats() (outlets, connections int, dropped uint64) {
if h == nil {
return 0, 0, 0
}
h.mu.RLock()
defer h.mu.RUnlock()
for _, set := range h.subscribers {
connections += len(set)
}
return len(h.subscribers), connections, h.dropped
}