// Package engine starts, watches and stops the Python recognition engine. // // Go cannot run ONNX, OpenCV or FAISS, so the engine stays Python and ships // frozen. What Go owns is its lifecycle: start it, keep it up, capture its // output, and stop it when the user asks — which is what the tray's start/stop // buttons actually drive. // // Deliberately not a Windows service. A service runs in session 0 and cannot // draw a tray icon, and spawning a child process needs no elevation while // controlling a service does. A service wrapper can be layered on later // without touching anything here. package engine import ( "bufio" "context" "encoding/json" "errors" "fmt" "io" "net/http" "os" "os/exec" "sync" "time" ) // State is what the tray icon colours itself from. type State string const ( Stopped State = "stopped" // not running, and not meant to be Starting State = "starting" // process spawned, not yet answering Running State = "running" // answering /api/health Backoff State = "backoff" // crashed, waiting to retry Failed State = "failed" // gave up ) const ( minBackoff = 1 * time.Second maxBackoff = 30 * time.Second // A run that lasted this long counts as healthy, so the next crash starts // its backoff from the bottom again. Without this a process that runs fine // for hours and then crashes once waits the full 30s to come back. stableRun = 60 * time.Second // How long a stopping process gets to exit on its own before it is killed. stopGrace = 10 * time.Second ) // Options configures a Supervisor. type Options struct { // Command builds the process to run. Injected rather than hardcoded so // tests can supervise /bin/sh instead of a 200 MB frozen engine. Command func(ctx context.Context) *exec.Cmd // LogWriter receives the engine's stdout and stderr. A crashed engine with // no captured output is undiagnosable, which on a customer site means a // site visit. LogWriter io.Writer // HealthURL, StatsURL, User, Password address the engine's own API. HealthURL string StatsURL string User string Password string // MaxRestarts of 0 means never give up. Non-zero is for tests. MaxRestarts int now func() time.Time } // Supervisor keeps one engine process running. Safe for concurrent use. type Supervisor struct { opts Options mu sync.Mutex state State lastErr error restarts int cancel context.CancelFunc done chan struct{} } func New(opts Options) *Supervisor { if opts.LogWriter == nil { opts.LogWriter = io.Discard } if opts.now == nil { opts.now = time.Now } return &Supervisor{opts: opts, state: Stopped} } // Start launches the engine and keeps it running until Stop. Calling it while // already running is a no-op rather than a second process — two engines on one // SQLite WAL and one camera is exactly the failure this package exists to // avoid. func (s *Supervisor) Start() { s.mu.Lock() if s.cancel != nil { s.mu.Unlock() return } ctx, cancel := context.WithCancel(context.Background()) s.cancel = cancel s.done = make(chan struct{}) s.state = Starting s.restarts = 0 done := s.done s.mu.Unlock() go s.supervise(ctx, done) } // Stop asks the engine to exit and waits for it. func (s *Supervisor) Stop() { s.mu.Lock() cancel, done := s.cancel, s.done s.cancel = nil s.mu.Unlock() if cancel == nil { return } cancel() if done != nil { <-done } s.setState(Stopped, nil) } // State reports what the supervisor is doing, plus the last error if any. func (s *Supervisor) State() (State, error) { s.mu.Lock() defer s.mu.Unlock() return s.state, s.lastErr } // Restarts counts crash-restarts since Start. func (s *Supervisor) Restarts() int { s.mu.Lock() defer s.mu.Unlock() return s.restarts } // -- the loop -------------------------------------------------------------- func (s *Supervisor) supervise(ctx context.Context, done chan struct{}) { defer close(done) backoff := minBackoff for { if ctx.Err() != nil { return } s.setState(Starting, nil) started := s.opts.now() err := s.runOnce(ctx) ran := s.opts.now().Sub(started) // A cancelled context means the user pressed Stop. Exiting then is // success, not a crash, and restarting would be the single most // annoying bug a tray app can have. if ctx.Err() != nil { return } s.mu.Lock() s.restarts++ restarts := s.restarts s.mu.Unlock() if s.opts.MaxRestarts > 0 && restarts >= s.opts.MaxRestarts { s.setState(Failed, err) return } if ran >= stableRun { backoff = minBackoff } s.setState(Backoff, err) select { case <-ctx.Done(): return case <-time.After(backoff): } if backoff < maxBackoff { backoff *= 2 if backoff > maxBackoff { backoff = maxBackoff } } } } func (s *Supervisor) runOnce(ctx context.Context) error { cmd := s.opts.Command(ctx) stdout, err := cmd.StdoutPipe() if err != nil { return err } cmd.Stderr = cmd.Stdout // Cancel ends the whole process tree, not just the process exec spawned. // `kill` is filled in after Start, once the tree is confined; until then // it is exec's own behaviour. var kill func() error cmd.Cancel = func() error { if kill == nil { return cmd.Process.Kill() } return kill() } prepare(cmd) if err := cmd.Start(); err != nil { return fmt.Errorf("engine failed to start: %w", err) } k, release, err := confine(cmd) if err != nil { // Not fatal: the engine runs, and stopping it falls back to killing // the one process. Logged because on Windows that fallback is the // bug this exists to fix. fmt.Fprintf(s.opts.LogWriter, "supervisor: could not confine engine process tree: %v\n", err) } kill = k defer release() pumped := make(chan struct{}) go func() { defer close(pumped) sc := bufio.NewScanner(stdout) sc.Buffer(make([]byte, 0, 64*1024), 1024*1024) for sc.Scan() { fmt.Fprintln(s.opts.LogWriter, sc.Text()) } }() s.setState(Running, nil) waitErr := cmd.Wait() <-pumped // A context cancel terminates the child through exec's own handling; the // resulting error is expected, not a fault. if ctx.Err() != nil { return nil } if waitErr != nil { return fmt.Errorf("engine exited: %w", waitErr) } return errors.New("engine exited unexpectedly with status 0") } func (s *Supervisor) setState(st State, err error) { s.mu.Lock() s.state = st if err != nil { s.lastErr = err } s.mu.Unlock() } // -- health ---------------------------------------------------------------- // Health is the subset of /api/health the tray and the server care about. type Health struct { Status string `json:"status"` RecognitionModel string `json:"recognition_model"` Cameras map[string]bool `json:"cameras"` // Paths is a STRUCT, not map[string]string, because `frozen` is a bool. // It was a map of strings, so decoding the engine's real reply failed with // "cannot unmarshal bool into Go struct field Health.paths" - and because // one bad field fails the whole document, a perfectly healthy engine was // reported unreachable: red tray, and a heartbeat carrying neither the // model nor the camera list, so head office showed 0 of 0 cameras for a // site that was watching one. Paths EnginePaths `json:"paths"` } // EnginePaths mirrors what `behavision paths` and /api/health report. Unknown // fields are ignored by encoding/json, so the engine can add to it freely. type EnginePaths struct { Frozen bool `json:"frozen"` InstallRoot string `json:"install_root"` StateRoot string `json:"state_root"` Config string `json:"config"` DataDir string `json:"data_dir"` ModelsDir string `json:"models_dir"` } // Health polls the engine's own API. A running process is not the same as a // working engine: the model can fail to load and the process stays up. func (s *Supervisor) Health(ctx context.Context) (*Health, error) { if s.opts.HealthURL == "" { return nil, errors.New("no health url configured") } var h Health if err := s.getJSON(ctx, s.opts.HealthURL, &h); err != nil { return nil, err } return &h, nil } // Stats is the slice of /api/stats the heartbeat carries. // // Only fraction_below_gate, because that is the number that decides whether a // site's footfall can be believed at all - the share of faces its cameras saw // and discarded before they ever became a visit. Everything else in /api/stats // is a local diagnostic and belongs on the local dashboard, not on the wire // every thirty seconds. type Stats struct { Cameras []struct { CameraID string `json:"camera_id"` Pipeline struct { BestQuality struct { N int `json:"n"` FractionBelowGate float64 `json:"fraction_below_gate"` } `json:"best_quality"` } `json:"pipeline"` } `json:"cameras"` } // WorstBelowGate returns the worst camera's figure, and whether any camera has // measured enough faces to have an opinion. // // Worst rather than average: one badly placed camera is a hole in the report, // and averaging it against three good ones hides the only camera anyone needs // to move. The sample floor is there because three faces is an anecdote - // reporting 1.00 from a single below-gate track would raise an alarm about a // camera nobody has walked past yet. func (s *Stats) WorstBelowGate() (float64, bool) { const minSamples = 10 worst, found := 0.0, false for _, c := range s.Cameras { if c.Pipeline.BestQuality.N < minSamples { continue } if !found || c.Pipeline.BestQuality.FractionBelowGate > worst { worst, found = c.Pipeline.BestQuality.FractionBelowGate, true } } return worst, found } // Stats polls the engine's pipeline counters. func (s *Supervisor) Stats(ctx context.Context) (*Stats, error) { if s.opts.StatsURL == "" { return nil, errors.New("no stats url configured") } var out Stats if err := s.getJSON(ctx, s.opts.StatsURL, &out); err != nil { return nil, err } return &out, nil } func (s *Supervisor) getJSON(ctx context.Context, url string, out any) error { req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) if err != nil { return err } if s.opts.User != "" { req.SetBasicAuth(s.opts.User, s.opts.Password) } resp, err := (&http.Client{Timeout: 5 * time.Second}).Do(req) if err != nil { return err } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { return fmt.Errorf("%s returned %s", url, resp.Status) } return json.NewDecoder(io.LimitReader(resp.Body, 4<<20)).Decode(out) } // LogFile opens the engine log, rotating aside anything already there so one // run's output cannot be mistaken for another's. func LogFile(path string) (*os.File, error) { if _, err := os.Stat(path); err == nil { os.Rename(path, path+".1") } return os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600) }