Files
Behavision/agent/main.go
Suriyakumarvijayanayagam 0b29dd4a50 The engine inherited whatever directory launched the app
Nothing ever set the child's working directory, so it took the parent's -
and an app started by double-clicking its bundle is handed "/", not
anywhere useful. On macOS the symptom was
`python: No module named behavision` repeating forever, because the dev
engine is invoked as `-m behavision` and that resolves against the
working directory.

The same app launched from a terminal inside the repo worked perfectly,
which is exactly the shape of a bug that survives every test a developer
runs. It only appeared when the app was started the way a user starts
one.

Config.EngineDir, empty meaning the install root, set by both launchers -
the desktop app and the headless agent, which had identical code and the
identical omission. It matters beyond this case: the shipped Windows
engine is a one-folder PyInstaller build whose relative paths should
resolve beside itself rather than beside Explorer's idea of a current
directory.

Verified by double-clicking the bundle with nothing in the environment:
engine up on 8010 (401, gated), w600k_r50 on CoreML, gallery 5/5
embeddings usable and none stranded, both office cameras connected and
streaming, and head office reporting cameras 2/2 one heartbeat later.

Two things that showed up while proving it, both the product being
honest rather than faults:

- The camera at .121 was genuinely unreachable for several minutes, and
  last_error said so in words an installer can act on - "cannot reach
  192.168.1.121:554 - No route" - rather than `connected: false`. That
  field was added yesterday for precisely this.
- Head office briefly showed cameras 0/1 against a local 2/2. That is a
  60-second heartbeat, not a disagreement; the next one read 2/2.

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

403 lines
14 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
cfg.SessionToken, cfg.SessionRefresh, cfg.SessionEmail = "", "", ""
caPath, err := enrol.SaveCA(b.CACert, paths.BrokerCA())
if err != nil {
return err
}
cfg.BrokerCAFile = caPath
// 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())
creds := config.NewCreds(paths.APICredentials(), cfg.APIUser, cfg.APIPassword)
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, Creds: creds,
})
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())
creds := config.NewCreds(paths.APICredentials(), cfg.APIUser, cfg.APIPassword)
// 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...)
// 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 = cfg.EngineDir
if cmd.Dir == "" {
cmd.Dir = paths.InstallRoot()
}
// 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, Creds: creds,
})
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)
eng.Creds = creds
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
}