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) // MQTT_USER is the name the ingest consumer and the env files use. This // read MQTT_USERNAME for a while, so the live stream connected to the // production broker with no credentials at all; MQTT_USERNAME is still // honoured for any deployment that set it. user := strings.TrimSpace(os.Getenv("MQTT_USER")) if user == "" { user = strings.TrimSpace(os.Getenv("MQTT_USERNAME")) } if 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 }