`main.go` only ever loaded `.env`; the `APP_ENV` switch described in `.env.local` / `.env.production` did not exist, and a missing variable surfaced one restart at a time as a log.Fatalf inside db.Connect. config.Load now picks `.env.<APP_ENV>` (default local) then `.env`, with real environment winning, reads every setting into one typed Config and reports everything missing in one message. Production insists on a POS signing secret; local warns when DB_HOST is not a local address. db, redis and the image store take the Config instead of reading env themselves. Also: - livehub read MQTT_USERNAME while everything else uses MQTT_USER, so the console stream connected to the broker unauthenticated. Both accepted. - .dockerignore: `COPY . .` was baking .env.production into the image. Dockerfile sets APP_ENV=production. - Drop utils/config.go (dead viper loader) and create_table.go (unused, hardcoded production DSN); go mod tidy removes viper. - .env.example lists every variable the code reads; docs/ENVIRONMENT.md. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
357 lines
12 KiB
Go
357 lines
12 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)
|
|
|
|
// 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
|
|
}
|