Files
Behavision/agent/pkg/cameras/clients.go
Suriyakumarvijayanayagam dad04e8cda Behavision: face recognition for retail, edge to head office
Five components that ship as one product:

- behavision/  the recognition engine. RTSP ingest, YuNet detection, IoU
               tracking, ArcFace embeddings, a FAISS/SQLite gallery, and a
               FastAPI dashboard. Identity is decided once per TRACK from an
               average of at least three embeddings, never per frame.
- agent/       the Go edge agent: supervises the engine, holds a durable
               spool, and drains it to MQTT. Nothing is acked before the
               broker confirms.
- desktop/     the shop PC application (Wails + React + tray).
- server/      the cloud API, MQTT consumer, reports and assistant.
- web/         platform.loyaly.ai, the head-office app, embedded in the
               server binary.

The gallery stores 512-float embeddings and timestamps - no images unless
`app.store_faces` is switched on. Those embeddings are biometric personal
data under GDPR and India's DPDP: template inversion reconstructs a
recognisable face from an ArcFace vector, so data/behavision.db is treated
as a biometric database and DELETE /api/visitors/{id} is a real erasure.

CLAUDE.md carries the reasoning behind every non-obvious decision here,
including the ones that were measured and the ones that were wrong first.

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

264 lines
8.3 KiB
Go

