Files
Behavision/desktop/app.go
Suriyakumarvijayanayagam 48a30d97db Live camera view in the app, and the green light that was lying about it
Two changes, and the second was found by verifying the first.

## Watching a camera from the app, in another building

Snapshots answer "is that camera working". They do not answer "what is
happening in my shop right now", which is what somebody who opens the app away
from the counter is asking. Head office's browser already had that answer -
LiveHub plus cameras.Live, where the shop PC asks outbound whether anybody is
watching and pushes JPEG frames for as long as somebody is - and the app could
not reach it.

cloud.CameraLive opens that feed and the app's own loopback relay re-emits it
as multipart MJPEG. That is the trick: frames arrive base64 over SSE, an <img>
cannot render that, and an <img> renders MJPEG natively - so a tile is an
ordinary <img> pointed at loopback whether the camera is in this room or
another city.

- Reconnecting happens in the relay, not the page. The server caps one push at
  five minutes, so doing it here means the <img> never sees the stream end.
- The headers are flushed before the first frame. Go writes them on the first
  body write, so without that the whole response waits for the shop PC to
  start pushing. Measured against production: 30 seconds and not even a
  Content-Type, which surfaces as the request timing out.
- One camera at a time. Watching makes a shop PC upload, so a grid that went
  live at once would put an estate's worth of cameras on the wire because
  somebody opened a page.
- live.mjpeg is behind the same per-run token as the engine routes, and a
  wrong token is a 404 that never reaches head office at all.
- CameraLive uses its own HTTP client: the shared one's 30s timeout covers the
  whole response and would sever a working view every thirty seconds - the
  trap that made the server set WriteTimeout to zero for its own SSE endpoint.

## A camera read "Connected" for 34 minutes after the shop PC went blind

Which is why the verification above looked like a failure: head office
registered the viewer and no frame ever came.

reportWith returns early when the engine is unreachable - correctly, it has
nothing to say - so the last state it sent stays in the database looking
current. Measured live: cam2 and entrance both reading Connected, in green,
with last_seen_at 34 minutes old, while the heartbeat from the same PC said
cameras_up 0 of 0. Two surfaces reading two stored fields and disagreeing.

false could not be the answer. It means "this camera is not connecting", which
sends an installer to check cabling on a camera that was working perfectly the
last time anybody could ask it. So there are four states and one function:

  connected       reported recently, and working
  not_connecting  reported recently, and the stream will not open
  waiting         no shop PC has ever reported this camera
  stale           reported once, and not lately

- Connected is CLEARED when stale or waiting. A stale true left in place stays
  available to every client reading the field directly, and leaves two fields
  on one object disagreeing - how the shops screen once came out labelled
  Working, in green, above "2 of 3 cameras not connecting".
- Computed in scanCamera, so every camera anybody reads passes through it. A
  state computed per handler is one a handler forgets, and this had already
  reached three screens.
- CameraStaleAfter is 5 minutes: five missed reports, not one. Same reasoning
  as three missed heartbeats - an indicator that cries wolf gets ignored.
- An unparseable last_seen_at is stale. It should be impossible, which is why
  it must not fall through to the state that says everything is fine.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-30 16:48:43 +05:30

967 lines
34 KiB
Go

