Files
Behavision/agent/main.go
Suriyakumarvijayanayagam 92573e9067 The installer, run on a clean machine, found two bugs in itself
Ran behavision-setup in a fresh Linux container: Python 3.12, nothing
else, the release contents mounted read-only the way Program Files or a
shared drive would be. It failed, and then it failed differently, and
both failures would have been the client's first experience.

1. `pip install <folder>` makes setuptools write behavision.egg-info
   INTO the folder. The folder is read-only wherever a release is
   sensibly unzipped, so: "could not create 'behavision.egg-info':
   Read-only file system". The release now ships a wheel - pure Python,
   buildable anywhere, nothing to build on the shop PC, and pip never
   touches the unzipped folder. Source stays as a fallback and is copied
   somewhere writable first.

2. The engine's paths.py knows two worlds - frozen (ProgramData) and a
   checkout (the repo root) - and a pip-installed engine is neither. It
   resolved its state root to site-packages: database there, camera
   list there, and its generated API credential in a folder the app
   never reads, while the app looked in ProgramData. Every call would be
   401 on a stock install, with nothing in either log saying why. The
   same disease as the Mac checkout two days ago, now in production
   shape.

   engine.ChildEnv is the one place the engine's environment is built,
   used by the desktop app, the headless agent and the installer's own
   smoke test. It passes BEHAVISION_DATA_DIR = this process's state
   root, which paths.py honours ahead of every other rule, so the two
   halves agree by construction however the engine was installed.

   It also seeds config/default.yaml into the state root: a package in
   site-packages has no config beside it to seed from.

Re-run on the same clean container: seven steps, all pass, models
downloaded, engine started and answered, and its data/ landed beside
agent.json - not in site-packages.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-11 12:26:54 +05:30

386 lines
13 KiB
Go

