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