// The methods bound to the frontend.
//
// Every one is a thin adapter: it talks to the local engine, the cloud, or the
// supervisor, and returns something JSON-shaped. No recognition logic lives
// here - the engine owns that, and duplicating any of it would give the UI a
// second opinion about who someone is.
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"os"
"os/exec"
"path/filepath"
"strings"
"sync"
"time"
agentbridge "github.com/loyaly/behavision-agent/pkg/bridge"
agentcameras "github.com/loyaly/behavision-agent/pkg/cameras"
agentcfg "github.com/loyaly/behavision-agent/pkg/config"
agentengine "github.com/loyaly/behavision-agent/pkg/engine"
"github.com/loyaly/behavision-agent/pkg/enrol"
agentmqtt "github.com/loyaly/behavision-agent/pkg/mqtt"
agentpaths "github.com/loyaly/behavision-agent/pkg/paths"
agentspool "github.com/loyaly/behavision-agent/pkg/spool"
"github.com/loyaly/behavision-desktop/internal/cloud"
"github.com/loyaly/behavision-desktop/internal/local"
)
type App struct {
ctx context.Context
mu sync.RWMutex
cfg agentcfg.Config
cloud *cloud.Client
local *local.Client
sup *agentengine.Supervisor
spool *agentspool.Spool
bridge *agentbridge.Bridge
broker *agentmqtt.Client
stopBridge func()
hookURL string
// Relays camera feeds to the webview so the engine's credential never has
// to travel in an <img> src, which a Chromium webview would strip anyway.
proxy *streamProxy
// Set once the operator logs in. Until then the UI shows the login sheet
// and nothing else is reachable.
onSessionChange func(bool)
}
func NewApp() *App {
cfg, _ := agentcfg.Load(agentpaths.AgentConfig())
// The engine invents its own Basic credential when none is configured,
// which is the default. Without this every call this app makes to the
// engine - health, cameras, the live feed - comes back 401, and the tray
// shows a healthy process the UI cannot talk to.
cfg = cfg.WithEngineCredentials(agentpaths.APICredentials())
base := cfg.APIBase
if base == "" {
base = "http://127.0.0.1:8010"
}
return &App{
cfg: cfg,
cloud: cloud.New(envOr("BEHAVISION_CLOUD", "https://mcp.loyaly.ai")),
local: localWithCreds(base, cfg),
proxy: newStreamProxy(),
}
}
func (a *App) startup(ctx context.Context) {
a.ctx = ctx
_ = agentpaths.EnsureState()
// Before any screen asks for a camera URL. A failure here is logged and
// not fatal: the rest of the app - people, cameras, the engine controls -
// works without a picture, and refusing to start over a broken tile would
// take a working shop offline.
if err := a.proxy.start(a.local.Base, a.local.User, a.local.Password); err != nil {
log.Printf("camera relay unavailable, tiles will not load: %v", err)
}
// And the other direction: watching a camera in another building, through
// head office's relay. Enabled unconditionally rather than only when a
// session already exists, because signing in is a thing that happens
// while the app is open - and CameraLive refuses without a session
// anyway, so there is nothing to gate.
if err := a.proxy.watchRemote(a.cloud.CameraLive); err != nil {
log.Printf("remote camera view unavailable: %v", err)
}
// A saved session means a shop PC that rebooted overnight comes back
// working instead of waiting for someone to log in.
if a.cfg.SessionToken != "" {
a.cloud.SetSession(cloud.Session{
Token: a.cfg.SessionToken, RefreshToken: a.cfg.SessionRefresh,
User: cloud.User{Email: a.cfg.SessionEmail},
})
}
// The server rotates the refresh token every time it is used, so a PC that
// refreshes and then reboots would come back holding one the server has
// already invalidated - it would look exactly like a normal expiry, twelve
// hours after anyone last touched the machine.
a.cloud.OnRefresh(func(s cloud.Session) { a.persistSession(s) })
exe := a.cfg.EngineExe
if exe != "" && !filepath.IsAbs(exe) {
exe = filepath.Join(agentpaths.InstallRoot(), exe)
}
logFile, _ := agentengine.LogFile(agentpaths.EngineLog())
a.sup = agentengine.New(agentengine.Options{
Command: func(c context.Context) *exec.Cmd {
cmd := exec.CommandContext(c, exe, a.cfg.EngineArgs...)
// Run the engine FROM a known directory rather than from whatever
// happened to launch us. A double-clicked bundle hands its child
// "/", and an engine invoked as `-m behavision` then cannot find
// itself - measured on macOS, where it retried forever.
cmd.Dir = a.cfg.EngineDir
if cmd.Dir == "" {
cmd.Dir = agentpaths.InstallRoot()
}
// How the engine learns where to post its detections. Its config
// already reads `events.webhook_url: ${BEHAVISION_WEBHOOK_URL}`,
// and python-dotenv does not override a variable the process
// already has, so this needs no new endpoint and no fixed port.
//
// Read here rather than captured, because the bridge picks its
// port after this closure is built and a restarted engine has to
// be told again. Without it the engine recognised people and the
// bridge received nothing: a claimed shop PC published heartbeats
// and zero visits.
cmd.Env = agentengine.ChildEnv(a.webhookURL())
return cmd
},
LogWriter: logFile,
HealthURL: strings.TrimRight(a.cfg.APIBase, "/") + "/api/health",
StatsURL: strings.TrimRight(a.cfg.APIBase, "/") + "/api/stats",
User: a.cfg.APIUser, Password: a.cfg.APIPassword,
})
a.startPipeline(ctx)
go a.watchConfig(ctx)
// Recognition starts with the app. Until this, the engine only ever
// started when somebody pressed Start - which meant a till that rebooted
// overnight came back with the window open, the tray icon showing, the
// session restored, and recognition off until a shop assistant noticed.
// That is the failure the tray colours exist to catch, and it should not
// be the default state every morning.
//
// Guarded on the interpreter actually being there: on a PC where setup has
// not run yet, starting the supervisor would loop on a missing executable
// with nothing useful to say. The Start button still exists for the one
// case where somebody has deliberately stopped it.
if _, err := os.Stat(exe); err == nil {
a.sup.Start()
} else {
log.Printf("engine not installed yet (%s); run behavision-setup, then Start", exe)
}
}
// webhookURL is the loopback address the bridge is listening on, or empty
// before it has started.
func (a *App) webhookURL() string {
a.mu.RLock()
defer a.mu.RUnlock()
return a.hookURL
}
// startPipeline connects detections to the server: the engine posts events to
// a loopback webhook, the bridge queues them durably, and the pump drains the
// queue to the broker. Without it the engine recognises people and nothing
// ever leaves the PC.
func (a *App) startPipeline(ctx context.Context) {
logger := log.New(os.Stdout, "", log.LstdFlags)
// A PC set up on its own has nothing to report to, and unlike an unclaimed
// one it never will. Queuing anyway would write up to SpoolMax visits to
// disk - each carrying a face template, which is biometric personal data -
// into a queue nothing is ever going to drain. Recognition, the gallery
// and the cameras are unaffected: they are the engine's, not the pump's.
//
// Deliberately distinct from the unclaimed case below, where the queue is
// exactly right: that PC is waiting for credentials, and its footfall from
// the day it was installed should survive until they arrive.
if a.cfg.Standalone && !a.cfg.Configured() {
logger.Print("standalone: recognition runs locally, nothing is reported")
go a.startLocalCameras(ctx, logger)
return
}
q, err := agentspool.Open(agentpaths.SpoolDir(), a.cfg.SpoolMax)
if err != nil {
logger.Printf("spool unavailable, detections will not be recorded: %v", err)
return
}
a.spool = q
// Created before the bridge and handed over unconditionally: an unclaimed
// PC has no pump reading it, which is harmless - the single slot fills
// once and later rings are dropped.
waker := agentmqtt.NewWaker()
a.bridge = &agentbridge.Bridge{
Queue: q,
Wake: waker.Wake,
Embeddings: agentbridge.NewEngineEmbeddings(
a.cfg.APIBase, a.cfg.APIUser, a.cfg.APIPassword),
TopicPrefix: topicPrefix(a.cfg),
Log: logger,
// Uploads face images through a URL the server mints, so this PC never
// holds bucket credentials. Harmless when the engine writes no images
// or the PC is not claimed: Upload reports "images off" and the visit
// queues without a photo.
Uploader: &agentbridge.SpacesUploader{
BaseURL: a.cfg.CloudBase, Token: a.cfg.AgentToken,
},
}
// Cameras, kept in step with head office. Started before the broker check
// because it does not need one: an unclaimed PC still reconciles (to
// nothing) and still keeps running its local cameras.
go a.startLocalCameras(ctx, logger)
url, stop, err := a.bridge.Listen(ctx)
if err != nil {
logger.Printf("event bridge failed to start: %v", err)
return
}
// Under the lock: the engine supervisor reads this from whichever goroutine
// launches the child, and Claim can run startPipeline again at any time.
a.mu.Lock()
a.stopBridge = stop
a.hookURL = url
a.mu.Unlock()
logger.Printf("event bridge on %s", url)
if !a.cfg.Configured() {
// Not claimed yet. The bridge still runs, so footfall from today is on
// disk waiting for the credentials rather than lost.
return
}
client, err := agentmqtt.NewClient(agentmqtt.ClientOptions{
BrokerURL: a.cfg.BrokerURL,
ClientID: "behavision-" + a.cfg.ClientID + "-" + a.cfg.SiteID,
Username: a.cfg.BrokerUsername, Password: a.cfg.BrokerPassword,
CAFile: a.cfg.BrokerCAFile, Log: logger,
})
if err != nil {
logger.Printf("broker unavailable, queuing locally: %v", err)
return
}
a.broker = client
go (&agentmqtt.Pump{
Queue: q, Publisher: client, Log: logger,
Wake: waker.C(),
HeartbeatTopic: topicPrefix(a.cfg) + "/heartbeat",
HeartbeatPayload: a.heartbeat,
}).Run(ctx)
logger.Print("broker pump running")
}
// startLocalCameras runs the reconciler that keeps this PC's cameras in step
// with head office. It is deliberately not conditional on being claimed: with
// no server to ask it reconciles against nothing and the locally configured
// cameras keep running, which is the whole of standalone operation.
func (a *App) startLocalCameras(ctx context.Context, logger *log.Logger) {
camUploader := &agentbridge.SpacesUploader{
BaseURL: a.cfg.CloudBase, Token: a.cfg.AgentToken,
}
camCloud := agentcameras.NewCloudClient(a.cfg.CloudBase, a.cfg.AgentToken)
camCloud.Upload = camUploader.UploadBytes
camEngine := agentcameras.NewEngineClient(
a.cfg.APIBase, a.cfg.APIUser, a.cfg.APIPassword)
// New(), not a struct literal: assembling the Syncer by hand here is how
// this app - and the headless agent - both ended up wiring configuration
// and forgetting the check runner, so "Test connection" at head office
// never completed on any shop PC.
// The live relay runs alongside the reconciler and uploads nothing until
// somebody at head office is actually watching a camera.
go agentcameras.NewLive(camEngine, camCloud, logger).Run(ctx)
agentcameras.New(camEngine, camCloud, logger).Run(ctx)
}
// topicPrefix must equal the MQTT username: the broker enforces
// `pattern write bv/%u/...`, so any other prefix is refused.
func topicPrefix(cfg agentcfg.Config) string {
if cfg.ClientID == "" || cfg.SiteID == "" {
return ""
}
return "bv/" + cfg.ClientID + "." + cfg.SiteID
}
func (a *App) heartbeat() []byte {
hb := map[string]any{"sent_at": time.Now().UTC().Format(time.RFC3339)}
if a.spool != nil {
hb["queued"] = a.spool.Len()
// Non-zero means this site's queue overflowed and it genuinely lost
// footfall. Reported rather than inferred from a dip in a graph.
hb["dropped"] = a.spool.Dropped()
}
s := a.EngineStatus()
hb["engine_state"] = s.State
if s.Model != "" {
hb["recognition_model"] = s.Model
}
if s.Cameras != nil {
hb["cameras"] = s.Cameras
}
b, _ := json.Marshal(hb)
return b
}
// PipelineStatus is what the UI shows about the link to head office.
type PipelineStatus struct {
WebhookURL string `json:"webhook_url"`
Queued int `json:"queued"`
Dropped uint64 `json:"dropped"`
Claimed bool `json:"claimed"`
// Standalone separates "nothing is being sent because this PC is set up on
// its own" from "nothing is being sent and something is wrong". They look
// identical from the counters alone, and only one of them is a fault.
Standalone bool `json:"standalone"`
BrokerUp bool `json:"broker_up"`
Accepted uint64 `json:"accepted"`
}
func (a *App) PipelineStatus() PipelineStatus {
out := PipelineStatus{
WebhookURL: a.hookURL,
Claimed: a.cfg.Configured(),
Standalone: a.cfg.Standalone && !a.cfg.Configured(),
}
if a.spool != nil {
out.Queued, out.Dropped = a.spool.Len(), a.spool.Dropped()
}
if a.bridge != nil {
out.Accepted = a.bridge.Accepted
}
if a.broker != nil {
out.BrokerUp = a.broker.Connected()
}
return out
}
// ---------------------------------------------------------------- session --
type SessionInfo struct {
LoggedIn bool `json:"logged_in"`
User cloud.User `json:"user"`
SiteName string `json:"site_name"`
Claimed bool `json:"claimed"`
// Standalone is a PC deliberately run on its own. The UI then shows only
// the screens that work without head office - the cameras and what this
// PC is seeing - rather than a sign-in form for an account that does not
// exist.
Standalone bool `json:"standalone"`
}
func (a *App) Session() SessionInfo {
a.mu.RLock()
defer a.mu.RUnlock()
return SessionInfo{
LoggedIn: a.cloud.LoggedIn(),
User: a.cloud.User(),
SiteName: a.cfg.SiteName,
Claimed: a.cfg.Configured(),
Standalone: a.cfg.Standalone && !a.cfg.Configured(),
}
}
// RunStandalone sets this PC up on its own, with no head office.
//
// Recognition, the cameras and the local gallery all work without a server -
// they always did - so refusing to open the app until somebody issues an
// enrolment code held the product hostage to a component it does not need. The
// choice is persisted because it has to survive a reboot, and it is reversible:
// Claim still works afterwards and clears the flag.
func (a *App) RunStandalone() (SessionInfo, error) {
a.mu.Lock()
a.cfg.Standalone = true
err := a.cfg.Save(agentpaths.AgentConfig())
a.mu.Unlock()
if err != nil {
// A choice that is not on disk works until the next restart and then
// silently is not made any more, which looks exactly like the app
// forgetting the setup step was ever done.
return SessionInfo{}, fmt.Errorf("could not save this choice: %w", err)
}
return a.Session(), nil
}
func (a *App) Login(email, password string) (SessionInfo, error) {
ctx, cancel := context.WithTimeout(a.ctx, 30*time.Second)
defer cancel()
sess, err := a.cloud.Login(ctx, email, password)
if err != nil {
return SessionInfo{}, err
}
a.persistSession(sess)
if a.onSessionChange != nil {
a.onSessionChange(true)
}
return a.Session(), nil
}
// persistSession writes the tokens to the DPAPI-protected config. Called on
// sign-in and on every silent refresh, so the two can never diverge.
func (a *App) persistSession(s cloud.Session) {
a.mu.Lock()
defer a.mu.Unlock()
a.cfg.SessionToken = s.Token
a.cfg.SessionRefresh = s.RefreshToken
if s.User.Email != "" {
a.cfg.SessionEmail = s.User.Email
}
_ = a.cfg.Save(agentpaths.AgentConfig())
}
// Claim links this PC to a shop, using the one-shot code an operator is given.
//
// This is the half of onboarding that had no way to happen. The server has had
// POST /api/agent/enrol since enrolment was built and `cloud.Client.Bootstrap`
// has existed to call it - and nothing called it, so a freshly installed PC
// displayed "Not linked to head office" and offered no way to link it. The
// only route was hand-editing a JSON file on a shop counter.
//
// Deliberately NOT session-authenticated, mirroring the endpoint: the person
// standing at a new shop PC has no account on it yet, and requiring a login
// first would mean shipping a password to every shop that installs the
// software.
func (a *App) Claim(code string) (SessionInfo, error) {
code = strings.TrimSpace(code)
if code == "" {
return SessionInfo{}, errors.New("type the installation code you were given")
}
ctx, cancel := context.WithTimeout(a.ctx, 30*time.Second)
defer cancel()
b, err := a.cloud.Bootstrap(ctx, code)
if err != nil {
return SessionInfo{}, err
}
a.mu.Lock()
// The slugs, not the uuids: the topic prefix is <client>.<site> and the
// broker's ACL is written against exactly that username.
a.cfg.ClientID = b.ClientSlug
a.cfg.SiteID = b.SiteSlug
a.cfg.SiteName = b.SiteName
a.cfg.BrokerURL = b.MQTTURL
a.cfg.BrokerUsername = b.MQTTUser
a.cfg.BrokerPassword = b.MQTTPass
a.cfg.AgentToken = b.AgentToken
a.cfg.CloudBase = a.cloud.Base
// A new head office: whoever was signed in was signed in somewhere else.
a.cfg.SessionToken, a.cfg.SessionRefresh, a.cfg.SessionEmail = "", "", ""
a.cloud.Clear()
caPath, err := enrol.SaveCA(b.CACert, agentpaths.BrokerCA())
if err != nil {
return SessionInfo{}, err
}
a.cfg.BrokerCAFile = caPath
// A PC that was running on its own and has now been linked is no longer
// standalone. Leaving the flag set would keep the head-office screens
// hidden on the one machine that just earned them.
a.cfg.Standalone = false
err = a.cfg.Save(agentpaths.AgentConfig())
a.mu.Unlock()
if err != nil {
// Reported, not swallowed. A claim that is not on disk works until the
// next restart and then silently is not claimed any more, which looks
// like the code was wrong when it was not.
return SessionInfo{}, fmt.Errorf("could not save the settings: %w", err)
}
// The pipeline was started unclaimed: no broker, no pump. Restart it so
// this PC begins publishing now rather than at the next launch - an
// installer who has to reboot to finish setting up will assume it failed.
a.restartPipeline()
return a.Session(), nil
}
// restartPipeline tears the bridge and broker down and builds them again from
// the current config. Only Claim needs it today; it exists as its own method
// because "stop everything that reads the config, then start it" is the part
// that is easy to get half right.
// watchConfig reloads agent.json when something else writes it.
//
// behavision-setup re-run on a PC with the app open re-claims the shop and
// rotates its API token; the running app kept the old one and every camera
// sync was refused from then on - heartbeats still flowed, so head office
// looked fine while the cameras went stale. A claim from `behavision-agent
// claim` does the same. Rather than ask people to restart the app, the app
// watches the file and picks the new credentials up itself.
func (a *App) watchConfig(ctx context.Context) {
path := agentpaths.AgentConfig()
last := mtime(path)
t := time.NewTicker(10 * time.Second)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
}
now := mtime(path)
if now.IsZero() || now.Equal(last) {
continue
}
last = now
fresh, err := agentcfg.Load(path)
if err != nil {
continue
}
fresh = fresh.WithEngineCredentials(agentpaths.APICredentials())
a.mu.Lock()
changed := fresh.AgentToken != a.cfg.AgentToken || fresh.SiteID != a.cfg.SiteID ||
fresh.BrokerPassword != a.cfg.BrokerPassword || fresh.CloudBase != a.cfg.CloudBase ||
fresh.Standalone != a.cfg.Standalone
if changed {
// Keep this process's live session; a claim clears it in the file
// deliberately, and that is honoured too.
a.cfg = fresh
if fresh.SessionToken == "" {
a.cloud.Clear()
}
}
a.mu.Unlock()
if changed {
a.restartPipeline()
}
}
}
func mtime(path string) time.Time {
st, err := os.Stat(path)
if err != nil {
return time.Time{}
}
return st.ModTime()
}
func (a *App) restartPipeline() {
if a.stopBridge != nil {
a.stopBridge()
a.stopBridge = nil
}
if a.broker != nil {
a.broker.Close()
a.broker = nil
}
a.startPipeline(a.ctx)
}
func (a *App) Logout() SessionInfo {
// Revoke server-side too. Clearing only the local copy leaves a live token
// on a machine somebody is about to hand back or resell.
ctx, cancel := context.WithTimeout(a.ctx, 10*time.Second)
defer cancel()
_ = a.cloud.Logout(ctx)
a.mu.Lock()
a.cfg.SessionToken, a.cfg.SessionRefresh, a.cfg.SessionEmail = "", "", ""
_ = a.cfg.Save(agentpaths.AgentConfig())
a.mu.Unlock()
if a.onSessionChange != nil {
a.onSessionChange(false)
}
return a.Session()
}
// ---------------------------------------------------------------- engine ---
type EngineStatus struct {
State string `json:"state"`
Error string `json:"error,omitempty"`
Restarts int `json:"restarts"`
Reachable bool `json:"reachable"`
Model string `json:"recognition_model,omitempty"`
Cameras map[string]bool `json:"cameras,omitempty"`
// Progress is the first-run model download, when one is happening.
Progress *agentengine.Progress `json:"progress,omitempty"`
}
func (a *App) EngineStatus() EngineStatus {
out := EngineStatus{State: "stopped"}
if a.sup == nil {
return out
}
st, err := a.sup.State()
out.State = string(st)
out.Restarts = a.sup.Restarts()
if p := a.sup.Progress(); p.What != "" {
out.Progress = &p
}
if err != nil {
out.Error = err.Error()
}
ctx, cancel := context.WithTimeout(a.ctx, 4*time.Second)
defer cancel()
// A running process is not a working engine: on a memory-starved box the
// large model loses the fallback chain and the process stays up regardless,
// so the UI reports which encoder actually loaded.
if h, herr := a.sup.Health(ctx); herr == nil {
out.Reachable = true
out.Model = h.RecognitionModel
out.Cameras = h.Cameras
}
return out
}
func (a *App) StartEngine() EngineStatus {
if a.sup != nil {
a.sup.Start()
}
return a.EngineStatus()
}
func (a *App) StopEngine() EngineStatus {
if a.sup != nil {
a.sup.Stop()
}
return a.EngineStatus()
}
// ---------------------------------------------------------------- cameras --
// Cameras lists this PC's cameras, or the company's if this PC has none of
// its own.
//
// The distinction is load-bearing and the UI is told which it got. A camera
// from the local engine is one THIS machine can reach, edit and stream. One
// from head office is a camera at a shop somewhere else: it has a snapshot
// and a connection state, and it cannot be edited from here because the shop
// PC on that LAN is the only thing that can reach it. Offering an Edit button
// that could not work would be worse than not showing the camera at all.
func (a *App) Cameras() ([]map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
defer cancel()
cams, err := a.local.Cameras(ctx)
if err == nil {
return cams, nil
}
if !a.cloud.LoggedIn() {
return nil, err
}
remote, rerr := a.cloud.RemoteCameras(ctx)
if rerr != nil {
return nil, err // the local failure is the one worth reporting
}
out := make([]map[string]any, 0, len(remote))
for _, c := range remote {
out = append(out, map[string]any{
"id": c.ID, "camera_id": c.CameraID, "label": c.Label,
"site": c.Site, "enabled": c.Enabled,
"connected": c.Connected, "last_seen_at": c.LastSeenAt,
"state": c.State, "state_note": c.StateNote,
"snapshot": c.Snapshot, "snapshot_at": c.SnapshotAt,
// What the screen keys off to hide Edit, Test and Check: this
// camera is on a network this PC cannot reach.
"remote": true,
})
}
return out, nil
}
// DiscoverCameras lists the cameras on this PC's network, so the add-camera
// form is a pick-list and not a request for an IP address nobody knows.
func (a *App) DiscoverCameras() (map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 30*time.Second)
defer cancel()
return a.local.DiscoverCameras(ctx)
}
func (a *App) TestCamera(cam map[string]any) (map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 60*time.Second)
defer cancel()
return a.local.TestCamera(ctx, cam)
}
func (a *App) SaveCamera(id string, cam map[string]any) (map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 60*time.Second)
defer cancel()
if id == "" {
return a.local.AddCamera(ctx, cam)
}
return a.local.UpdateCamera(ctx, id, cam)
}
func (a *App) DeleteCamera(id string) error {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.local.DeleteCamera(ctx, id)
}
// StartPlacementCheck begins the guided commissioning walk. This is the step
// that stops a site being signed off with a camera that recognises nobody.
func (a *App) StartPlacementCheck(id string, seconds float64) (map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
defer cancel()
return a.local.StartPlacementCheck(ctx, id, seconds)
}
func (a *App) PlacementResult(id string) (map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
defer cancel()
return a.local.PlacementResult(ctx, id)
}
// StreamURL is the MJPEG endpoint for a camera tile.
//
// It points at this app's own loopback relay, not at the engine directly. The
// previous version put the engine's Basic credentials inline in the URL, with
// a comment saying they were there "so an <img> tag can load it" - which a
// browser will not do. Chromium strips credentials from subresource URLs, and
// WebView2 is Chromium, so every camera tile on a shop PC was a broken image.
// See stream_proxy.go for the measurement.
//
// The relay is also why no password appears in the page any more. If it is not
// running the fallback is the bare engine URL with no credential: correct for
// an engine configured without auth, and for one with auth a tile that fails
// to load rather than a password sitting in the DOM.
func (a *App) StreamURL(cameraID string) string {
if u := a.proxy.urlFor(cameraID, "stream.mjpeg"); u != "" {
return u
}
base := strings.TrimPrefix(strings.TrimPrefix(a.local.Base, "http://"), "https://")
return fmt.Sprintf("http://%s/api/cameras/%s/stream.mjpeg", base, cameraID)
}
// RemoteStreamURL is the live view of a camera in another building.
//
// The picture comes from head office's relay - the shop PC pushes frames
// outbound because nothing can reach in - and this app re-emits them as MJPEG
// on its own loopback, so a tile is an ordinary <img> either way. A screen
// therefore never has to know which building it is looking at.
//
// Empty when the relay is not running, and the caller shows the last snapshot
// instead. There is no useful fallback URL: the head-office endpoint needs
// this session's bearer, which an <img> cannot send.
func (a *App) RemoteStreamURL(cameraID string) string {
if !a.cloud.LoggedIn() {
return ""
}
return a.proxy.urlFor(cameraID, "live.mjpeg")
}
// ------------------------------------------------------------------- live --
type LiveSnapshot struct {
Stats map[string]any `json:"stats"`
Events []map[string]any `json:"events"`
// Viewing is true when none of this came from an engine on THIS PC. The
// screen must say so: the numbers are the company's, not this machine's,
// and a laptop in a hotel showing "2 cameras live" without that word
// would be claiming to be watching a shop it cannot see.
Viewing bool `json:"viewing"`
}
// Live is what the shop PC sees, and falls back to what HEAD OFFICE sees.
//
// A PC with no engine is not necessarily broken - it is somebody signed in on
// a laptop away from the shop, which is the ordinary way an owner looks at
// their estate. Until now that produced "engine not reachable at
// 127.0.0.1:8010", an accurate sentence and a useless one when the reader was
// never expecting an engine on that machine.
//
// The local engine always wins when it is there: it is this shop's own
// ground truth and it is live rather than a heartbeat old.
func (a *App) Live() (LiveSnapshot, error) {
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
defer cancel()
stats, err := a.local.Stats(ctx)
if err == nil {
events, eerr := a.local.Events(ctx, 40)
if eerr == nil {
return LiveSnapshot{Stats: stats, Events: events}, nil
}
}
// No engine here. If nobody is signed in either, the honest answer is
// still the local error - there is nothing else to show and the person
// is most likely setting this PC up.
if !a.cloud.LoggedIn() {
return LiveSnapshot{}, err
}
return a.liveFromCloud(ctx)
}
// liveFromCloud builds the same shape the Live screen already renders, out of
// the estate's own feed, so the view needs no second code path.
func (a *App) liveFromCloud(ctx context.Context) (LiveSnapshot, error) {
sites, err := a.cloud.Sites(ctx)
if err != nil {
return LiveSnapshot{}, err
}
arrivals, err := a.cloud.Arrivals(ctx, 40)
if err != nil {
return LiveSnapshot{}, err
}
// The counters are summed across the estate, and fraction_below_gate
// takes the WORST site rather than an average - one badly placed camera
// is a hole in the numbers, and averaging it against three good ones
// hides the only site anyone needs to visit. Same rule the heartbeat
// already follows.
var up, total int
worst := 0.0
people := map[string]struct{}{}
for _, s := range sites {
up, total = up+s.CamerasUp, total+s.CamerasTotal
if s.FractionBelowGate > worst {
worst = s.FractionBelowGate
}
}
events := make([]map[string]any, 0, len(arrivals))
for _, v := range arrivals {
if v.VisitorID != "" {
people[v.VisitorID] = struct{}{}
}
events = append(events, map[string]any{
"type": map[bool]string{true: "person.new", false: "person.seen"}[v.IsNew],
"ts": v.OccurredAt, "camera_id": v.CameraID,
"data": map[string]any{
"label": v.Label, "ref": v.Ref, "site": v.Site,
"similarity": v.Similarity, "attributes": v.Attributes,
},
})
}
return LiveSnapshot{
Viewing: true,
Events: events,
Stats: map[string]any{
"cameras": []map[string]any{},
"gallery": map[string]any{"identities": len(people), "sightings": len(arrivals)},
"cameras_up": up, "cameras_total": total,
"fraction_below_gate": worst,
},
}, nil
}
// ---------------------------------------------------------------- reports --
func (a *App) Footfall(from, to, bucket string) (cloud.FootfallReport, error) {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.Footfall(ctx, from, to, bucket)
}
func (a *App) Sales(from, to string) (cloud.SalesReport, error) {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.Sales(ctx, from, to)
}
func (a *App) Customers(query string, limit int) ([]cloud.Customer, error) {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
if limit <= 0 {
limit = 100
}
return a.cloud.Customers(ctx, query, limit)
}
// Sites is the health of every store this account can see. It is what makes
// "no customers today" distinguishable from "that shop's PC has been unplugged
// for a week" - two identical rows of zeroes with completely different answers.
func (a *App) Sites() ([]cloud.SiteHealth, error) {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.Sites(ctx)
}
// VisitorHistory is one customer's timeline, for the customer record screen.
// Ask is the help panel. It needs head office: the assistant runs there,
// against this company's own data, as this signed-in user. A PC running on
// its own has nobody to ask, and the panel says so rather than erroring.
func (a *App) Ask(history []cloud.AssistantTurn) (cloud.AssistantAnswer, error) {
ctx, cancel := context.WithTimeout(a.ctx, 90*time.Second)
defer cancel()
return a.cloud.Ask(ctx, history)
}
func (a *App) VisitorHistory(id string, limit int) ([]cloud.Visit, error) {
if limit <= 0 {
limit = 100
}
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.VisitorHistory(ctx, id, limit)
}
// VisitorPhoto returns a link to this customer's face image, valid for a few
// minutes. "There is no photo" comes back as a Photo with Available false and
// a sentence explaining why, not as an error - see cloud.Photo.
func (a *App) VisitorPhoto(id string) (cloud.Photo, error) {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.VisitorImage(ctx, id)
}
// ForgetCustomer erases a person at the request of that person.
//
// Bound as its own method rather than folded into SaveProfile because it is
// not an edit: it destroys the face template, the photo and the profile, and
// cannot be undone.
func (a *App) ForgetCustomer(id string) error {
// Longer than the usual 20s: the server deletes every stored image from
// object storage before it touches the database, and refuses the whole
// request if any one of them fails.
ctx, cancel := context.WithTimeout(a.ctx, 60*time.Second)
defer cancel()
return a.cloud.ForgetVisitor(ctx, id)
}
func (a *App) SaveProfile(p cloud.Profile) error {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.SaveProfile(ctx, p)
}
func (a *App) RecordPurchase(visitorID string, amount float64,
items []string, notes string) error {
ctx, cancel := context.WithTimeout(a.ctx, 20*time.Second)
defer cancel()
return a.cloud.RecordPurchase(ctx, visitorID, amount, items, notes)
}
// --------------------------------------------------------------- identity --
// LocalIdentities reads the engine's own gallery. Shown alongside the cloud
// customer list because they answer different questions: this is who this PC
// can recognise right now, that is who the business knows.
func (a *App) LocalIdentities(limit int) ([]map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
defer cancel()
if limit <= 0 {
limit = 50
}
return a.local.Identities(ctx, limit)
}
func (a *App) LocalSightings(limit int) ([]map[string]any, error) {
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
defer cancel()
if limit <= 0 {
limit = 50
}
return a.local.Sightings(ctx, limit)
}
func envOr(key, def string) string {
if v := osGetenv(key); v != "" {
return v
}
return def
}
// localWithCreds builds the engine client with a credential resolver, so a
// first run - where the engine writes its credential after the app has looked
// for it - recovers by itself instead of 401ing for the life of the process.
func localWithCreds(base string, cfg agentcfg.Config) *local.Client {
c := local.New(base, cfg.APIUser, cfg.APIPassword)
c.Creds = agentcfg.NewCreds(agentpaths.APICredentials(), cfg.APIUser, cfg.APIPassword)
return c
}