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) }