live stock update
This commit is contained in:
348
messaging/livehub.go
Normal file
348
messaging/livehub.go
Normal file
@@ -0,0 +1,348 @@
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user