Files
Behavision/agent/pkg/cameras/live.go
Suriyakumarvijayanayagam 50d122e5f0 An unused constant was a Stop() that could hang forever
Ran staticcheck across all three Go modules for the first time. server
(23k lines) and desktop came back clean. agent had seven findings, and
one of them was not tidiness.

`stopGrace = 10 * time.Second` was declared and wired to nothing.
Stop() cancels the context, cmd.Cancel kills the process tree, and then
Stop() blocks on cmd.Wait() - which, with no WaitDelay set, waits not
just for the process but for every writer of its stdout pipe to close.
One grandchild still holding that pipe hangs Wait, hangs Stop, and on the
desktop app that is the tray's Quit never returning. The constant named
the intent and nothing read it. cmd.WaitDelay = stopGrace is the line
that was missing.

The rest were real but small: an unused field in the live relay, an
unused sleep helper in the pump, and "net/url" imported twice under two
names - both genuinely used, in two functions doing the same job for the
same reason, so they are unified rather than one deleted. My first pass
deleted the wrong one on a bad grep and the build caught it immediately.

Three findings are suppressed rather than fixed, with the reason stated:

- Two "error strings should not end with punctuation". Both are
  multi-line messages a shop operator reads at a counter, not errors
  anything wraps. ST1005 exists because wrapped errors concatenate
  mid-sentence; stripping the full stops would run three sentences
  together to satisfy a rule that does not apply.
- A deliberately nil context in a pump test - the point of the test is
  that an unconnected client does not panic. It already carried
  //nolint:staticcheck, which is golangci-lint's directive and
  staticcheck ignores, which is why it kept being reported.

Also tidied agent/go.mod, which had paho and x/sys marked indirect while
being imported directly.

All three modules clean, all suites pass: 21 Go packages, 226 engine
tests.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-30 15:14:09 +05:30

243 lines
7.3 KiB
Go

