The server has always sent the broker's CA certificate in the enrolment response, precisely so it never has to ship in an installer. Nothing on the receiving end wrote it anywhere: the agent read the field under the wrong name (ca_pem, the server says ca_cert) and the desktop app read it correctly and dropped it. Every claimed PC therefore dialled tls://mcp.loyaly.ai:8883 with the system trust store, the private CA failed verification, and the agent reported 'the broker did not accept this PC' - a TLS failure is indistinguishable from a refusal at that layer. No real site could ever have published a visit. Found by claiming this Mac as a real shop against production; fixed by writing the CA to broker-ca.crt beside agent.json on both claim paths. Verified: broker connected over TLS, camera pushed from head office, engine streaming it. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
736 lines
26 KiB
Go
736 lines
26 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: local.New(base, cfg.APIUser, cfg.APIPassword),
|
|
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)
|
|
}
|
|
|
|
// 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...)
|
|
// 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)
|
|
|
|
// 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
|
|
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.
|
|
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"`
|
|
}
|
|
|
|
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 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 --
|
|
|
|
func (a *App) Cameras() ([]map[string]any, error) {
|
|
ctx, cancel := context.WithTimeout(a.ctx, 15*time.Second)
|
|
defer cancel()
|
|
return a.local.Cameras(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)
|
|
}
|
|
|
|
// ------------------------------------------------------------------- live --
|
|
|
|
type LiveSnapshot struct {
|
|
Stats map[string]any `json:"stats"`
|
|
Events []map[string]any `json:"events"`
|
|
}
|
|
|
|
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 {
|
|
return LiveSnapshot{}, err
|
|
}
|
|
events, err := a.local.Events(ctx, 40)
|
|
if err != nil {
|
|
return LiveSnapshot{}, err
|
|
}
|
|
return LiveSnapshot{Stats: stats, Events: events}, 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.
|
|
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
|
|
}
|