I got this wrong first time. "Head office cannot show live video cheaply" conflated TRUE VIDEO with SEEING THE CAMERA NOW, and only the first needs WebRTC and a TURN server. The shop PC is behind a router with no inbound route, so head office cannot pull the engine's MJPEG. It can answer the agent's outbound requests, which is the shape of everything else here: the server holds a poll open, the agent asks "is anyone watching?", and pushes JPEGs up for exactly as long as somebody is. Measured on the office camera: 98 KB full frame, 20.8 KB re-encoded at 640/q60, so one watcher costs ~83 KB/s. 47 frames arrived in 12 seconds - 4 fps, as configured. The UI says "about 4 frames a second" rather than letting anyone conclude the camera stutters. Nothing is uploaded when nobody is looking, which is the whole cost argument: Publish returns false once the last viewer goes, interest lapses on a timer each viewer refreshes as it reads (so a closed tab stops the upload within seconds), one push is capped at five minutes, and the UI streams one camera at a time. LiveHub is deliberately the opposite of the arrivals Hub. There a doorbell pushes nothing because nothing may be lost; here a dropped frame is the correct outcome, so each viewer has a one-slot buffer that is overwritten - the only frame worth having is the newest, and a queue would show an ever-growing delay behind the shop instead of dropping back to live. Ownership is proved once, before anything streams: the relay is keyed on a camera id, a hub does not know whose camera it holds, and a camera id is not a secret. Verified: another tenant gets 404, no session gets 401, and an agent cannot push into another site's camera. Also fixes a bug I introduced with it - the Live button was gated on `connected`, which is head office's last report and up to two minutes stale, so it hid itself during every reconnect. "Is that camera really down?" is exactly when somebody wants to look, and a hidden control says "you cannot" where the honest answer is "here is why". Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
304 lines
10 KiB
Go
304 lines
10 KiB
Go
// Package cameras keeps a shop PC's cameras in step with head office.
|
|
//
|
|
// The split is forced by the network, not by taste: only this PC is on the
|
|
// camera's LAN, so only this PC can connect to it — but the person onboarding a
|
|
// camera is often in an office somewhere else. So head office holds the DESIRED
|
|
// configuration and the agent pulls it.
|
|
//
|
|
// Pull, never push. A shop PC sits behind a router with no inbound route, so it
|
|
// has to ask; and asking makes the whole thing idempotent — a sync that fails
|
|
// halfway is fixed by the next one rather than leaving two systems disagreeing.
|
|
package cameras
|
|
|
|
import (
|
|
"context"
|
|
"log"
|
|
"time"
|
|
)
|
|
|
|
// Engine is the local recognition engine's camera API.
|
|
type Engine interface {
|
|
List(ctx context.Context) ([]Local, error)
|
|
Add(ctx context.Context, cam Local) error
|
|
Update(ctx context.Context, id string, cam Local) error
|
|
Remove(ctx context.Context, id string) error
|
|
Snapshot(ctx context.Context, id string) ([]byte, error)
|
|
}
|
|
|
|
// Cloud is head office.
|
|
type Cloud interface {
|
|
Desired(ctx context.Context) ([]Desired, error)
|
|
Report(ctx context.Context, rep Report) error
|
|
// UploadSnapshot puts a JPEG in object storage and names it. Used for the
|
|
// placement check's proof picture, which is transient.
|
|
UploadSnapshot(ctx context.Context, jpeg []byte) (key string, err error)
|
|
// PutSnapshot gets a camera's latest frame to head office by whichever
|
|
// route this deployment has - the bucket, or the server itself when there
|
|
// is none. An empty key means the server already stored it, so the state
|
|
// report has nothing to carry.
|
|
PutSnapshot(ctx context.Context, cameraID string, jpeg []byte) (key string, err error)
|
|
}
|
|
|
|
// Local is a camera as the engine holds it.
|
|
type Local struct {
|
|
ID string `json:"id"`
|
|
Label string `json:"label,omitempty"`
|
|
Host string `json:"host,omitempty"`
|
|
Port int `json:"port,omitempty"`
|
|
Path string `json:"path,omitempty"`
|
|
Username string `json:"username,omitempty"`
|
|
Password string `json:"password,omitempty"`
|
|
MaxWidth int `json:"max_width,omitempty"`
|
|
Tuning map[string]any `json:"tuning,omitempty"`
|
|
Connected bool `json:"connected"`
|
|
}
|
|
|
|
// Desired is a camera as head office holds it.
|
|
type Desired struct {
|
|
// ID is head office's uuid for this camera; CameraID is the name the
|
|
// engine on this PC knows it by. The live relay translates between them.
|
|
ID string `json:"id"`
|
|
CameraID string `json:"camera_id"`
|
|
Label string `json:"label"`
|
|
Host string `json:"host"`
|
|
Port int `json:"port"`
|
|
Path string `json:"path"`
|
|
Username string `json:"username"`
|
|
Password string `json:"password"`
|
|
MaxWidth int `json:"max_width"`
|
|
Tuning map[string]any `json:"tuning"`
|
|
Enabled bool `json:"enabled"`
|
|
Revision int64 `json:"revision"`
|
|
Deleted bool `json:"deleted"`
|
|
}
|
|
|
|
type State struct {
|
|
CameraID string `json:"camera_id"`
|
|
Connected bool `json:"connected"`
|
|
SnapshotKey string `json:"snapshot_key,omitempty"`
|
|
}
|
|
|
|
type Report struct {
|
|
State []State `json:"state,omitempty"`
|
|
Adopt []Desired `json:"adopt,omitempty"`
|
|
}
|
|
|
|
// Syncer reconciles the two, on a timer.
|
|
type Syncer struct {
|
|
Engine Engine
|
|
Cloud Cloud
|
|
Log *log.Logger
|
|
|
|
// Every how often to reconcile configuration. Cameras change rarely, and
|
|
// each sync is a database read on the server for every site in the estate,
|
|
// so this is minutes rather than seconds.
|
|
Interval time.Duration
|
|
// How often to send a fresh picture of each camera. A shop floor does not
|
|
// change much, and each frame is a few tens of kilobytes uploaded over the
|
|
// same connection the visits have to travel on.
|
|
SnapshotEvery time.Duration
|
|
|
|
// Checks and Prober are the "prove this camera works" half. Both nil on a
|
|
// PC that has never been claimed, and the syncer simply skips that work
|
|
// rather than treating it as a failure.
|
|
Checks Checks
|
|
Prober Prober
|
|
|
|
// applied remembers the revision last pushed into the engine, so an
|
|
// unchanged site costs one request and no engine calls at all.
|
|
applied map[string]int64
|
|
}
|
|
|
|
const (
|
|
DefaultInterval = 2 * time.Minute
|
|
DefaultSnapshotEvery = 60 * time.Second
|
|
)
|
|
|
|
// New builds a fully wired Syncer from the two clients every caller already
|
|
// has.
|
|
//
|
|
// It exists because the four fields were assembled by hand at each call site
|
|
// and both of them - the headless agent and the desktop app - set Engine and
|
|
// Cloud and forgot Checks and Prober. runChecks returns silently when either
|
|
// is nil (correct: an unclaimed PC has neither), so pressing "Test connection"
|
|
// at head office left the camera saying "checking..." until the five-minute
|
|
// stale release, and then said nothing at all. No error, on either side.
|
|
//
|
|
// The same two objects satisfy all four interfaces, so there was never a
|
|
// reason for a caller to choose.
|
|
func New(eng *EngineClient, cloud *CloudClient, log *log.Logger) *Syncer {
|
|
return &Syncer{Engine: eng, Cloud: cloud, Checks: cloud, Prober: eng, Log: log}
|
|
}
|
|
|
|
// Run reconciles until ctx is cancelled.
|
|
func (s *Syncer) Run(ctx context.Context) {
|
|
interval, snapEvery := s.Interval, s.SnapshotEvery
|
|
if interval <= 0 {
|
|
interval = DefaultInterval
|
|
}
|
|
if snapEvery <= 0 {
|
|
snapEvery = DefaultSnapshotEvery
|
|
}
|
|
// Immediately on start, so a PC that has just been claimed picks up its
|
|
// cameras now rather than in two minutes.
|
|
s.Once(ctx)
|
|
|
|
config := time.NewTicker(interval)
|
|
defer config.Stop()
|
|
snaps := time.NewTicker(snapEvery)
|
|
defer snaps.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-config.C:
|
|
s.Once(ctx)
|
|
case <-snaps.C:
|
|
s.report(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Once performs one full reconcile: pull desired, apply, then report back.
|
|
func (s *Syncer) Once(ctx context.Context) {
|
|
if s.applied == nil {
|
|
s.applied = map[string]int64{}
|
|
}
|
|
desired, err := s.Cloud.Desired(ctx)
|
|
if err != nil {
|
|
// Not fatal and not even unusual: an unclaimed PC has no credentials
|
|
// and a disconnected one has no network. The engine keeps running the
|
|
// cameras it already has, which is the whole point of the local store.
|
|
s.logf("camera sync: %v", err)
|
|
return
|
|
}
|
|
local, err := s.Engine.List(ctx)
|
|
if err != nil {
|
|
s.logf("camera sync: engine unavailable: %v", err)
|
|
return
|
|
}
|
|
|
|
have := map[string]Local{}
|
|
for _, c := range local {
|
|
have[c.ID] = c
|
|
}
|
|
known := map[string]bool{}
|
|
|
|
for _, d := range desired {
|
|
known[d.CameraID] = true
|
|
_, exists := have[d.CameraID]
|
|
|
|
switch {
|
|
case d.Deleted || !d.Enabled:
|
|
if exists {
|
|
if err := s.Engine.Remove(ctx, d.CameraID); err != nil {
|
|
s.logf("camera %s: remove failed: %v", d.CameraID, err)
|
|
continue
|
|
}
|
|
s.logf("camera %s removed (head office)", d.CameraID)
|
|
}
|
|
delete(s.applied, d.CameraID)
|
|
|
|
case !exists:
|
|
if err := s.Engine.Add(ctx, toLocal(d)); err != nil {
|
|
s.logf("camera %s: add failed: %v", d.CameraID, err)
|
|
continue
|
|
}
|
|
s.applied[d.CameraID] = d.Revision
|
|
s.logf("camera %s added from head office", d.CameraID)
|
|
|
|
case s.applied[d.CameraID] != d.Revision:
|
|
// The revision is what keeps this cheap. Without it every sync
|
|
// would PATCH every camera, and a PATCH restarts the connection —
|
|
// so a healthy site would drop its own video every two minutes.
|
|
if err := s.Engine.Update(ctx, d.CameraID, toLocal(d)); err != nil {
|
|
s.logf("camera %s: update failed: %v", d.CameraID, err)
|
|
continue
|
|
}
|
|
s.applied[d.CameraID] = d.Revision
|
|
s.logf("camera %s updated to revision %d", d.CameraID, d.Revision)
|
|
}
|
|
}
|
|
|
|
// Anything running here that head office has never heard of gets offered
|
|
// up. Without this, switching the feature on would delete every camera an
|
|
// existing site is already running — including the one it was commissioned
|
|
// with. The server refuses to overwrite its own config with these, and
|
|
// keeps tombstones, so a deleted camera is not resurrected.
|
|
var adopt []Desired
|
|
for id, c := range have {
|
|
if known[id] {
|
|
continue
|
|
}
|
|
adopt = append(adopt, Desired{
|
|
CameraID: id, Label: orElse(c.Label, id), Host: c.Host,
|
|
Port: c.Port, Path: c.Path, Username: c.Username,
|
|
Password: c.Password, MaxWidth: c.MaxWidth, Tuning: c.Tuning,
|
|
Enabled: true,
|
|
})
|
|
}
|
|
s.reportWith(ctx, adopt)
|
|
// Last, so a check requested against a camera added in the same breath
|
|
// finds it already applied. Inside Once() rather than beside it in the
|
|
// loop, so it also runs on startup and cannot be called twice a tick.
|
|
s.runChecks(ctx, desired)
|
|
}
|
|
|
|
func (s *Syncer) report(ctx context.Context) { s.reportWith(ctx, nil) }
|
|
|
|
// reportWith sends observed state, and a fresh picture from each camera.
|
|
func (s *Syncer) reportWith(ctx context.Context, adopt []Desired) {
|
|
local, err := s.Engine.List(ctx)
|
|
if err != nil {
|
|
s.logf("camera report: engine unavailable: %v", err)
|
|
return
|
|
}
|
|
rep := Report{Adopt: adopt}
|
|
for _, c := range local {
|
|
st := State{CameraID: c.ID, Connected: c.Connected}
|
|
if c.Connected {
|
|
// A snapshot failure never blocks the state report. Knowing a
|
|
// camera is down matters far more than having a picture of it,
|
|
// and the picture is the part most likely to fail.
|
|
if jpeg, err := s.Engine.Snapshot(ctx, c.ID); err == nil && len(jpeg) > 0 {
|
|
// An empty key is not a failure: it means this deployment has
|
|
// no object storage and the server stored the picture itself.
|
|
if key, err := s.Cloud.PutSnapshot(ctx, c.ID, jpeg); err == nil {
|
|
st.SnapshotKey = key
|
|
} else {
|
|
s.logf("camera %s: snapshot upload failed: %v", c.ID, err)
|
|
}
|
|
}
|
|
}
|
|
rep.State = append(rep.State, st)
|
|
}
|
|
if len(rep.State) == 0 && len(rep.Adopt) == 0 {
|
|
return
|
|
}
|
|
if err := s.Cloud.Report(ctx, rep); err != nil {
|
|
s.logf("camera report: %v", err)
|
|
}
|
|
}
|
|
|
|
func toLocal(d Desired) Local {
|
|
return Local{
|
|
ID: d.CameraID, Label: d.Label, Host: d.Host, Port: d.Port,
|
|
Path: d.Path, Username: d.Username, Password: d.Password,
|
|
MaxWidth: d.MaxWidth, Tuning: d.Tuning,
|
|
}
|
|
}
|
|
|
|
func orElse(s, fallback string) string {
|
|
if s == "" {
|
|
return fallback
|
|
}
|
|
return s
|
|
}
|
|
|
|
func (s *Syncer) logf(format string, args ...any) {
|
|
if s.Log != nil {
|
|
s.Log.Printf(format, args...)
|
|
}
|
|
}
|