package api import ( "encoding/base64" "encoding/binary" "errors" "fmt" "io" "net/http" "time" ) const ( // liveMaxFrame bounds one frame. The agent re-encodes to ~640 px before // sending, so a frame is tens of KB; 1 MB is a ceiling, not a target. liveMaxFrame = 1 << 20 // liveSession caps one push. A browser tab left open for a week must not // leave a shop uploading for a week; the agent simply asks again while // anyone is still watching, so the cap costs a reconnect, not the stream. liveSession = 5 * time.Minute // liveWaitForWork is how long the agent's poll is held open. Long enough // that an idle site makes ~2 requests a minute; short enough to sit well // inside any proxy's idle timeout. liveWaitForWork = 25 * time.Second ) // handleWatchLive streams one camera's frames to a signed-in user over SSE. // // SSE rather than serving MJPEG directly, for the same reason the snapshot is // not a signed link: an cannot send an Authorization header, and minting // a URL that works without a session - for LIVE video of a shop floor, no less // - would be a much worse trade than the 33% base64 costs. func (s *Server) handleWatchLive(w http.ResponseWriter, r *http.Request) { p := PrincipalFrom(r.Context()) id, ok := s.resolveCamera(w, r, r.PathValue("id")) if !ok { return } // Ownership is checked HERE, once, before anything is streamed. Everything // after this point is keyed on a camera id, and a hub does not know whose // camera it is holding. if _, _, err := s.Store.CameraRef(r.Context(), p.ClientID, id); err != nil { writeErr(w, http.StatusNotFound, "not_found", "No such camera.") return } flusher, ok := w.(http.Flusher) if !ok { s.serverError(w, "live", errors.New("this server cannot stream")) return } frames, release := s.live().Watch(id) defer release() h := w.Header() h.Set("Content-Type", "text/event-stream") h.Set("Cache-Control", "no-store") h.Set("Connection", "keep-alive") // Without this a proxy buffers the stream into one response that arrives // when the connection closes - which for live video means never. h.Set("X-Accel-Buffering", "no") w.WriteHeader(http.StatusOK) // Told up front, so a viewer can say "waiting for the shop PC" rather than // showing an empty box while the agent is still being asked. fmt.Fprint(w, "event: waiting\ndata: {}\n\n") flusher.Flush() ctx := r.Context() // Refreshed as we go rather than once at the start: this is what tells the // agent somebody is still there, and a viewer that has gone away stops a // shop uploading within seconds without having to announce anything. keep := time.NewTicker(liveIdle / 3) defer keep.Stop() for { select { case <-ctx.Done(): return case <-keep.C: s.live().Keep(id) case frame, ok := <-frames: if !ok { return } s.live().Keep(id) if _, err := fmt.Fprintf(w, "event: frame\ndata: %s\n\n", base64.StdEncoding.EncodeToString(frame)); err != nil { return } flusher.Flush() } } } // handleAgentLiveWanted is the shop PC asking whether anyone is watching. // // Held open rather than answered immediately: an agent polling every few // seconds would put a floor under how quickly a live view can start, and one // polling slowly would put a ceiling on it. Holding the request means pressing // "Live" reaches the shop PC at once, and an idle site costs about two requests // a minute. func (s *Server) handleAgentLiveWanted(w http.ResponseWriter, r *http.Request, ap AgentPrincipal) { ids, err := s.Store.SiteCameraIDs(r.Context(), ap.SiteID) if err != nil { s.serverError(w, "live wanted", err) return } if wanted := s.live().WantedAmong(ids); len(wanted) > 0 { writeJSON(w, http.StatusOK, map[string]any{"cameras": wanted}) return } select { case <-r.Context().Done(): return case <-s.live().Bell(ids): case <-time.After(liveWaitForWork): } writeJSON(w, http.StatusOK, map[string]any{ "cameras": s.live().WantedAmong(ids)}) } // handleAgentPushLive receives frames for as long as somebody is watching. // // One request carrying many frames, each prefixed with its length, rather than // a request per frame: at a few frames a second the per-request overhead and // the TLS handshakes would cost more than the pictures. func (s *Server) handleAgentPushLive(w http.ResponseWriter, r *http.Request, ap AgentPrincipal) { id := r.PathValue("camera") if !looksLikeUUID(id) { writeErr(w, http.StatusNotFound, "not_found", "No such camera.") return } // The camera must belong to the AGENT's own site. Without this an agent // could push its own pictures into another site's live view - a camera id // is not a secret, and the agent supplies this one. siteID, _, err := s.Store.CameraRefBySite(r.Context(), ap.SiteID, id) if err != nil || siteID != ap.SiteID { writeErr(w, http.StatusNotFound, "not_found", "No such camera.") return } deadline := time.Now().Add(liveSession) body := r.Body var header [4]byte frames := 0 for { if time.Now().After(deadline) { break } if _, err := io.ReadFull(body, header[:]); err != nil { break } n := binary.BigEndian.Uint32(header[:]) if n == 0 || n > liveMaxFrame { // A length this side cannot trust ends the stream rather than // allocating what it was told to. writeErr(w, http.StatusBadRequest, "bad_frame", "Frame size out of range.") return } frame := make([]byte, n) if _, err := io.ReadFull(body, frame); err != nil { break } frames++ if !s.live().Publish(id, frame) { // Nobody is watching any more. Saying so in the response is what // stops the shop uploading; the agent goes back to waiting. break } } writeJSON(w, http.StatusOK, map[string]any{"frames": frames}) }