package cameras
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"strings"
"time"
)
// EngineClient talks to the recognition engine on this PC's loopback.
type EngineClient struct {
Base string
User string
Password string
Client *http.Client
}
func NewEngineClient(base, user, password string) *EngineClient {
return &EngineClient{
Base: strings.TrimRight(base, "/"), User: user, Password: password,
// Generous, because adding a camera makes the engine dial it, and a
// wrong address takes the full RTSP timeout to fail. Shorter than that
// and every genuinely-bad camera looks like an engine fault instead.
Client: &http.Client{Timeout: 30 * time.Second},
}
}
func (e *EngineClient) do(ctx context.Context, method, path string, body, out any) error {
var rdr io.Reader
if body != nil {
b, err := json.Marshal(body)
if err != nil {
return err
}
rdr = bytes.NewReader(b)
}
req, err := http.NewRequestWithContext(ctx, method, e.Base+path, rdr)
if err != nil {
return err
}
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
if e.User != "" {
req.SetBasicAuth(e.User, e.Password)
}
resp, err := e.Client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
// The engine's message, not just a status. "camera stored but failed to
// start: connection refused" is something an operator can act on;
// "500" is not, and this string ends up in the agent log a support
// engineer reads.
msg, _ := io.ReadAll(io.LimitReader(resp.Body, 2048))
return fmt.Errorf("engine %s: %s", resp.Status, strings.TrimSpace(string(msg)))
}
if out == nil {
return nil
}
return json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(out)
}
func (e *EngineClient) List(ctx context.Context) ([]Local, error) {
var out []Local
return out, e.do(ctx, http.MethodGet, "/api/cameras", nil, &out)
}
func (e *EngineClient) Add(ctx context.Context, cam Local) error {
return e.do(ctx, http.MethodPost, "/api/cameras", cam, nil)
}
func (e *EngineClient) Update(ctx context.Context, id string, cam Local) error {
// The engine takes the id from the path on PATCH and refuses it in the
// body, so it is cleared here rather than at the call site.
cam.ID = ""
return e.do(ctx, http.MethodPatch, "/api/cameras/"+url.PathEscape(id), cam, nil)
}
func (e *EngineClient) Remove(ctx context.Context, id string) error {
return e.do(ctx, http.MethodDelete, "/api/cameras/"+url.PathEscape(id), nil, nil)
}
// Snapshot fetches the most recent frame the engine holds.
//
// Not a fresh capture: the engine already keeps the latest frame in memory for
// its own MJPEG stream, so this costs a memory copy rather than a camera round
// trip. A camera that has not produced a frame yet answers 503, which is a
// normal state on a just-added camera and not an error worth logging loudly.
func (e *EngineClient) Snapshot(ctx context.Context, id string) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet,
e.Base+"/api/cameras/"+url.PathEscape(id)+"/frame.jpg", nil)
if err != nil {
return nil, err
}
if e.User != "" {
req.SetBasicAuth(e.User, e.Password)
}
resp, err := e.Client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("engine %s", resp.Status)
}
// Bounded. A frame is tens of kilobytes; anything near this cap means
// something other than a JPEG is on the other end.
return io.ReadAll(io.LimitReader(resp.Body, 8<<20))
}
// CloudClient talks to head office with this agent's own token.
type CloudClient struct {
Base string
Token string
Client *http.Client
// Upload is the existing image path: the server mints a presigned URL and
// the agent PUTs to it. Reused rather than reimplemented, so a shop PC
// still never holds bucket credentials — the reason that path exists.
Upload func(ctx context.Context, jpeg []byte) (string, error)
}
func NewCloudClient(base, token string) *CloudClient {
return &CloudClient{
Base: strings.TrimRight(base, "/"), Token: token,
Client: &http.Client{Timeout: 20 * time.Second},
}
}
func (c *CloudClient) do(ctx context.Context, method, path string, body, out any) error {
if c.Token == "" {
// An unclaimed PC. Said plainly, because this is the normal state
// between installing the software and typing an enrolment code, and it
// must not read as a fault in the log.
return fmt.Errorf("this PC is not claimed by a company yet")
}
var rdr io.Reader
if body != nil {
b, err := json.Marshal(body)
if err != nil {
return err
}
rdr = bytes.NewReader(b)
}
req, err := http.NewRequestWithContext(ctx, method, c.Base+path, rdr)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+c.Token)
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
resp, err := c.Client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("head office: %s", resp.Status)
}
if out == nil || resp.StatusCode == http.StatusNoContent {
return nil
}
return json.NewDecoder(io.LimitReader(resp.Body, 4<<20)).Decode(out)
}
func (c *CloudClient) Desired(ctx context.Context) ([]Desired, error) {
var body struct {
Cameras []Desired `json:"cameras"`
}
if err := c.do(ctx, http.MethodGet, "/api/agent/cameras", nil, &body); err != nil {
return nil, err
}
return body.Cameras, nil
}
func (c *CloudClient) Report(ctx context.Context, rep Report) error {
return c.do(ctx, http.MethodPost, "/api/agent/cameras", rep, nil)
}
func (c *CloudClient) UploadSnapshot(ctx context.Context, jpeg []byte) (string, error) {
if c.Upload == nil {
return "", fmt.Errorf("images are not enabled for this deployment")
}
return c.Upload(ctx, jpeg)
}
// ---------------------------------------------------------------- probing
// Test opens the candidate stream once, without saving it.
//
// The engine does this as a sync handler in its threadpool because
// cv2.VideoCapture blocks hard, and it checks TCP reachability first - so a
// wrong address, which is the single most likely thing anybody types, answers
// in milliseconds rather than the ~75 s an FFmpeg connect takes to give up.
func (e *EngineClient) Test(ctx context.Context, cam Local) (TestResult, error) {
var out TestResult
// A generous ceiling: the engine's own deadline is 12 s for the frame plus
// 3 s to connect, and cutting it off earlier would report a timeout of our
// own making as if it were the camera's.
ctx, cancel := context.WithTimeout(ctx, 45*time.Second)
defer cancel()
return out, e.do(ctx, http.MethodPost, "/api/cameras/test", cam, &out)
}
// Placement runs the engine's commissioning watch and returns its report whole.
//
// Polled rather than awaited: the engine starts the watch and answers
// immediately, so the run survives this request being retried, and the report
// arrives with `running: true` until it does not.
func (e *EngineClient) Placement(ctx context.Context, cameraID string, seconds int) (
map[string]any, error) {
path := "/api/cameras/" + url.PathEscape(cameraID) + "/commission"
var report map[string]any
if err := e.do(ctx, http.MethodPost, path,
map[string]any{"seconds": seconds}, &report); err != nil {
return nil, err
}
deadline := time.Now().Add(time.Duration(seconds+20) * time.Second)
for time.Now().Before(deadline) {
select {
case <-ctx.Done():
return report, ctx.Err()
case <-time.After(2 * time.Second):
}
var latest map[string]any
if err := e.do(ctx, http.MethodGet, path, nil, &latest); err != nil {
// Keep the last good report rather than losing the whole run to
// one failed poll - the engine may simply have been busy.
continue
}
report = latest
if running, _ := latest["running"].(bool); !running {
return report, nil
}
}
return report, nil
}
// ---------------------------------------------------------------- check jobs
func (c *CloudClient) Pending(ctx context.Context) ([]Job, error) {
var body struct {
Checks []Job `json:"checks"`
}
if err := c.do(ctx, http.MethodGet, "/api/agent/checks", nil, &body); err != nil {
return nil, err
}
return body.Checks, nil
}
func (c *CloudClient) Submit(ctx context.Context, res Result) error {
return c.do(ctx, http.MethodPost, "/api/agent/checks", res, nil)
}