// 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 run supervise the engine and report to head office (default) claim 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 \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 . 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 }