From dc9c049bb9a7d214d98d722b4ead356215143dad Mon Sep 17 00:00:00 2001 From: abhishek Date: Thu, 27 Aug 2026 11:17:17 +0530 Subject: [PATCH] live stock update --- controllers/liveController.go | 165 ++++++++++++++ facade/container.go | 7 + main.go | 6 + messaging/livehub.go | 348 ++++++++++++++++++++++++++++++ repositories/posUserRepository.go | 8 +- routes/posroutes.go | 16 ++ 6 files changed, 548 insertions(+), 2 deletions(-) create mode 100644 controllers/liveController.go create mode 100644 messaging/livehub.go diff --git a/controllers/liveController.go b/controllers/liveController.go new file mode 100644 index 0000000..586c96c --- /dev/null +++ b/controllers/liveController.go @@ -0,0 +1,165 @@ +package controllers + +import ( + "bufio" + "encoding/json" + "fmt" + "net/http" + "strconv" + "strings" + "time" + + "nearle/messaging" + "nearle/services" + + "github.com/gofiber/fiber/v2" + "github.com/valyala/fasthttp" +) + +// The console's live event stream. +// +// Server-Sent Events rather than a WebSocket, deliberately. The traffic is one +// way — the console never sends anything up this pipe — and SSE is plain HTTP, +// so it inherits the reverse proxy, the TLS termination and the load balancer +// already in front of this service with no upgrade handshake to configure. The +// browser's own `EventSource` also reconnects on its own, which is a reconnect +// loop nobody has to write or get wrong. +// +// What goes down it is a nudge, never a figure: "outlet 1185 sold something, +// products 42 and 77". The console re-reads through the normal endpoints. See +// `messaging/livehub.go` for why. + +type LiveController struct { + posService services.PosService +} + +func NewLiveController(posService services.PosService) *LiveController { + return &LiveController{posService: posService} +} + +const ( + // Below every idle timeout worth worrying about: nginx and most cloud load + // balancers cut an idle connection at 60s, and a quiet shop produces no + // events for hours. + liveHeartbeat = 25 * time.Second + + // A connection is recycled rather than held forever, so a replica being + // drained empties out on its own and a client that has silently gone away + // stops being written to. The browser reconnects immediately; the console + // refetches on reconnect anyway, so the seam is invisible. + liveMaxAge = 30 * time.Minute +) + +// Stream opens the event stream for one outlet. +// +// GET /live/api/v1/web/live/events?tenantid=1147&locationid=1185 +// +// Scoped exactly like every other console POS read: the tenant must own the +// outlet. Without that check any caller could name any locationid and watch a +// competitor's tills report in. +func (ctl *LiveController) Stream(c *fiber.Ctx) error { + tenantID, _ := strconv.Atoi(strings.TrimSpace(c.Query("tenantid"))) + locationID, _ := strconv.Atoi(strings.TrimSpace(c.Query("locationid"))) + + if tenantID <= 0 || locationID <= 0 { + return c.Status(http.StatusBadRequest).JSON(fiber.Map{ + "status": false, "code": http.StatusBadRequest, + "message": "tenantid and locationid are required", + }) + } + + allowed, err := ctl.posService.LocationAllowed(tenantID, locationID) + if err != nil { + return c.Status(http.StatusInternalServerError).JSON(fiber.Map{ + "status": false, "code": http.StatusInternalServerError, + "message": "could not verify the outlet", + }) + } + if !allowed { + return c.Status(http.StatusForbidden).JSON(fiber.Map{ + "status": false, "code": http.StatusForbidden, + "message": fmt.Sprintf("outlet %d does not belong to tenant %d", locationID, tenantID), + }) + } + + events, release := messaging.Hub.Subscribe(locationID) + + c.Set("Content-Type", "text/event-stream") + c.Set("Cache-Control", "no-cache") + c.Set("Connection", "keep-alive") + // Nginx buffers proxied responses by default, which holds each event until + // the buffer fills — for a stream of 80-byte messages, indefinitely. + c.Set("X-Accel-Buffering", "no") + + c.Context().SetBodyStreamWriter(fasthttp.StreamWriter(func(w *bufio.Writer) { + // The writer owns the subscription from here: this closure outlives the + // handler, so releasing in a defer above would unsubscribe instantly. + defer release() + + // Tells EventSource how long to wait before reconnecting, and gives the + // client something to prove the stream is open. + if _, err := fmt.Fprintf(w, "retry: 3000\nevent: open\ndata: {\"locationid\":%d}\n\n", locationID); err != nil { + return + } + if err := w.Flush(); err != nil { + return + } + + heartbeat := time.NewTicker(liveHeartbeat) + defer heartbeat.Stop() + deadline := time.NewTimer(liveMaxAge) + defer deadline.Stop() + + for { + select { + case event, open := <-events: + if !open { + return + } + payload, err := json.Marshal(event) + if err != nil { + continue + } + if _, err := fmt.Fprintf(w, "event: %s\ndata: %s\n\n", event.Type, payload); err != nil { + return + } + // A failed flush is how a disconnected client is discovered: + // fasthttp surfaces the broken pipe here and nowhere else. + if err := w.Flush(); err != nil { + return + } + + case <-heartbeat.C: + // An SSE comment. EventSource ignores it; every proxy in the + // path sees traffic and keeps the connection open. + if _, err := fmt.Fprint(w, ": ping\n\n"); err != nil { + return + } + if err := w.Flush(); err != nil { + return + } + + case <-deadline.C: + fmt.Fprint(w, "event: bye\ndata: {}\n\n") + _ = w.Flush() + return + } + } + })) + + return nil +} + +// Health reports what the stream is doing. Useful for answering "is anyone +// actually connected" without reading logs. +func (ctl *LiveController) Health(c *fiber.Ctx) error { + outlets, connections, dropped := messaging.Hub.Stats() + return c.JSON(fiber.Map{ + "status": true, "code": http.StatusOK, + "details": fiber.Map{ + "outlets_watched": outlets, + "connections": connections, + "dropped": dropped, + }, + }) +} diff --git a/facade/container.go b/facade/container.go index daf22bb..0ceba4f 100644 --- a/facade/container.go +++ b/facade/container.go @@ -20,6 +20,7 @@ type Facade struct { StockRequestController *controllers.StockRequestController CatalogueController *controllers.CatalogueController PosController *controllers.PosController + LiveController *controllers.LiveController // Held so the NATS consumer can reach the ingest without going through // HTTP. Unexported: everything else should use the controller. @@ -94,6 +95,11 @@ func NewFacade(db *gorm.DB, catalogueDB *gorm.DB) *Facade { posService := services.NewPosService(posRepo, posPresence) posController := controllers.NewPosController(posService) + // Shares the POS service purely for its outlet-ownership check — the + // stream itself reads no database and holds no state beyond its + // subscribers. + liveController := controllers.NewLiveController(posService) + return &Facade{ UserController: userController, ProductController: productController, @@ -106,6 +112,7 @@ func NewFacade(db *gorm.DB, catalogueDB *gorm.DB) *Facade { StockRequestController: stockRequestController, CatalogueController: catalogueController, PosController: posController, + LiveController: liveController, posService: posService, } } diff --git a/main.go b/main.go index c04fedb..5e98d47 100644 --- a/main.go +++ b/main.go @@ -118,6 +118,12 @@ func main() { // would be a cycle — so they hold an interface and this supplies it. Left // unset the notification is simply skipped, which is what happens when the // broker is unavailable and must not stop the API booting. + // The console's live stream listens to the same broker, read-only, and on + // EVERY replica — unlike the ingest consumer below, which is elected. A + // console connected to a non-elected replica must still see its tills. + // Never fatal: no broker simply means the console keeps polling. + messaging.StartLiveHub() + posMqtt, err := messaging.StartPosMqttConsumer(f.PosService()) if err != nil { log.Fatal("POS MQTT consumer failed to start:", err) diff --git a/messaging/livehub.go b/messaging/livehub.go new file mode 100644 index 0000000..695e04e --- /dev/null +++ b/messaging/livehub.go @@ -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 +} diff --git a/repositories/posUserRepository.go b/repositories/posUserRepository.go index f734777..841d15e 100644 --- a/repositories/posUserRepository.go +++ b/repositories/posUserRepository.go @@ -487,9 +487,13 @@ func (r *posRepository) ListPosUsers(tenantID, locationID int, includeInactive b params := []interface{}{tenantID, locationID, models.PosRoleSupervisor, models.PosRoleCashier} if !includeInactive { - query += ` AND LOWER(COALESCE(status,'active')) <> 'inactive'` + // `a.status`, qualified. `staffshifts` carries a `status` column too, so + // the bare name made Postgres refuse the whole statement with + // `column reference "status" is ambiguous` (42702) — a 500 on every + // console call, since the console never asks for inactive rows. + query += ` AND LOWER(COALESCE(a.status,'active')) <> 'inactive'` } - query += ` ORDER BY userid` + query += ` ORDER BY a.userid` if err := r.db.Raw(query, params...).Scan(&rows).Error; err != nil { return nil, err diff --git a/routes/posroutes.go b/routes/posroutes.go index f1b3445..69c3097 100644 --- a/routes/posroutes.go +++ b/routes/posroutes.go @@ -76,6 +76,7 @@ func RegisterPosRoutes(api fiber.Router, f *facade.Facade) { registerPosStaffConsoleRoutes(api, f) registerPosReadConsoleRoutes(api, f) + registerLiveRoutes(api, f) } // The same counter-sales reads, for callers that are not a terminal. @@ -157,3 +158,18 @@ func registerPosStaffConsoleRoutes(api fiber.Router, f *facade.Facade) { g.Put("/updatestaffshift", f.PosController.WebUpdateStaffShift) } } + +// The console's live event stream. +// +// Separate from the POS groups above because the caller is the back office, +// not a terminal, and because it is the one route here that holds a +// connection open — worth being obvious about when reading the route table. +// +// `/events` is Server-Sent Events and never returns until the client goes +// away or the age limit fires. `/events/health` is an ordinary JSON read that +// says how many consoles are attached to THIS replica. +func registerLiveRoutes(api fiber.Router, f *facade.Facade) { + g := api.Group("/v1/web/live") + g.Get("/events", f.LiveController.Stream) + g.Get("/events/health", f.LiveController.Health) +}