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
292 lines
11 KiB
Go
292 lines
11 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/loyaly/behavision-server/internal/api"
|
|
)
|
|
|
|
// ErrNoSecrets is the API package's sentinel, aliased rather than redeclared.
|
|
//
|
|
// Two variables with the same text would compare unequal under errors.Is, so
|
|
// the handler's check would silently fall through to a 500 - the failure this
|
|
// error exists to replace with a sentence an operator can act on.
|
|
var ErrNoSecrets = api.ErrNoSecrets
|
|
|
|
const cameraCols = `
|
|
c.id::text, c.site_id::text, si.name, c.camera_id, c.label,
|
|
c.host, c.port, c.path, c.username, (c.password_enc IS NOT NULL),
|
|
c.max_width, c.tuning, c.enabled, c.revision,
|
|
c.connected, c.last_seen_at, c.snapshot_key, c.snapshot_at,
|
|
c.check_kind, c.check_requested_at, c.check_started_at, c.check_finished_at,
|
|
c.check_seconds, c.check_result, c.check_image_key`
|
|
|
|
func scanCamera(row pgx.Row) (api.Camera, error) {
|
|
var c api.Camera
|
|
var lastSeen, snapAt *time.Time
|
|
var snapKey string
|
|
var checkKind *string
|
|
var reqAt, startAt, finAt *time.Time
|
|
var checkSeconds int
|
|
var checkResult []byte
|
|
var checkImage string
|
|
if err := row.Scan(&c.ID, &c.SiteID, &c.Site, &c.CameraID, &c.Label,
|
|
&c.Host, &c.Port, &c.Path, &c.Username, &c.HasPassword,
|
|
&c.MaxWidth, &c.Tuning, &c.Enabled, &c.Revision,
|
|
&c.Connected, &lastSeen, &snapKey, &snapAt,
|
|
&checkKind, &reqAt, &startAt, &finAt,
|
|
&checkSeconds, &checkResult, &checkImage); err != nil {
|
|
return c, err
|
|
}
|
|
kind := ""
|
|
if checkKind != nil {
|
|
kind = *checkKind
|
|
}
|
|
// Carried on the camera rather than fetched separately: "is this camera set
|
|
// up" and "has anyone proved it works" are the same question to the person
|
|
// asking, and two requests to answer it is two chances for the screen to
|
|
// show a camera and its verdict from different moments.
|
|
c.Check = checkOf(kind, reqAt, startAt, finAt, checkSeconds, checkResult, checkImage)
|
|
if lastSeen != nil {
|
|
c.LastSeenAt = lastSeen.UTC().Format(time.RFC3339)
|
|
}
|
|
if snapAt != nil {
|
|
c.SnapshotAt = snapAt.UTC().Format(time.RFC3339)
|
|
}
|
|
// The KEY travels in ImageKey, which is json:"-", and the handler swaps it
|
|
// for a signed link. Same rule as an arrival's face.
|
|
c.Snapshot.Key = snapKey
|
|
return c, nil
|
|
}
|
|
|
|
// Cameras lists a tenant's cameras, optionally for one site.
|
|
//
|
|
// Never returns a password, and structurally cannot: the column is not in the
|
|
// select list at all, only whether it is set.
|
|
func (s *Store) Cameras(ctx context.Context, clientID, siteID string) ([]api.Camera, error) {
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT `+cameraCols+`
|
|
FROM site_cameras c
|
|
JOIN sites si ON si.id = c.site_id
|
|
WHERE c.client_id = $1 AND c.deleted_at IS NULL
|
|
AND ($2 = '' OR c.site_id = $2::uuid)
|
|
ORDER BY si.name, c.camera_id`, clientID, siteID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []api.Camera
|
|
for rows.Next() {
|
|
c, err := scanCamera(rows)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, c)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// SaveCamera creates or updates one camera and bumps its revision.
|
|
//
|
|
// The revision bump is what makes the agent's reconcile cheap: it compares one
|
|
// integer instead of diffing every field, so a sync on an unchanged site costs
|
|
// a single query and no engine calls.
|
|
func (s *Store) SaveCamera(ctx context.Context, clientID, siteID, cameraID string,
|
|
in api.CameraInput) (api.Camera, error) {
|
|
|
|
var out api.Camera
|
|
if in.Password != nil && *in.Password != "" && s.secrets == nil {
|
|
return out, ErrNoSecrets
|
|
}
|
|
|
|
// Site must belong to this tenant. Checked in SQL rather than trusted from
|
|
// the path: a site id is caller-supplied and this would otherwise write a
|
|
// camera into somebody else's shop.
|
|
var owns bool
|
|
if err := s.pool.QueryRow(ctx,
|
|
`SELECT EXISTS (SELECT 1 FROM sites WHERE id = $1::uuid AND client_id = $2::uuid)`,
|
|
siteID, clientID).Scan(&owns); err != nil {
|
|
return out, err
|
|
}
|
|
if !owns {
|
|
return out, pgx.ErrNoRows
|
|
}
|
|
|
|
var sealed []byte
|
|
if in.Password != nil && *in.Password != "" {
|
|
// Sealed with the SITE id as additional data, so a row copied between
|
|
// sites in the database does not decrypt into a working credential.
|
|
b, err := s.secrets.SealString(*in.Password, siteID)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
sealed = b
|
|
}
|
|
|
|
// COALESCE on every field: a nil pointer means "leave this alone". An
|
|
// operator editing a label must not blank the password, and the form does
|
|
// not send one because the API never gave it back.
|
|
row := s.pool.QueryRow(ctx, `
|
|
INSERT INTO site_cameras (client_id, site_id, camera_id, label, host, port,
|
|
path, username, password_enc, max_width, tuning, enabled)
|
|
VALUES ($1::uuid, $2::uuid, $3,
|
|
COALESCE($4, ''), COALESCE($5, ''), COALESCE($6, 554),
|
|
COALESCE($7, '/'), COALESCE($8, ''), $9,
|
|
COALESCE($10, 1280), COALESCE($11, '{}'::jsonb), COALESCE($12, true))
|
|
ON CONFLICT (site_id, camera_id) DO UPDATE SET
|
|
label = COALESCE($4, site_cameras.label),
|
|
host = COALESCE($5, site_cameras.host),
|
|
port = COALESCE($6, site_cameras.port),
|
|
path = COALESCE($7, site_cameras.path),
|
|
username = COALESCE($8, site_cameras.username),
|
|
password_enc = COALESCE($9, site_cameras.password_enc),
|
|
max_width = COALESCE($10, site_cameras.max_width),
|
|
tuning = COALESCE($11, site_cameras.tuning),
|
|
enabled = COALESCE($12, site_cameras.enabled),
|
|
revision = site_cameras.revision + 1,
|
|
updated_at = now(),
|
|
-- Re-saving a deleted camera revives it. An operator adding back a
|
|
-- camera they removed should get their camera, not a unique-key
|
|
-- error about a row they cannot see.
|
|
deleted_at = NULL
|
|
RETURNING id`, clientID, siteID, cameraID,
|
|
in.Label, in.Host, in.Port, in.Path, in.Username, sealed,
|
|
in.MaxWidth, in.Tuning, in.Enabled)
|
|
|
|
var id string
|
|
if err := row.Scan(&id); err != nil {
|
|
return out, err
|
|
}
|
|
return s.CameraByID(ctx, clientID, id)
|
|
}
|
|
|
|
func (s *Store) CameraByID(ctx context.Context, clientID, id string) (api.Camera, error) {
|
|
return scanCamera(s.pool.QueryRow(ctx, `
|
|
SELECT `+cameraCols+`
|
|
FROM site_cameras c
|
|
JOIN sites si ON si.id = c.site_id
|
|
WHERE c.client_id = $1 AND c.id = $2::uuid AND c.deleted_at IS NULL`,
|
|
clientID, id))
|
|
}
|
|
|
|
// DeleteCamera tombstones a camera.
|
|
//
|
|
// A tombstone rather than a DELETE, because the agent adopts cameras it finds
|
|
// configured on the shop PC. A hard delete here would be undone on the next
|
|
// sync by the very camera the operator just removed - and they would have no
|
|
// idea why it kept coming back.
|
|
func (s *Store) DeleteCamera(ctx context.Context, clientID, id string) (api.Camera, error) {
|
|
cam, err := s.CameraByID(ctx, clientID, id)
|
|
if err != nil {
|
|
return cam, err
|
|
}
|
|
_, err = s.pool.Exec(ctx, `
|
|
UPDATE site_cameras
|
|
SET deleted_at = now(), revision = revision + 1, updated_at = now()
|
|
WHERE client_id = $1 AND id = $2::uuid`, clientID, id)
|
|
return cam, err
|
|
}
|
|
|
|
// AgentCameras is the desired configuration for one site, WITH passwords.
|
|
//
|
|
// The only route that decrypts them, and it is reachable only with that site's
|
|
// own agent token. Deleted cameras are included, flagged: the agent cannot
|
|
// distinguish "head office removed this" from "head office has not seen this
|
|
// yet" by absence, and would re-adopt what was just deleted.
|
|
func (s *Store) AgentCameras(ctx context.Context, siteID string) ([]api.AgentCamera, error) {
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT id::text, camera_id, label, host, port, path, username, password_enc,
|
|
max_width, tuning, enabled, revision, (deleted_at IS NOT NULL)
|
|
FROM site_cameras
|
|
WHERE site_id = $1::uuid
|
|
ORDER BY camera_id`, siteID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []api.AgentCamera
|
|
for rows.Next() {
|
|
var c api.AgentCamera
|
|
var sealed []byte
|
|
if err := rows.Scan(&c.ID, &c.CameraID, &c.Label, &c.Host, &c.Port, &c.Path,
|
|
&c.Username, &sealed, &c.MaxWidth, &c.Tuning, &c.Enabled,
|
|
&c.Revision, &c.Deleted); err != nil {
|
|
return nil, err
|
|
}
|
|
if len(sealed) > 0 && s.secrets != nil {
|
|
// A password that will not decrypt is sent as empty rather than
|
|
// failing the whole sync: one unreadable camera must not stop the
|
|
// other three being configured. The agent reports the connection
|
|
// failure, which is the symptom an operator can actually act on.
|
|
if pw, err := s.secrets.OpenString(sealed, siteID); err == nil {
|
|
c.Password = pw
|
|
} else {
|
|
s.auditFailed("camera password decrypt", err)
|
|
}
|
|
}
|
|
out = append(out, c)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ApplyAgentReport records what a shop PC observes, and adopts any camera it
|
|
// is running that head office does not know about.
|
|
//
|
|
// Adoption is what makes turning this on safe. Every existing site already has
|
|
// cameras configured locally - including the office camera this was tested with
|
|
// - and a reconcile that only pushed downwards would delete all of them on
|
|
// first sync.
|
|
func (s *Store) ApplyAgentReport(ctx context.Context, clientID, siteID string,
|
|
rep api.AgentCameraReport) error {
|
|
|
|
tx, err := s.pool.Begin(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer tx.Rollback(ctx) //nolint:errcheck
|
|
|
|
for _, cam := range rep.Adopt {
|
|
var sealed []byte
|
|
if cam.Password != "" && s.secrets != nil {
|
|
if b, err := s.secrets.SealString(cam.Password, siteID); err == nil {
|
|
sealed = b
|
|
}
|
|
}
|
|
// DO NOTHING on conflict, deliberately. Adoption must never overwrite
|
|
// head office's configuration with what the shop PC happens to hold -
|
|
// that would make an edit here silently revert on the next sync. It
|
|
// only fills in cameras nobody has configured centrally, tombstones
|
|
// included, so a deleted camera stays deleted.
|
|
if _, err := tx.Exec(ctx, `
|
|
INSERT INTO site_cameras (client_id, site_id, camera_id, label, host,
|
|
port, path, username, password_enc,
|
|
max_width, tuning, enabled)
|
|
VALUES ($1::uuid, $2::uuid, $3, $4, $5, $6, $7, $8, $9, $10,
|
|
COALESCE($11, '{}'::jsonb), $12)
|
|
ON CONFLICT (site_id, camera_id) DO NOTHING`,
|
|
clientID, siteID, cam.CameraID, cam.Label, cam.Host, cam.Port,
|
|
cam.Path, cam.Username, sealed, cam.MaxWidth, cam.Tuning,
|
|
cam.Enabled); err != nil {
|
|
return fmt.Errorf("adopt %q: %w", cam.CameraID, err)
|
|
}
|
|
}
|
|
|
|
for _, st := range rep.State {
|
|
if _, err := tx.Exec(ctx, `
|
|
UPDATE site_cameras
|
|
SET connected = $3, last_seen_at = now(),
|
|
snapshot_key = CASE WHEN $4 = '' THEN snapshot_key ELSE $4 END,
|
|
snapshot_at = CASE WHEN $4 = '' THEN snapshot_at ELSE now() END
|
|
WHERE site_id = $1::uuid AND camera_id = $2`,
|
|
siteID, st.CameraID, st.Connected, st.SnapshotKey); err != nil {
|
|
return fmt.Errorf("state %q: %w", st.CameraID, err)
|
|
}
|
|
}
|
|
return tx.Commit(ctx)
|
|
}
|