166 lines
5.0 KiB
Go
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,
|
|
},
|
|
})
|
|
}
|