Files
backend_fiesta/controllers/liveController.go
2026-08-27 11:17:17 +05:30

166 lines
5.0 KiB
Go

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,
},
})
}