// Command behavision-agent is the Go half of the edge install: it supervises
// the Python recognition engine and moves its events to the server.
//
// Modes:
//
// run supervise the engine and drain the spool (what the tray runs)
// status one-shot health report, for support and for the installer
// paths where this agent thinks state lives
//
// The tray and the Wails UI wrap this; none of the logic below assumes a
// window exists, so `run` works headless over SSH or from a scheduled task.
package main
import (
"context"
"encoding/json"
"flag"
"fmt"
"log"
"os"
"os/exec"
"os/signal"
"path/filepath"
"strings"
"syscall"
"time"
"github.com/loyaly/behavision-agent/pkg/bridge"
"github.com/loyaly/behavision-agent/pkg/cameras"
"github.com/loyaly/behavision-agent/pkg/config"
"github.com/loyaly/behavision-agent/pkg/engine"
"github.com/loyaly/behavision-agent/pkg/enrol"
"github.com/loyaly/behavision-agent/pkg/mqtt"
"github.com/loyaly/behavision-agent/pkg/paths"
"github.com/loyaly/behavision-agent/pkg/spool"
)
var version = "dev"
func main() {
flag.Usage = func() {
fmt.Fprintf(os.Stderr, `behavision-agent %s
usage: %s <command>
run supervise the engine and report to head office (default)
claim <code> link this PC to a shop, using an installation code
status what this PC is and whether it is claimed
paths where this install reads and writes
`, version, filepath.Base(os.Args[0]))
}
flag.Parse()
mode := "run"
if flag.NArg() > 0 {
mode = flag.Arg(0)
}
var err error
switch mode {
case "run":
err = cmdRun()
case "status":
err = cmdStatus()
case "claim":
err = cmdClaim(flag.Args()[1:])
case "paths":
err = cmdPaths()
default:
flag.Usage()
os.Exit(2)
}
if err != nil {
log.Fatalf("behavision-agent: %v", err)
}
}
// cmdClaim is the headless half of onboarding.
//
// The desktop app has had a Setup screen for this; a back-office PC with no
// window had nothing at all, so the only way to claim one was to hand-edit
// agent.json - which is the state that screen was built to end.
func cmdClaim(args []string) error {
if len(args) == 0 {
return fmt.Errorf("usage: behavision-agent claim <installation code>\n" +
"Ask whoever manages your shops for one - they can create it from\n" +
"the Behavision platform, under the shop.")
}
// Joined rather than requiring quotes: the code is printed in groups for
// reading aloud, and an operator pasting it will paste the spaces too.
code := strings.Join(args, "")
if err := paths.EnsureState(); err != nil {
return err
}
cfg, err := config.Load(paths.AgentConfig())
if err != nil {
return err
}
base := cfg.CloudBase
if v := os.Getenv("BEHAVISION_CLOUD"); v != "" {
base = v
}
if base == "" {
base = "https://mcp.loyaly.ai"
}
b, err := enrol.Claim(context.Background(), base, code)
if err != nil {
return err
}
// The slugs, not the uuids: the topic prefix is <client>.<site> and the
// broker's ACL is written against exactly that username.
cfg.ClientID = b.ClientSlug
cfg.SiteID = b.SiteSlug
cfg.SiteName = b.SiteName
cfg.BrokerURL = b.MQTTURL
cfg.BrokerUsername = b.MQTTUser
cfg.BrokerPassword = b.MQTTPass
cfg.AgentToken = b.AgentToken
cfg.CloudBase = base
// A PC that was running on its own and has now been linked is no longer
// standalone.
cfg.Standalone = false
if err := cfg.Save(paths.AgentConfig()); err != nil {
// Reported, never swallowed: a claim that is not on disk works until
// the next restart and then silently is not claimed any more, which
// looks exactly like a wrong code.
return fmt.Errorf("could not save the settings: %w", err)
}
fmt.Printf("linked to %s (%s.%s)\n", b.SiteName, b.ClientSlug, b.SiteSlug)
fmt.Printf("settings written to %s\n", paths.AgentConfig())
fmt.Println("restart the agent for it to take effect.")
return nil
}
func cmdPaths() error {
return json.NewEncoder(os.Stdout).Encode(map[string]string{
"version": version,
"install_root": paths.InstallRoot(),
"state_root": paths.StateRoot(),
"agent_config": paths.AgentConfig(),
"spool": paths.SpoolDir(),
"engine_log": paths.EngineLog(),
})
}
func cmdStatus() error {
cfg, err := config.Load(paths.AgentConfig())
if err != nil {
return err
}
// The engine invents its own Basic credential when none is configured,
// which is the default. Reading it here is what stops every call the agent
// makes to the engine coming back 401 on a stock install.
cfg = cfg.WithEngineCredentials(paths.APICredentials())
q, err := spool.Open(paths.SpoolDir(), cfg.SpoolMax)
if err != nil {
return err
}
sup := engine.New(engine.Options{
Command: func(context.Context) *exec.Cmd { return nil },
HealthURL: strings.TrimRight(cfg.APIBase, "/") + "/api/health",
StatsURL: strings.TrimRight(cfg.APIBase, "/") + "/api/stats",
User: cfg.APIUser, Password: cfg.APIPassword,
})
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
out := map[string]any{
"version": version,
"configured": cfg.Configured(),
"secrets_protected": config.SecretsProtected(),
"queued": q.Len(),
"dropped": q.Dropped(),
}
if h, err := sup.Health(ctx); err != nil {
out["engine"] = map[string]any{"reachable": false, "error": err.Error()}
} else {
out["engine"] = h
}
enc := json.NewEncoder(os.Stdout)
enc.SetIndent("", " ")
return enc.Encode(out)
}
func cmdRun() error {
logger := log.New(os.Stdout, "", log.LstdFlags|log.LUTC)
if err := paths.EnsureState(); err != nil {
return err
}
cfg, err := config.Load(paths.AgentConfig())
if err != nil {
return err
}
// The engine invents its own Basic credential when none is configured,
// which is the default. Reading it here is what stops every call the agent
// makes to the engine coming back 401 on a stock install.
cfg = cfg.WithEngineCredentials(paths.APICredentials())
// Opened before the engine starts: detections arriving in the first second
// must have somewhere to land.
q, err := spool.Open(paths.SpoolDir(), cfg.SpoolMax)
if err != nil {
return fmt.Errorf("spool: %w", err)
}
logFile, err := engine.LogFile(paths.EngineLog())
if err != nil {
return err
}
defer logFile.Close()
exe := cfg.EngineExe
if !filepath.IsAbs(exe) {
// Resolved against the install root, not the working directory: a
// service or a shortcut can start us anywhere.
exe = filepath.Join(paths.InstallRoot(), exe)
}
// hookURL is read when the engine is LAUNCHED, not when the supervisor is
// built, because the bridge has not picked its port yet and because a
// restarted engine has to be told again.
var hookURL string
sup := engine.New(engine.Options{
Command: func(ctx context.Context) *exec.Cmd {
cmd := exec.CommandContext(ctx, exe, cfg.EngineArgs...)
// How the engine learns where to send detections. The engine's
// 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.
//
// Without it the engine recognised people and the bridge received
// nothing - the URL was returned, logged and even exposed on the
// desktop's status object, and never actually given to the engine.
// A claimed shop PC published heartbeats and zero visits.
cmd.Env = engine.ChildEnv(hookURL)
return cmd
},
LogWriter: logFile,
HealthURL: strings.TrimRight(cfg.APIBase, "/") + "/api/health",
StatsURL: strings.TrimRight(cfg.APIBase, "/") + "/api/stats",
User: cfg.APIUser, Password: cfg.APIPassword,
})
ctx, stop := signal.NotifyContext(context.Background(),
os.Interrupt, syscall.SIGTERM)
defer stop()
// The bridge always runs, claimed or not: a PC that is set up before its
// tenant credentials arrive must still record the footfall it sees, and
// the spool is what holds it until the broker is configured.
// Created before the bridge and handed over unconditionally. On a PC that
// is not claimed yet there is no pump reading it, which costs nothing: the
// waker's single slot fills once and every later ring is dropped.
waker := mqtt.NewWaker()
br := &bridge.Bridge{
Queue: q,
Wake: waker.Wake,
Embeddings: bridge.NewEngineEmbeddings(cfg.APIBase, cfg.APIUser, cfg.APIPassword),
TopicPrefix: topicPrefix(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: &bridge.SpacesUploader{
BaseURL: cfg.CloudBase, Token: cfg.AgentToken,
},
}
// Cameras, kept in step with head office. Runs whether or not this PC is
// claimed: unclaimed it simply logs that it has no credentials yet, and the
// engine carries on with the cameras already in its own store.
uploader := &bridge.SpacesUploader{BaseURL: cfg.CloudBase, Token: cfg.AgentToken}
cloud := cameras.NewCloudClient(cfg.CloudBase, cfg.AgentToken)
cloud.Upload = uploader.UploadBytes
eng := cameras.NewEngineClient(cfg.APIBase, cfg.APIUser, cfg.APIPassword)
go cameras.New(eng, cloud, logger).Run(ctx)
// The live relay, which uploads nothing until somebody at head office is
// actually watching a camera.
go cameras.NewLive(eng, cloud, logger).Run(ctx)
// Before the engine starts, so the engine can be launched already knowing
// where to post its detections.
//
// A PC set up to run on its own is the one case where it should not: it
// has nothing to report to and, unlike an unclaimed one, never will, so
// queuing would write up to SpoolMax visits - each carrying a face
// template, which is biometric personal data - into a queue nothing is
// going to drain. Recognition and the cameras are unaffected; they belong
// to the engine, not the pump.
if cfg.Standalone && !cfg.Configured() {
logger.Print("standalone: recognition runs locally, nothing is reported")
} else {
url, stopBridge, err := br.Listen(ctx)
if err != nil {
return fmt.Errorf("event bridge: %w", err)
}
defer stopBridge()
hookURL = url
logger.Printf("event bridge listening on %s", hookURL)
}
logger.Printf("starting engine: %s %s", exe, strings.Join(cfg.EngineArgs, " "))
sup.Start()
if cfg.Configured() {
logger.Printf("tenant %s / site %s; broker %s",
cfg.ClientID, cfg.SiteID, cfg.BrokerURL)
client, err := mqtt.NewClient(mqtt.ClientOptions{
BrokerURL: cfg.BrokerURL,
ClientID: "behavision-" + cfg.ClientID + "-" + cfg.SiteID,
Username: cfg.BrokerUsername, Password: cfg.BrokerPassword,
CAFile: cfg.BrokerCAFile, Log: logger,
})
if err != nil {
// Not fatal. Events keep accumulating on disk and go out when the
// link returns - which is the entire point of the spool.
logger.Printf("broker unavailable, queuing locally: %v", err)
} else {
defer client.Close()
pump := &mqtt.Pump{
Queue: q, Publisher: client, Log: logger,
Wake: waker.C(),
HeartbeatTopic: topicPrefix(cfg) + "/heartbeat",
HeartbeatPayload: func() []byte {
return heartbeat(q, sup)
},
}
go pump.Run(ctx)
logger.Print("broker pump running")
}
} else {
logger.Print("not claimed by a tenant yet - recording locally only")
}
<-ctx.Done()
logger.Print("stopping engine")
sup.Stop()
return nil
}
// topicPrefix is the site's MQTT namespace. The broker enforces
// `pattern write bv/%u/...`, so this must equal the credential's username or
// every publish is refused.
func topicPrefix(cfg config.Config) string {
if cfg.ClientID == "" || cfg.SiteID == "" {
return ""
}
return "bv/" + cfg.ClientID + "." + cfg.SiteID
}
// heartbeat says the site is alive and what shape it is in.
//
// `dropped` matters most: non-zero means this site's queue overflowed and it
// genuinely lost footfall the customer paid for. Reporting it is the only way
// that becomes visible rather than being inferred from a dip in a graph.
func heartbeat(q *spool.Spool, sup *engine.Supervisor) []byte {
hb := map[string]any{
"sent_at": time.Now().UTC().Format(time.RFC3339),
"agent_version": version,
"queued": q.Len(),
"dropped": q.Dropped(),
}
if sup != nil {
state, _ := sup.State()
hb["engine_state"] = string(state)
hctx, cancel := context.WithTimeout(context.Background(), 4*time.Second)
defer cancel()
if h, err := sup.Health(hctx); err == nil {
hb["recognition_model"] = h.RecognitionModel
hb["cameras"] = h.Cameras
}
// The share of faces this site's cameras saw and discarded before they
// could become visits. It is the difference between "a quiet week" and
// "the camera is pointed at the ceiling", which are the same row of
// numbers on a footfall report without it. Measured on the Office1
// camera it was 0.727.
if st, err := sup.Stats(hctx); err == nil {
if worst, ok := st.WorstBelowGate(); ok {
hb["fraction_below_gate"] = worst
}
}
}
b, _ := json.Marshal(hb)
return b
}