package api import ( "encoding/base64" "fmt" "net/http" "strconv" "strings" "time" ) const ( // A poller asking for more than this is either paging history through the // wrong endpoint or has lost its cursor. Capped rather than rejected: a // mobile app coming back from a tunnel should get a big catch-up page and // a fresh cursor, not an error it has no way to recover from. maxArrivals = 200 defaultArrivals = 50 // How often a live stream re-queries even if no doorbell rings. This is the // safety net for a second server instance whose ingest this process cannot // hear, so it is slow on purpose - it is the fallback, not the mechanism. streamFallback = 15 * time.Second // SSE comment sent on an idle connection. Mobile networks and reverse // proxies both close a stream that has been silent for a minute or two, // and a client that reconnects every ninety seconds is a client that // re-queries constantly. streamKeepalive = 20 * time.Second ) // handleArrivals is the live feed: who walked in, with their photos, in one // request. // // This is the endpoint a mobile app or a shop-floor screen actually needs, and // it is the one thing the API could not previously answer. `GET /api/visitors` // searches a customer list by name; it cannot tell you that four people just // came through the door, and until it could, a client had no way to know which // customer ids to ask about. func (s *Server) handleArrivals(w http.ResponseWriter, r *http.Request) { p := PrincipalFrom(r.Context()) q := ArrivalQuery{ ClientID: p.ClientID, SiteID: siteParam(r), Limit: queryInt(r, "limit", defaultArrivals, maxArrivals), } if q.SiteID != "" { var ok bool if q.SiteID, ok = s.resolveSiteFilter(w, r, q.SiteID); !ok { return } } // The site is still filtered by client_id in SQL as well. A site_id from // the query string is caller-controlled, and this is a read of other // people's customers if it is ever trusted on its own. if cur := trim(r.URL.Query().Get("cursor")); cur != "" { seq, err := decodeCursor(cur) if err != nil { badRequest(w, "that cursor is not one of ours - drop it and poll again without one") return } q.AfterSeq = &seq } else if since := trim(r.URL.Query().Get("since")); since != "" { at, err := time.Parse(time.RFC3339, since) if err != nil { badRequest(w, "since must look like 2026-09-02T10:30:00Z") return } at = at.UTC() q.Since = &at } page, err := s.arrivalPage(r, q) if err != nil { s.serverError(w, "arrivals", err) return } writeJSON(w, http.StatusOK, page) } // arrivalPage runs the query, attaches photos and returns the next cursor. // Shared by the poll and the stream so the two cannot answer differently. func (s *Server) arrivalPage(r *http.Request, q ArrivalQuery) (ArrivalPage, error) { rows, err := s.Store.Arrivals(r.Context(), q) if err != nil { return ArrivalPage{}, err } s.attachImages(r, rows) page := ArrivalPage{ Arrivals: rows, PolledAt: s.now().UTC().Format(time.RFC3339Nano), } if page.Arrivals == nil { // An empty list, never null. A client looping over the response should // not have to special-case a quiet minute. page.Arrivals = []Arrival{} } if n := len(rows); n > 0 { // The LAST row, and the rows are ascending by position, so this is the // highest position the caller has now seen. page.Cursor = encodeCursor(rows[n-1].Seq) } else if q.AfterSeq != nil { // Nothing new. Hand the caller its own position back rather than an // empty string, so a poll that returns nothing does not reset the feed // to the beginning on the next request. page.Cursor = encodeCursor(*q.AfterSeq) } return page, nil } // attachImages swaps each row's object key for a short-lived signed link. // // One audit row for the whole page, not one per photo. Every read of a face is // worth recording - "who looked at my customers" has to be answerable - but a // tablet polling this feed every two seconds would write tens of thousands of // rows a day and bury the single deliberate look that an investigation is // actually after. The row records how many faces were surfaced and to whom, // which is the fact worth keeping. func (s *Server) attachImages(r *http.Request, rows []Arrival) { p := PrincipalFrom(r.Context()) seen := make([]string, 0, len(rows)) for i := range rows { key := rows[i].ImageKey rows[i].ImageKey = "" rows[i].Image = s.imageFor(key) if rows[i].Image.Available && rows[i].VisitorID != "" { seen = append(seen, rows[i].VisitorID) } } if len(seen) == 0 { return } s.Store.Audit(r.Context(), AuditEntry{ ClientID: p.ClientID, ActorID: p.UserID, ActorKind: "user", Action: "image.view.feed", Entity: "visits", Detail: map[string]any{"count": len(seen), "visitor_ids": seen}, }) } // handleArrivalStream is the same feed pushed instead of polled. // // Server-sent events rather than websockets: this direction is one-way, SSE is // stdlib with no dependency, it survives the reverse proxy in front of this // server unchanged, and browsers and mobile HTTP clients reconnect it on their // own. A websocket would buy bidirectionality that nothing here wants. func (s *Server) handleArrivalStream(w http.ResponseWriter, r *http.Request) { flusher, ok := w.(http.Flusher) if !ok { // Without flushing this is not a stream, it is a response that arrives // at the end of the day. Say so rather than appearing to work. writeErr(w, http.StatusInternalServerError, "server_error", "streaming is not available on this connection") return } p := PrincipalFrom(r.Context()) q := ArrivalQuery{ ClientID: p.ClientID, SiteID: siteParam(r), Limit: queryInt(r, "limit", defaultArrivals, maxArrivals), } if q.SiteID != "" { var ok bool if q.SiteID, ok = s.resolveSiteFilter(w, r, q.SiteID); !ok { return } } // Last-Event-ID is what the browser's EventSource resends automatically on // a dropped connection, so honouring it is what makes a reconnect lossless // without the client writing any recovery code. ?cursor= is the same thing // for a native client that cannot set the header. cur := trim(r.Header.Get("Last-Event-ID")) if cur == "" { cur = trim(r.URL.Query().Get("cursor")) } if cur != "" { if seq, err := decodeCursor(cur); err == nil { q.AfterSeq = &seq } // A cursor we cannot read is not worth failing a reconnect over: the // client falls back to the recent window, which is the same thing it // would get on a fresh connection. } h := w.Header() h.Set("Content-Type", "text/event-stream") h.Set("Cache-Control", "no-cache") h.Set("Connection", "keep-alive") // Traefik and nginx both buffer by default, which turns an event stream // into a file that arrives when the connection closes. h.Set("X-Accel-Buffering", "no") w.WriteHeader(http.StatusOK) flusher.Flush() bell, release := s.hub().Subscribe() defer release() ctx := r.Context() fallback := time.NewTicker(streamFallback) defer fallback.Stop() keepalive := time.NewTicker(streamKeepalive) defer keepalive.Stop() // Send whatever is already there before waiting for a doorbell, so a client // that connects after people have walked in is not blind until the next // one does. send := func() bool { page, err := s.arrivalPage(r, q) if err != nil { s.logf("ERROR arrival stream: %v", err) // Keep the connection: a transient database error should not log // a shop's screen out and start a reconnect storm across an estate. return true } if len(page.Arrivals) == 0 { return true } if page.Cursor != "" { // The SSE id becomes the client's Last-Event-ID, so the cursor // rides the protocol's own reconnect machinery instead of needing // application-level recovery. fmt.Fprintf(w, "id: %s\n", page.Cursor) // Advance our own position from the rows we just sent, so the next // wake-up asks for what comes after them. last := page.Arrivals[len(page.Arrivals)-1].Seq q.AfterSeq = &last } fmt.Fprint(w, "event: arrivals\ndata: ") if err := writeCompactJSON(w, page); err != nil { return false } fmt.Fprint(w, "\n\n") flusher.Flush() return true } if !send() { return } for { select { case <-ctx.Done(): return case id, open := <-bell: if !open { return } // One process serves every tenant, so a doorbell for somebody // else's shop must not cost this connection a query. if id != p.ClientID { continue } if !send() { return } case <-fallback.C: if !send() { return } case <-keepalive.C: // A comment line. Keeps proxies and mobile radios from deciding // the connection is dead, and is ignored by every SSE client. fmt.Fprint(w, ": keepalive\n\n") flusher.Flush() } } } func (s *Server) hub() *Hub { if s.Hub == nil { // A server built without one still streams; it just has no doorbell, // so it falls back to the slow tick. A nil map panic on an endpoint // somebody forgot to wire is a worse outcome than a slower feed. return nil } return s.Hub } // ---------------------------------------------------------------- cursors // Cursors are opaque on purpose, and base64 is what makes them look it. A // caller that reads one starts depending on the ordering column, and changing // that later then breaks every deployed mobile app rather than just this file - // which is exactly what happened once already, when the feed was ordered by a // timestamp and a random uuid. func encodeCursor(seq int64) string { return base64.RawURLEncoding.EncodeToString( []byte("v1:" + strconv.FormatInt(seq, 10))) } func decodeCursor(s string) (int64, error) { raw, err := base64.RawURLEncoding.DecodeString(s) if err != nil { return 0, fmt.Errorf("cursor is not base64: %w", err) } // The version prefix is what lets the ordering change again without // silently misreading cursors already held by deployed clients: an old // cursor fails to parse and the client restarts cleanly from the recent // window, rather than resuming at a position that now means something else. body, ok := strings.CutPrefix(string(raw), "v1:") if !ok { return 0, fmt.Errorf("cursor is not a v1 cursor") } seq, err := strconv.ParseInt(body, 10, 64) if err != nil || seq < 0 { return 0, fmt.Errorf("cursor position is not a number") } return seq, nil }