package cameras
import (
"bytes"
"context"
"crypto/sha256"
"encoding/binary"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"time"
)
// Live relays camera frames to head office, but only while somebody is
// watching.
//
// The engine serves MJPEG on this PC's loopback and this PC sits behind a
// router with no inbound route, so head office cannot pull it. It can answer
// our outbound requests, which is the shape of everything else here: we ask
// "is anyone watching?", and push frames for as long as the answer is yes.
//
// It is a few frames a second of re-encoded JPEG, not 25 fps video. True video
// needs WebRTC and a TURN server; this needs neither, and answers the question
// somebody at head office is actually asking - what does that camera see right
// now - at a cost a shop's uplink can carry.
//
// **Nothing is uploaded when nobody is looking.** That is the entire cost
// argument, and it is why the wanted-check comes first and the push stops the
// moment the server says the last viewer has gone.
type Live struct {
Engine *EngineClient
Cloud *CloudClient
Log *log.Logger
FPS float64
Width int
Quality int
}
// Defaults, measured against the office camera rather than guessed.
//
// The engine produces ~12 distinct frames a second, so asking for more than
// that only re-sends pictures the viewer already has - which is why the poll
// runs slightly ahead of it and identical frames are dropped rather than sent.
// 640 px at quality 60 is ~20 KB, so a watcher costs ~200 KB/s at the full
// rate, and a camera nobody is watching costs nothing at all.
const (
DefaultLiveFPS = 15.0
DefaultLiveWidth = 640
DefaultLiveQuality = 60
)
func NewLive(eng *EngineClient, cloud *CloudClient, logger *log.Logger) *Live {
return &Live{Engine: eng, Cloud: cloud, Log: logger,
FPS: DefaultLiveFPS, Width: DefaultLiveWidth, Quality: DefaultLiveQuality}
}
// Run waits for viewers and serves them until the context ends.
func (l *Live) Run(ctx context.Context) {
if l.Engine == nil || l.Cloud == nil {
return
}
for {
if ctx.Err() != nil {
return
}
wanted, err := l.Cloud.LiveWanted(ctx)
if err != nil {
// Unclaimed, offline, or head office is down. All three mean the
// same thing here - nobody can be watching - so back off rather
// than hammering, and keep the shop's own recognition untouched.
if ctx.Err() != nil {
return
}
l.sleep(ctx, 15*time.Second)
continue
}
if len(wanted) == 0 {
// The poll is held open by the server, so an empty answer already
// means ~25 s passed. No extra delay.
continue
}
for _, id := range wanted {
if ctx.Err() != nil {
return
}
l.serve(ctx, id)
}
}
}
// serve pushes frames for one camera until the server says stop.
func (l *Live) serve(ctx context.Context, cameraID string) {
engineID, err := l.Cloud.LiveEngineID(ctx, cameraID)
if err != nil {
l.logf("live %s: %v", cameraID, err)
l.sleep(ctx, 2*time.Second)
return
}
interval := time.Duration(float64(time.Second) / l.fps())
// A pipe so frames can be written as they are grabbed while one request
// carries all of them. A request per frame would spend more on handshakes
// and headers than on pictures.
pr, pw := io.Pipe()
done := make(chan error, 1)
go func() { done <- l.Cloud.PushLive(ctx, cameraID, pr) }()
tick := time.NewTicker(interval)
defer tick.Stop()
// The engine re-serves its latest frame until the pipeline produces a new
// one, so polling faster than it encodes returns the SAME picture again.
// Measured: 93 polls in 6 s yielded 72 distinct frames. Sending the
// duplicates would cost a fifth of the bandwidth for nothing, so the poll
// runs a little ahead of the engine and the repeats are dropped - which is
// what lets the rate follow the camera instead of a guess.
var lastSum [32]byte
for {
select {
case <-ctx.Done():
_ = pw.CloseWithError(context.Canceled)
<-done
return
case err := <-done:
// The server closed the request: the last viewer went away, or the
// session cap was reached. Either way stop grabbing frames.
_ = pw.Close()
if err != nil {
l.logf("live %s ended: %v", cameraID, err)
}
return
case <-tick.C:
}
jpeg, err := l.Engine.Frame(ctx, engineID, l.Width, l.Quality)
if err != nil || len(jpeg) == 0 {
// A camera that is reconnecting has no frame. Keep the request
// open - the viewer sees the last frame rather than a dropped
// stream, and the next tick may well have one.
continue
}
if sum := sha256.Sum256(jpeg); sum == lastSum {
continue
} else {
lastSum = sum
}
var hdr [4]byte
binary.BigEndian.PutUint32(hdr[:], uint32(len(jpeg)))
if _, err := pw.Write(hdr[:]); err != nil {
<-done
return
}
if _, err := pw.Write(jpeg); err != nil {
<-done
return
}
}
}
func (l *Live) fps() float64 {
if l.FPS <= 0 || l.FPS > 25 {
// A ceiling rather than a target: duplicate frames are dropped, so
// polling above what the engine encodes costs requests and no
// bandwidth - but it is still work, on the PC doing the recognition.
return DefaultLiveFPS
}
return l.FPS
}
func (l *Live) sleep(ctx context.Context, d time.Duration) {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
case <-t.C:
}
}
func (l *Live) logf(format string, args ...any) {
if l.Log != nil {
l.Log.Printf(format, args...)
}
}
// ------------------------------------------------------------------ wire --
// LiveWanted asks head office which of this site's cameras are being watched.
// The server holds the request open, so this returns promptly when somebody
// presses Live and after ~25 s when nobody has.
func (c *CloudClient) LiveWanted(ctx context.Context) ([]string, error) {
var body struct {
Cameras []string `json:"cameras"`
}
// Longer than the server's own wait, so a held request is not cut off by
// our own client timeout and reported as a failure.
ctx, cancel := context.WithTimeout(ctx, 60*time.Second)
defer cancel()
if err := c.do(ctx, http.MethodGet, "/api/agent/live", nil, &body); err != nil {
return nil, err
}
return body.Cameras, nil
}
// LiveEngineID maps head office's camera uuid to the name the engine knows,
// which is the only name this PC can ask for a frame with.
func (c *CloudClient) LiveEngineID(ctx context.Context, cameraID string) (string, error) {
desired, err := c.Desired(ctx)
if err != nil {
return "", err
}
for _, d := range desired {
if d.ID == cameraID {
return d.CameraID, nil
}
}
return "", fmt.Errorf("camera %s is not one of this site's", cameraID)
}
// PushLive streams frames until the server stops reading.
func (c *CloudClient) PushLive(ctx context.Context, cameraID string, body io.Reader) error {
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
c.Base+"/api/agent/cameras/"+cameraID+"/live", body)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+c.Token)
req.Header.Set("Content-Type", "application/octet-stream")
resp, err := c.Client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
blob, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10))
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("head office: %s: %s", resp.Status, bytes.TrimSpace(blob))
}
var out struct {
Frames int `json:"frames"`
}
_ = json.Unmarshal(blob, &out)
return nil
}