Files
Behavision/agent/pkg/cameras/clients.go
Suriyakumarvijayanayagam 16f0e69cec A fresh shop PC could never authenticate to its own engine
The agent read the engine's generated credential file once, at startup.
On a brand new install that file does not exist yet: the agent starts the
engine, and the engine writes its credential seconds later. So the agent
held an empty credential for the life of the process and every call it
makes - health, stats, camera sync, the embedding for a visit - came back
401, with a tray showing a red engine that was running perfectly.

Measured on a fresh state directory today: three 401s, no camera ever
reconciled, and the engine left running the YAML-seeded main stream
instead of the sub-stream head office holds. The install script hid this
on Windows because setup runs the engine once before the app starts.

config.Creds resolves lazily and re-reads on a rejection; the camera
client, the supervisor and the desktop app's engine client all retry once
when it changes. A configured BEHAVISION_API_USER is never re-read - an
operator who set one means it. Tests pin the actual first-run ordering.

Also adds demo/, a one-screen live console for showing the whole chain:
camera, the six steps with a measured camera-to-cloud latency, the
customer editable in place, and the raw JSON a phone and a dashboard
receive from production side by side.

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

363 lines
12 KiB
Go

package cameras
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"strconv"
"strings"
"time"
"github.com/loyaly/behavision-agent/pkg/bridge"
"github.com/loyaly/behavision-agent/pkg/config"
)
// EngineClient talks to the recognition engine on this PC's loopback.
type EngineClient struct {
Base string
User string
Password string
// Creds re-reads the engine's generated credential when one is rejected.
// Without it a fresh install is 401 for the life of the process: the agent
// starts the engine, and the engine writes its credential file seconds
// after the agent has already read (and failed to find) it.
Creds *config.Creds
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")
}
user, pass := e.User, e.Password
if e.Creds != nil {
user, pass = e.Creds.Get()
}
if user != "" {
req.SetBasicAuth(user, pass)
}
resp, err := e.Client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusUnauthorized && e.Creds != nil && e.Creds.Refresh() {
// The engine generated its credential after we last looked. Read it
// and try once more rather than failing for the life of the process.
resp.Body.Close()
return e.do(ctx, method, path, body, out)
}
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) {
return e.Frame(ctx, id, 0, 0)
}
// Frame fetches the latest frame, optionally re-encoded smaller.
//
// The live relay asks for ~640 px at quality 60 - about a third the bytes of
// the full frame - because it sends several a second up a shop's uplink, where
// the snapshot sends one a minute and can afford the detail. The engine does
// the re-encode: it already has OpenCV open and the frame in memory, and
// shipping a scaler into the agent to redo that would be the same work twice.
func (e *EngineClient) Frame(ctx context.Context, id string, width, quality int) ([]byte, error) {
q := url.Values{}
if width > 0 {
q.Set("width", strconv.Itoa(width))
}
if quality > 0 {
q.Set("quality", strconv.Itoa(quality))
}
target := e.Base + "/api/cameras/" + url.PathEscape(id) + "/frame.jpg"
if len(q) > 0 {
target += "?" + q.Encode()
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, target, 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)
}
// PutSnapshot gets a camera's latest frame to head office by whichever route
// that deployment has.
//
// The bucket first: a presigned PUT goes straight to object storage and never
// passes through the API, which is what makes it the right route at estate
// scale. When there is no bucket the picture goes to the server itself, which
// stores one row per camera. Without this second route head office reported
// "This system is not storing images" for every camera forever, on the screen
// whose entire job is to show the camera.
//
// The returned key is empty for the direct route - there is no object to name -
// and the server records the picture as it stores it, so the state report has
// nothing to carry.
func (c *CloudClient) PutSnapshot(ctx context.Context, cameraID string, jpeg []byte) (string, error) {
if c.Upload != nil {
key, err := c.Upload(ctx, jpeg)
if err == nil {
return key, nil
}
// A bucket that is configured here but disabled at the server is the
// ordinary case on a self-hosted install: fall through rather than
// giving up, and let the direct route decide.
if !isImagesDisabled(err) {
return "", err
}
}
return "", c.putSnapshotDirect(ctx, cameraID, jpeg)
}
func (c *CloudClient) putSnapshotDirect(ctx context.Context, cameraID string, jpeg []byte) error {
req, err := http.NewRequestWithContext(ctx, http.MethodPut,
c.Base+"/api/agent/cameras/"+url.PathEscape(cameraID)+"/snapshot",
bytes.NewReader(jpeg))
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+c.Token)
req.Header.Set("Content-Type", "image/jpeg")
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)
}
return nil
}
// isImagesDisabled recognises the server saying it has no object storage.
//
// A sentinel, not a string match on the message: this decides whether to take a
// completely different route, and getting it wrong from prose that somebody
// later rewords would silently stop every camera picture in the estate.
func isImagesDisabled(err error) bool {
return errors.Is(err, bridge.ErrImagesOff)
}
// ---------------------------------------------------------------- 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)
}