package api import ( "sync" "time" ) // LiveHub carries camera frames from a shop PC to whoever is watching. // // The shop PC's engine serves MJPEG on its own loopback, behind a router with // no inbound route, so head office cannot pull it. What head office CAN do is // answer the agent's outbound requests - which is the whole shape of this // product already - so the agent asks "is anyone watching?", and pushes frames // up for as long as somebody is. // // This is NOT true video. It is a few frames a second of re-encoded JPEG, which // is what an outbound HTTP relay can carry honestly. Real 25 fps video needs // WebRTC and a TURN server; this needs neither, and for "is that camera pointed // at the right place, and is someone in the shop" a few frames a second is what // the question actually requires. // // It is deliberately the OPPOSITE of the arrivals Hub, which is a doorbell that // pushes nothing because nothing may be lost. Here a dropped frame is the // correct outcome: a slow viewer must never stall the pump or accumulate a // backlog of stale pictures, because the only frame worth having is the newest // one. So each viewer gets a one-slot buffer and a full slot is overwritten. type LiveHub struct { mu sync.Mutex cameras map[string]*liveCamera } type liveCamera struct { viewers map[chan []byte]struct{} // wanted is refreshed by every watching viewer. The agent stops pushing // when it lapses, which is what keeps a shop's uplink idle when nobody is // looking - the entire cost argument for this feature. wanted time.Time // bell fires when the first viewer arrives, so an agent long-polling for // work is answered immediately instead of on its next tick. bell chan struct{} } // liveIdle is how long a camera stays "wanted" after the last viewer refreshed // it. Longer than the viewer's refresh interval so an ordinary pause between // refreshes does not stop the stream, short enough that a browser that // vanished stops a shop uploading within seconds. const liveIdle = 12 * time.Second func NewLiveHub() *LiveHub { return &LiveHub{cameras: map[string]*liveCamera{}} } // Watch registers a viewer and returns its frame channel plus a release func. func (h *LiveHub) Watch(cameraID string) (<-chan []byte, func()) { h.mu.Lock() defer h.mu.Unlock() c := h.cameras[cameraID] if c == nil { c = &liveCamera{viewers: map[chan []byte]struct{}{}, bell: make(chan struct{}, 1)} h.cameras[cameraID] = c } ch := make(chan []byte, 1) c.viewers[ch] = struct{}{} c.wanted = time.Now().Add(liveIdle) select { case c.bell <- struct{}{}: default: } return ch, func() { h.mu.Lock() defer h.mu.Unlock() if cam := h.cameras[cameraID]; cam != nil { delete(cam.viewers, ch) if len(cam.viewers) == 0 { // Dropped entirely rather than left empty: an estate's worth of // cameras nobody is watching would otherwise accumulate here // for the life of the process. delete(h.cameras, cameraID) } } close(ch) } } // Keep extends a camera's interest window. Called by each viewer as it reads, // so interest expires on its own when a browser goes away without saying so - // which is the normal way a tab closes. func (h *LiveHub) Keep(cameraID string) { h.mu.Lock() defer h.mu.Unlock() if c := h.cameras[cameraID]; c != nil { c.wanted = time.Now().Add(liveIdle) } } // Wanted reports whether anyone is watching this camera right now. func (h *LiveHub) Wanted(cameraID string) bool { h.mu.Lock() defer h.mu.Unlock() c := h.cameras[cameraID] return c != nil && len(c.viewers) > 0 && time.Now().Before(c.wanted) } // WantedAmong filters a site's cameras down to the ones being watched. The // agent asks with the cameras it has, so this never has to know a site's // inventory. func (h *LiveHub) WantedAmong(ids []string) []string { var out []string for _, id := range ids { if h.Wanted(id) { out = append(out, id) } } return out } // Bell returns a channel that fires when a viewer starts watching one of these // cameras, so an agent waiting for work wakes at once rather than on a tick. // A nil channel blocks forever, which is the right behaviour for the caller's // select when none of the cameras is known here yet. func (h *LiveHub) Bell(ids []string) <-chan struct{} { h.mu.Lock() defer h.mu.Unlock() for _, id := range ids { if c := h.cameras[id]; c != nil { return c.bell } } // Register a placeholder for the first camera so a later Watch can ring // something. Cheap: one struct per camera an agent asked about. if len(ids) == 0 { return nil } c := &liveCamera{viewers: map[chan []byte]struct{}{}, bell: make(chan struct{}, 1)} h.cameras[ids[0]] = c return c.bell } // Publish hands one frame to every viewer of a camera and reports whether any // remain. The agent uses that answer to stop pushing. // // A viewer whose slot is full has its pending frame REPLACED, never queued. The // newest frame is the only one worth having, and a queue here would show a // viewer an ever-growing delay behind the shop rather than dropping back to // live. func (h *LiveHub) Publish(cameraID string, frame []byte) bool { h.mu.Lock() defer h.mu.Unlock() c := h.cameras[cameraID] if c == nil || len(c.viewers) == 0 { return false } for ch := range c.viewers { select { case ch <- frame: default: select { case <-ch: default: } select { case ch <- frame: default: } } } return time.Now().Before(c.wanted) } // live returns the hub, creating it on first use. // // Lazily, and stored on the Server, so a server built without one still works: // unlike the arrivals doorbell there is no degraded mode to fall back to here, // and a nil map panic on an endpoint somebody forgot to wire is the worst way // to find out. func (s *Server) live() *LiveHub { s.liveOnce.Do(func() { if s.Live == nil { s.Live = NewLiveHub() } }) return s.Live }