4 Commits

Author SHA1 Message Date
e262fc8482 The live picture was chained to the recognition pipeline
Reported from the first Windows install: the camera feed lags. It did,
and not because of the network, the proxy or the webview.

The MJPEG stream served _annotated_jpeg - the frame the pipeline had
most recently FINISHED with, encoded after detection, quality scoring,
tracking and identification had all run on it. On a modest shop PC that
is a few frames a second, and every picture was already as old as that
processing. It looked like lag because it was lag. On the fast machine
it was developed on the pipeline kept up with the stream's own 10 fps
cap, which is why nobody here ever saw it.

Two more things compounded it. Every processed frame was JPEG-encoded
whether or not a viewer existed - CPU spent on precisely the machine
short of it. And ffmpeg ran its RTSP demuxer with default buffering,
which holds a comfortable queue of frames before handing over the first:
half a second to two seconds a live view can never recover.

Now the picture and the boxes are decoupled. latest_jpeg_since takes the
capture thread's freshest frame at the camera's own rate and draws the
boxes from the last processed frame over it - encoded on demand, per
request, so a camera nobody watches costs no encode at all. The stream
sends a frame only when the camera has a newer one, capped at 15 fps;
nothing is sent twice. Boxes older than a second are not drawn, so a
stalled pipeline cannot leave one floating over an empty spot.
_publish_annotated becomes _remember_tracks: a handful of tuples under
the lock, no copy, no encode. ffmpeg gets nobuffer / low_delay /
max_delay.

Measured on cam2's sub-stream, same machine, ten seconds each:

  before   99 frames sent,  98 distinct    9.8 new pictures/s
  after   141 frames sent, 141 distinct   14.0 new pictures/s

against a 15 fps camera, with the pipeline still processing 166 of 181
captured frames alongside - and engine CPU DOWN from 90% with no viewer
to 62% with one attached.

Engine version 1.0.0 -> 1.1.0 so a re-run of setup reinstalls it rather
than pip deciding the requirement is already satisfied.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-11 16:28:05 +05:30
b59e667a68 A demo release with the office cameras sealed inside it
Wanted: install it and the two office cameras are already there - but
without the release carrying their admin password where anyone with the
zip can read it. "Encode it" does not achieve that; anything the
installer can decode, anyone holding the installer can decode.

pkg/demo seals the camera list with AES-256-GCM under a key that is NOT
in the package: a 120-bit unlock code minted when the bundle is sealed,
given to whoever runs setup by voice or message, typed once. The code
is random, so it is key material directly through SHA-256; a human-
chosen passphrase would need a KDF and a dependency, 120 random bits do
not. The sealed file contains the format marker and noise. Tested: the
password and the host do not appear in it, a wrong code and a flipped
byte are both refused as ErrWrongCode, every seal differs.

behavision-demo-pack seals; it runs on the build machine and is never
shipped. The code is printed once and stored nowhere.

behavision-setup, on finding demo-cameras.enc beside the engine source,
asks for the code BEFORE the ten-minute download so a mistyped one costs
seconds, and adds the cameras at the end - through the running engine's
own Add Camera endpoint, not by writing its file. The store's save() is
what applies DPAPI to the password on Windows, so this is how the
credential ends up encrypted and machine-bound on the demo PC rather
than in cameras.json for anyone who can read ProgramData. It then marks
the PC standalone, so the app opens on Live instead of asking for an
installation code it will never get.

Which found the gap that DPAPI only works if pywin32 is importable, and
nothing had ever pulled it in - every Windows install to date would have
logged the warning and written camera passwords in the clear. Added as
a Windows-only dependency.

Verified in a clean container: a wrong code refused, the right one
unlocks two cameras, every install step passes, both cameras added
through the API, standalone set.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-11 16:07:49 +05:30
70c447873d "Session expired" on a screen where nobody had signed in
The first Windows install reached the setup screen, typed an
installation code, and was told the session had expired. There was no
session. The code had been minted on a different head office, and the
server said so - 401 bad_token, "That installation code is not valid.
Ask for a new one." - and the client threw the message away, because it
mapped every 401 to the string "session expired".

A 401 on a call that carried a session is a session problem. A 401 on a
call that carried none is about the request, and the server's message is
the answer. The client now tells them apart by whether it sent a token.
Two tests, one for each side of the rule.

Also: a launcher for pointing a Windows PC at a head office on the LAN,
with the two settings that needs and a comment saying why neither is
acceptable outside a demo.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-11 15:31:55 +05:30
719ba2c7f5 Recognition starts with the app, not with a button
The engine only ever started when somebody pressed Start. So 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 existing: on a PC where setup has
not run yet, the supervisor would loop on a missing executable with
nothing useful to say. Start and Stop remain for the case where somebody
has deliberately stopped it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-11 12:49:56 +05:30
16 changed files with 745 additions and 28 deletions

View File

@@ -0,0 +1,80 @@
// Command behavision-demo-pack seals a camera list into demo-cameras.enc for a
// demo release. It runs on the machine that builds the release and is never
// shipped.
//
// behavision-demo-pack -cameras cameras.json -out demo-cameras.enc
//
// Prints the unlock code exactly once. It is not stored anywhere; a code you
// can look up later is a code anyone with access to the build machine holds.
// Lose it and seal again.
package main
import (
"encoding/json"
"flag"
"fmt"
"os"
"github.com/loyaly/behavision-agent/pkg/demo"
)
func main() {
in := flag.String("cameras", "", "JSON array of cameras (id, host, port, path, username, password)")
out := flag.String("out", "demo-cameras.enc", "sealed bundle to write")
flag.Parse()
if *in == "" {
fmt.Fprintln(os.Stderr, "usage: behavision-demo-pack -cameras cameras.json [-out demo-cameras.enc]")
os.Exit(2)
}
raw, err := os.ReadFile(*in)
if err != nil {
die("read cameras: %v", err)
}
var cams []demo.Camera
if err := json.Unmarshal(raw, &cams); err != nil {
die("cameras.json: %v", err)
}
if len(cams) == 0 {
die("no cameras in %s", *in)
}
for i, c := range cams {
switch {
case c.ID == "":
die("camera %d has no id", i)
case c.Host == "":
die("camera %q has no host", c.ID)
case c.Path == "":
die("camera %q has no path - the stream path is the field nobody can guess", c.ID)
}
}
// Re-marshal so only the fields the engine accepts travel, in a stable
// shape, whatever extra keys the input happened to carry.
plain, err := json.Marshal(cams)
if err != nil {
die("marshal: %v", err)
}
code, err := demo.NewCode()
if err != nil {
die("code: %v", err)
}
sealed, err := demo.Seal(code, plain)
if err != nil {
die("seal: %v", err)
}
if err := os.WriteFile(*out, sealed, 0o644); err != nil {
die("write: %v", err)
}
fmt.Printf("\n sealed %d camera(s) into %s (%d bytes)\n\n", len(cams), *out, len(sealed))
fmt.Printf(" unlock code: %s\n\n", code)
fmt.Println(" Shown once. Give it to whoever runs behavision-setup, by voice")
fmt.Println(" or message - not in the same place as the zip.")
fmt.Println()
}
func die(format string, args ...any) {
fmt.Fprintf(os.Stderr, " "+format+"\n", args...)
os.Exit(1)
}

View File

@@ -32,7 +32,11 @@ import (
"strings" "strings"
"time" "time"
"bytes"
"encoding/json"
"github.com/loyaly/behavision-agent/pkg/config" "github.com/loyaly/behavision-agent/pkg/config"
"github.com/loyaly/behavision-agent/pkg/demo"
"github.com/loyaly/behavision-agent/pkg/engine" "github.com/loyaly/behavision-agent/pkg/engine"
"github.com/loyaly/behavision-agent/pkg/paths" "github.com/loyaly/behavision-agent/pkg/paths"
) )
@@ -69,6 +73,17 @@ func run() error {
return fmt.Errorf("could not create %s: %w", state, err) return fmt.Errorf("could not create %s: %w", state, err)
} }
// A demo release ships its cameras sealed. Ask for the code NOW, before
// the ten-minute download, so a mistyped one costs seconds; the cameras
// are actually added at the end, through the running engine.
demoCams, err := unlockDemo(src)
if err != nil {
return err
}
if demoCams != nil {
step("Demo cameras", fmt.Sprintf("%d unlocked", len(demoCams)))
}
py, ver, err := findPython() py, ver, err := findPython()
if err != nil { if err != nil {
return err return err
@@ -113,10 +128,19 @@ func run() error {
// Proving it starts is the point. An installer that reports success and // Proving it starts is the point. An installer that reports success and
// leaves a shop with an engine that will not run has done worse than // leaves a shop with an engine that will not run has done worse than
// failing: the failure surfaces later, to someone who did not install it. // failing: the failure surfaces later, to someone who did not install it.
if err := smokeTest(vpy); err != nil { if err := smokeTest(vpy, demoCams); err != nil {
return fmt.Errorf("the engine installed but would not start: %w", err) return fmt.Errorf("the engine installed but would not start: %w", err)
} }
step("Engine starts and answers", "verified") step("Engine starts and answers", "verified")
if demoCams != nil {
step("Demo cameras", "added to the engine")
// No head office in a demo. Without this the app opens on "type an
// installation code" and sits there; with it, it opens on Live.
if err := markStandalone(); err != nil {
return err
}
step("Head office", "none - running on this PC only")
}
fmt.Println() fmt.Println()
fmt.Println(" Done. Start Behavision from the Start menu or the desktop icon.") fmt.Println(" Done. Start Behavision from the Start menu or the desktop icon.")
@@ -347,8 +371,8 @@ func writeConfig(vpy string) error {
// smokeTest starts the engine exactly as the app will and waits for its API to // smokeTest starts the engine exactly as the app will and waits for its API to
// answer. Any reply counts, including 401: the engine invents its own // answer. Any reply counts, including 401: the engine invents its own
// credential when none is configured, and a refusal proves it is serving. // credential when none is configured, and a refusal proves it is serving.
func smokeTest(vpy string) error { func smokeTest(vpy string, demoCams []demo.Camera) error {
ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second) ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
defer cancel() defer cancel()
cmd := exec.CommandContext(ctx, vpy, "-m", "behavision", "run") cmd := exec.CommandContext(ctx, vpy, "-m", "behavision", "run")
@@ -370,7 +394,15 @@ func smokeTest(vpy string) error {
if err == nil { if err == nil {
_, _ = io.Copy(io.Discard, resp.Body) _, _ = io.Copy(io.Discard, resp.Body)
resp.Body.Close() resp.Body.Close()
return nil if demoCams == nil {
return nil
}
// Through the engine's own Add Camera, not written to its file:
// the store is what applies DPAPI to the password on Windows, so
// this is how the credential ends up encrypted on disk rather
// than sitting in cameras.json for anyone who can read
// ProgramData.
return addCameras(demoCams)
} }
if cmd.ProcessState != nil && cmd.ProcessState.Exited() { if cmd.ProcessState != nil && cmd.ProcessState.Exited() {
break break
@@ -410,3 +442,102 @@ func pause() {
fmt.Print(" Press Enter to close. ") fmt.Print(" Press Enter to close. ")
_, _ = bufio.NewReader(os.Stdin).ReadString('\n') _, _ = bufio.NewReader(os.Stdin).ReadString('\n')
} }
// unlockDemo returns the sealed cameras a demo release ships, or nil when this
// is not a demo release. Asks for the unlock code on the console; three tries,
// because a code is read down a phone and typed by hand.
func unlockDemo(src string) ([]demo.Camera, error) {
sealed, err := os.ReadFile(filepath.Join(src, "demo-cameras.enc"))
if err != nil {
return nil, nil // not a demo release
}
fmt.Println()
fmt.Println(" This is a demo release with the cameras already set up.")
fmt.Println(" It needs the unlock code you were given.")
fmt.Println()
in := bufio.NewReader(os.Stdin)
for attempt := 1; attempt <= 3; attempt++ {
fmt.Print(" Unlock code: ")
line, _ := in.ReadString('\n')
plain, err := demo.Open(line, sealed)
if err == nil {
var cams []demo.Camera
if err := json.Unmarshal(plain, &cams); err != nil {
return nil, fmt.Errorf("the bundle unlocked but did not parse: %w", err)
}
fmt.Println()
return cams, nil
}
fmt.Printf(" %v\n", err)
}
return nil, errors.New("no valid unlock code after three tries. Check it " +
"with whoever gave you this release and run setup again")
}
// addCameras posts each demo camera to the running engine, with the credential
// the engine generated for itself on first start.
func addCameras(cams []demo.Camera) error {
user, pass, err := engineCredential()
if err != nil {
return err
}
client := &http.Client{Timeout: 30 * time.Second}
for _, c := range cams {
if c.Port == 0 {
c.Port = 554
}
body, _ := json.Marshal(c)
req, _ := http.NewRequest(http.MethodPost, "http://127.0.0.1:8010/api/cameras",
bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
if user != "" {
req.SetBasicAuth(user, pass)
}
resp, err := client.Do(req)
if err != nil {
return fmt.Errorf("adding camera %s: %w", c.ID, err)
}
msg, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
resp.Body.Close()
// 409 is "already there" - a re-run of setup, which is allowed.
if resp.StatusCode >= 300 && resp.StatusCode != http.StatusConflict {
return fmt.Errorf("adding camera %s: %s: %s", c.ID, resp.Status,
strings.TrimSpace(string(msg)))
}
}
return nil
}
// engineCredential reads the Basic credential the engine wrote on its first
// start. Empty when the engine is configured without one.
func engineCredential() (string, string, error) {
b, err := os.ReadFile(paths.APICredentials())
if err != nil {
if os.IsNotExist(err) {
return "", "", nil
}
return "", "", err
}
var user, pass string
for _, line := range strings.Split(string(b), "\n") {
if v, ok := strings.CutPrefix(line, "username="); ok {
user = strings.TrimSpace(v)
}
if v, ok := strings.CutPrefix(line, "password="); ok {
pass = strings.TrimSpace(v)
}
}
return user, pass, nil
}
// markStandalone records that this PC runs on its own, through the same
// config type the app reads.
func markStandalone() error {
path := paths.AgentConfig()
cfg, err := config.Load(path)
if err != nil {
return err
}
cfg.Standalone = true
return cfg.Save(path)
}

114
agent/pkg/demo/bundle.go Normal file
View File

@@ -0,0 +1,114 @@
// Package demo seals a camera list so a release can carry it without carrying
// the credentials in any usable form.
//
// The need: a demo build that installs with the office cameras already set up,
// handed to people who should not be able to read the cameras' admin password
// out of the zip. "Encode it" does not do that - anything the installer can
// decode, anyone holding the installer can decode. So the bundle is encrypted
// with a key that is NOT in the package: a short unlock code, generated when
// the bundle is sealed, spoken or messaged to whoever runs setup, and typed
// once. Without it the file is noise.
//
// The code is random, not chosen, so it is used as key material directly
// (through SHA-256) rather than stretched with a KDF. A human-chosen
// passphrase would need argon2 and a dependency; 120 random bits do not.
package demo
import (
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"crypto/sha256"
"encoding/base32"
"errors"
"fmt"
"strings"
)
// Magic identifies the file and the format version, so a future change can be
// told apart from corruption instead of failing as "authentication failed".
const magic = "BVDEMO1\n"
// Camera is one entry as the engine's Add Camera endpoint accepts it.
type Camera struct {
ID string `json:"id"`
Label string `json:"label,omitempty"`
Host string `json:"host"`
Port int `json:"port"`
Path string `json:"path"`
Username string `json:"username"`
Password string `json:"password"`
MaxWidth int `json:"max_width,omitempty"`
}
// NewCode mints an unlock code: 15 random bytes as 24 base32 characters in
// four groups, the same shape as an installation code, for the same reason -
// it gets read down a phone.
func NewCode() (string, error) {
raw := make([]byte, 15)
if _, err := rand.Read(raw); err != nil {
return "", err
}
s := base32.StdEncoding.WithPadding(base32.NoPadding).EncodeToString(raw)
return fmt.Sprintf("%s-%s-%s-%s", s[0:6], s[6:12], s[12:18], s[18:24]), nil
}
// NormalizeCode makes the typed and the printed form hash the same: case,
// spaces and dashes are all noise a person adds or drops.
func NormalizeCode(code string) string {
code = strings.ToUpper(code)
code = strings.NewReplacer("-", "", " ", "", "\t", "", "\r", "", "\n", "").Replace(code)
return code
}
func keyFor(code string) []byte {
sum := sha256.Sum256([]byte("behavision-demo-bundle:" + NormalizeCode(code)))
return sum[:]
}
// Seal encrypts plaintext under the code. Output is magic || nonce || ciphertext.
func Seal(code string, plaintext []byte) ([]byte, error) {
block, err := aes.NewCipher(keyFor(code))
if err != nil {
return nil, err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, err
}
nonce := make([]byte, gcm.NonceSize())
if _, err := rand.Read(nonce); err != nil {
return nil, err
}
out := append([]byte(magic), nonce...)
return gcm.Seal(out, nonce, plaintext, []byte(magic)), nil
}
// ErrWrongCode is what a mistyped code looks like. GCM cannot tell a wrong key
// from a corrupted file, and neither can we, so both read as this.
var ErrWrongCode = errors.New("that unlock code does not open this bundle")
// Open decrypts a sealed bundle.
func Open(code string, sealed []byte) ([]byte, error) {
if !strings.HasPrefix(string(sealed), magic) {
return nil, errors.New("not a Behavision demo bundle")
}
body := sealed[len(magic):]
block, err := aes.NewCipher(keyFor(code))
if err != nil {
return nil, err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, err
}
if len(body) < gcm.NonceSize() {
return nil, errors.New("bundle is truncated")
}
nonce, ct := body[:gcm.NonceSize()], body[gcm.NonceSize():]
plain, err := gcm.Open(nil, nonce, ct, []byte(magic))
if err != nil {
return nil, ErrWrongCode
}
return plain, nil
}

View File

@@ -0,0 +1,87 @@
package demo
import (
"bytes"
"errors"
"strings"
"testing"
)
func TestSealedBundleRoundTripsWithTheCodeAsTyped(t *testing.T) {
code, err := NewCode()
if err != nil {
t.Fatal(err)
}
if len(NormalizeCode(code)) != 24 {
t.Fatalf("code should be 24 base32 chars, got %q", code)
}
secret := []byte(`[{"id":"cam1","password":"the-camera-admin-password"}]`)
sealed, err := Seal(code, secret)
if err != nil {
t.Fatal(err)
}
// People type codes in lower case, with the dashes dropped, with a space
// where a dash was. All of those are the same code.
for _, typed := range []string{
code,
strings.ToLower(code),
strings.ReplaceAll(code, "-", ""),
strings.ReplaceAll(code, "-", " "),
" " + code + "\n",
} {
got, err := Open(typed, sealed)
if err != nil {
t.Fatalf("open with %q: %v", typed, err)
}
if !bytes.Equal(got, secret) {
t.Fatalf("round trip changed the contents")
}
}
}
// The whole point of the file: the password is not in it.
func TestTheSealedFileDoesNotContainTheSecret(t *testing.T) {
code, _ := NewCode()
sealed, _ := Seal(code, []byte(`{"password":"the-camera-admin-password","host":"192.168.1.121"}`))
for _, leak := range []string{"the-camera-admin-password", "192.168.1.121", "password"} {
if bytes.Contains(sealed, []byte(leak)) {
t.Fatalf("sealed bundle contains %q in the clear", leak)
}
}
}
func TestAWrongCodeIsRefusedNotMisread(t *testing.T) {
code, _ := NewCode()
other, _ := NewCode()
sealed, _ := Seal(code, []byte("secret"))
if _, err := Open(other, sealed); !errors.Is(err, ErrWrongCode) {
t.Fatalf("a different code should be ErrWrongCode, got %v", err)
}
// One flipped byte in the ciphertext is the same answer: GCM refuses
// rather than returning garbage that then gets written into cameras.json.
tampered := append([]byte{}, sealed...)
tampered[len(tampered)-1] ^= 0x01
if _, err := Open(code, tampered); !errors.Is(err, ErrWrongCode) {
t.Fatalf("a tampered bundle should be refused, got %v", err)
}
}
func TestSomethingThatIsNotABundleSaysSo(t *testing.T) {
if _, err := Open("ABCDEF-GHIJKL-MNOPQR-STUVWX", []byte("hello")); err == nil ||
errors.Is(err, ErrWrongCode) {
t.Fatalf("a non-bundle should be named as such, not blamed on the code: %v", err)
}
}
// Two seals of the same plaintext under the same code must differ: a fixed
// nonce would let two releases' bundles be compared byte for byte.
func TestEverySealIsDifferent(t *testing.T) {
code, _ := NewCode()
a, _ := Seal(code, []byte("same"))
b, _ := Seal(code, []byte("same"))
if bytes.Equal(a, b) {
t.Fatal("nonce is not random")
}
}

View File

@@ -6,6 +6,7 @@ models finish loading without a single unguarded None dereference.
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import time
import logging import logging
import secrets import secrets
from pathlib import Path from pathlib import Path
@@ -370,11 +371,23 @@ def create_app(engine: Engine) -> FastAPI:
# Stop when the camera is deleted or its worker dies - otherwise a # Stop when the camera is deleted or its worker dies - otherwise a
# removed camera leaves this generator running for the life of the # removed camera leaves this generator running for the life of the
# process, holding a reference to a worker nothing else can see. # process, holding a reference to a worker nothing else can see.
# Driven by the camera, not a timer: a frame goes out when the
# capture thread has one newer than the last one sent, so nothing
# is sent twice and nothing waits on the recognition pipeline.
# Capped at 15 fps - the office cameras' own rate - so a viewer
# never costs more encodes than the camera produces pictures.
last_ts, min_gap, sent_at = 0.0, 1.0 / 15, 0.0
while engine.workers.get(camera_id) is worker and worker.is_alive(): while engine.workers.get(camera_id) is worker and worker.is_alive():
jpeg = worker.latest_jpeg() now = time.time()
if jpeg is not None: if now - sent_at < min_gap:
yield boundary + jpeg + b"\r\n" await asyncio.sleep(min_gap - (now - sent_at))
await asyncio.sleep(0.1) # ~10 fps to the browser continue
jpeg, ts = worker.latest_jpeg_since(last_ts)
if jpeg is None:
await asyncio.sleep(0.02)
continue
last_ts, sent_at = ts, time.time()
yield boundary + jpeg + b"\r\n"
return StreamingResponse( return StreamingResponse(
generate(), generate(),

View File

@@ -17,10 +17,20 @@ import numpy as np
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
# Force TCP transport and a 5s socket timeout for RTSP before OpenCV loads # Set before OpenCV loads ffmpeg, which reads this once.
# ffmpeg. UDP is the default and silently drops frames on lossy Wi-Fi. #
# rtsp_transport=tcp: UDP is the default and silently drops frames on lossy
# Wi-Fi. stimeout: a 5s socket timeout so a dead camera is noticed.
#
# fflags=nobuffer and flags=low_delay: without them ffmpeg's RTSP demuxer
# holds a comfortable queue of frames before handing over the first, which
# on a live feed is half a second to two seconds of latency that no amount of
# work downstream can recover - the frame is already old when we get it. A
# recorder wants that buffer; a live view does not. max_delay caps the
# reorder wait for the same reason.
os.environ.setdefault( os.environ.setdefault(
"OPENCV_FFMPEG_CAPTURE_OPTIONS", "rtsp_transport;tcp|stimeout;5000000" "OPENCV_FFMPEG_CAPTURE_OPTIONS",
"rtsp_transport;tcp|stimeout;5000000|fflags;nobuffer|flags;low_delay|max_delay;200000",
) )

View File

@@ -162,7 +162,11 @@ class CameraWorker(threading.Thread):
# every test using a stubbed worker passed. # every test using a stubbed worker passed.
self._stopping = threading.Event() self._stopping = threading.Event()
self._lock = threading.Lock() self._lock = threading.Lock()
self._annotated_jpeg: Optional[bytes] = None # What the live view draws over the freshest frame: the boxes from
# the most recent processed frame, and when they were computed. NOT a
# pre-rendered JPEG - see latest_jpeg for why.
self._overlay: "list[tuple[tuple[int, int, int, int], tuple[int, int, int], str]]" = []
self._overlay_ts = 0.0
self._last_frame_ts = 0.0 self._last_frame_ts = 0.0
self._was_connected = False self._was_connected = False
self.frames_processed = 0 self.frames_processed = 0
@@ -187,8 +191,46 @@ class CameraWorker(threading.Thread):
self.source.stop() self.source.stop()
def latest_jpeg(self) -> Optional[bytes]: def latest_jpeg(self) -> Optional[bytes]:
jpeg, _ = self.latest_jpeg_since(0.0)
return jpeg
def latest_jpeg_since(self, known_ts: float) -> "tuple[Optional[bytes], float]":
"""The freshest captured frame with the latest boxes drawn on it, or
(None, known_ts) if the camera has produced nothing newer.
The live picture is deliberately NOT the frame the pipeline last
finished with. That version advanced only when detection, tracking and
identification had all completed on a frame - a few times a second on a
modest shop PC - and every picture it showed was already as old as that
processing. It looked like lag because it was lag. Here the picture runs
at the camera's rate off the capture thread's latest frame, and the
boxes - which genuinely can only update at pipeline rate - are drawn
over it from the last processed frame. Boxes may trail a fast walker by
one pipeline period; the picture never does.
Encoded on demand, per request, so a camera nobody is watching pays for
no JPEG at all. The old path encoded every processed frame whether or
not a viewer existed - CPU spent on precisely the machine short of it.
"""
frame, ts = self.source.latest_since(known_ts)
if frame is None:
return None, known_ts
with self._lock: with self._lock:
return self._annotated_jpeg overlay, overlay_ts = list(self._overlay), self._overlay_ts
# A stalled pipeline must not leave a box floating over an empty spot.
# Older than a second and the person has walked out from under it.
draw = overlay if (time.time() - overlay_ts) < 1.0 else []
if draw:
frame = frame.copy()
for (x1, y1, x2, y2), color, text in draw:
cv2.rectangle(frame, (x1, y1), (x2, y2), color, 2)
if text:
cv2.putText(frame, text, (x1, max(20, y1 - 8)),
cv2.FONT_HERSHEY_SIMPLEX, 0.55, color, 2)
ok, buf = cv2.imencode(".jpg", frame, [int(cv2.IMWRITE_JPEG_QUALITY), 80])
if not ok:
return None, known_ts
return buf.tobytes(), ts
def stats(self) -> dict: def stats(self) -> dict:
return { return {
@@ -236,7 +278,7 @@ class CameraWorker(threading.Thread):
for track in ended: for track in ended:
self._finish_track(track, ts) self._finish_track(track, ts)
self._publish_annotated(frame, active) self._remember_tracks(active)
self.frames_processed += 1 self.frames_processed += 1
except Exception: except Exception:
log.exception("[%s] frame processing failed", self.cam_cfg.id) log.exception("[%s] frame processing failed", self.cam_cfg.id)
@@ -426,12 +468,14 @@ class CameraWorker(threading.Thread):
track.quality, rcfg=self.rcfg): track.quality, rcfg=self.rcfg):
track.reinforcements += 1 track.reinforcements += 1
def _publish_annotated(self, frame: np.ndarray, tracks: "list[Track]") -> None: def _remember_tracks(self, tracks: "list[Track]") -> None:
canvas = frame.copy() """Record what to draw. Cheap: a handful of tuples under the lock,
no frame copy and no encode. The encode happens in latest_jpeg_since,
only when somebody is looking."""
overlay = []
for t in tracks: for t in tracks:
if t.misses > 0: if t.misses > 0:
continue # only draw tracks matched in this frame continue # only draw tracks matched in this frame
x1, y1, x2, y2 = t.box
if t.state == "resolved": if t.state == "resolved":
color = _COLORS["known"] if t.label and not str(t.label).startswith( color = _COLORS["known"] if t.label and not str(t.label).startswith(
"Visitor") else _COLORS["new"] "Visitor") else _COLORS["new"]
@@ -440,15 +484,10 @@ class CameraWorker(threading.Thread):
color, text = _COLORS["ambiguous"], "?" color, text = _COLORS["ambiguous"], "?"
else: else:
color, text = _COLORS["pending"], "" color, text = _COLORS["pending"], ""
cv2.rectangle(canvas, (x1, y1), (x2, y2), color, 2) overlay.append((tuple(t.box), color, text))
if text: with self._lock:
cv2.putText(canvas, text, (x1, max(20, y1 - 8)), self._overlay = overlay
cv2.FONT_HERSHEY_SIMPLEX, 0.55, color, 2) self._overlay_ts = time.time()
ok, buf = cv2.imencode(".jpg", canvas,
[int(cv2.IMWRITE_JPEG_QUALITY), 80])
if ok:
with self._lock:
self._annotated_jpeg = buf.tobytes()
class Engine: class Engine:

View File

@@ -124,6 +124,23 @@ func (a *App) startup(ctx context.Context) {
}) })
a.startPipeline(ctx) 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 // webhookURL is the loopback address the bridge is listening on, or empty

View File

@@ -180,9 +180,16 @@ func (c *Client) send(ctx context.Context, method, path string, raw []byte, out
switch { switch {
case resp.StatusCode == http.StatusUnauthorized && e.Error == "token_expired": case resp.StatusCode == http.StatusUnauthorized && e.Error == "token_expired":
return errTokenExpired return errTokenExpired
case resp.StatusCode == http.StatusUnauthorized: case resp.StatusCode == http.StatusUnauthorized && tok != "":
// A 401 on a call we sent a session with: the session is the problem.
return ErrUnauthorized return ErrUnauthorized
case resp.StatusCode >= 400: case resp.StatusCode >= 400:
// Every other 4xx/5xx - including a 401 on a call that carried NO
// session, such as redeeming an installation code - is about the
// request, and the server wrote its message for exactly this moment.
// Mapping those to "session expired" told an installer their session
// had lapsed on a screen where they had never signed in, and hid
// "That installation code is not valid" behind it.
msg := e.Message msg := e.Message
if msg == "" { if msg == "" {
msg = fmt.Sprintf("%s %s: %s", method, path, resp.Status) msg = fmt.Sprintf("%s %s: %s", method, path, resp.Status)

View File

@@ -6,6 +6,7 @@ import (
"errors" "errors"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"strings"
"testing" "testing"
) )
@@ -137,3 +138,43 @@ func TestVisitorIDIsPathEscaped(t *testing.T) {
t.Errorf("path = %q", got) t.Errorf("path = %q", got)
} }
} }
// Redeeming an installation code is the one call a fresh PC makes before it
// has any session. When the server refuses it - wrong code, wrong head office -
// it answers 401 with a message written for the installer. That message must
// reach them: "session expired" on a screen where nobody has signed in sent a
// real installer looking for a login problem that did not exist.
func TestARefusedInstallationCodeSaysWhyNotSessionExpired(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Header.Get("Authorization") != "" {
t.Errorf("enrol must not carry a session, got %q", r.Header.Get("Authorization"))
}
fail(w, http.StatusUnauthorized, "bad_token",
"That installation code is not valid. Ask for a new one.")
}))
t.Cleanup(srv.Close)
c := New(srv.URL) // deliberately no session
_, err := c.Bootstrap(context.Background(), "KWFH5S-EH46LT-EE4X47-OSOH7D")
if err == nil {
t.Fatal("a refused code must be an error")
}
if errors.Is(err, ErrUnauthorized) {
t.Fatalf("a refused code is not a session problem, got %v", err)
}
if !strings.Contains(err.Error(), "installation code is not valid") {
t.Fatalf("the server's own words should reach the installer, got %v", err)
}
}
// The other side of the same rule: a 401 on a call that DID carry a session is
// a session problem, and must still read as one.
func TestARejectedSessionStillReadsAsSessionExpired(t *testing.T) {
c := serve(t, func(w http.ResponseWriter, r *http.Request) {
fail(w, http.StatusUnauthorized, "unauthorized", "Sign in again.")
})
err := c.do(context.Background(), http.MethodGet, "/api/auth/me", nil, nil)
if !errors.Is(err, ErrUnauthorized) {
t.Fatalf("a 401 with a session should be ErrUnauthorized, got %v", err)
}
}

View File

@@ -38,6 +38,13 @@ SETTING UP
This takes several minutes. Leave the window open until it says Done. This takes several minutes. Leave the window open until it says Done.
If anything fails it prints why, and running it again is safe. If anything fails it prints why, and running it again is safe.
DEMO RELEASE ONLY: if the release came with the cameras already set up,
setup first asks for an unlock code. Type the code you were given. The
camera details are sealed inside the release and cannot be read without
it; with it, both cameras are added and the PC is set to run on its own,
with no head office. Skip the installation-code screen - it will not
appear.
3. Double-click Behavision.exe 3. Double-click Behavision.exe
The window opens and an icon appears in the system tray, next to the The window opens and an icon appears in the system tray, next to the

View File

@@ -0,0 +1,21 @@
@echo off
rem Start Behavision against a head office running on another PC on this LAN,
rem instead of the production server it uses by default.
rem
rem For demos and pilots only. Two things are deliberately weaker than
rem production and both are named here so nobody copies this into a shop:
rem
rem - head office over plain http, not https
rem - the message broker over plain tcp. The app REFUSES plaintext MQTT to
rem any address that is not its own machine, by design - the payloads are
rem customer visit records - so the second line below is the documented
rem escape hatch and must not be set anywhere that is not a demo.
rem
rem Edit the address to the PC running head office, then double-click this
rem instead of Behavision.exe. Everything else - the installation code, the
rem sign-in, the cameras - works exactly as INSTALL.txt describes.
set BEHAVISION_CLOUD=http://192.168.1.117:8088
set BEHAVISION_ALLOW_PLAINTEXT_MQTT=1
start "" "%~dp0Behavision.exe"

View File

@@ -1,6 +1,6 @@
[project] [project]
name = "behavision" name = "behavision"
version = "1.0.0" version = "1.1.0"
description = "Production face recognition over RTSP" description = "Production face recognition over RTSP"
requires-python = ">=3.10" requires-python = ">=3.10"
dependencies = [ dependencies = [
@@ -14,6 +14,10 @@ dependencies = [
"python-dotenv>=1.0", "python-dotenv>=1.0",
"faiss-cpu>=1.7.4", "faiss-cpu>=1.7.4",
"requests>=2.31", "requests>=2.31",
# DPAPI for camera passwords at rest (behavision/cameras.py). Without it the
# store logs a warning and writes them in the clear - which is what every
# Windows install had been doing, since nothing pulled this in.
"pywin32>=306; sys_platform == 'win32'",
] ]
[project.optional-dependencies] [project.optional-dependencies]

View File

@@ -8,3 +8,4 @@ PyYAML>=6.0
python-dotenv>=1.0 python-dotenv>=1.0
faiss-cpu>=1.7.4 faiss-cpu>=1.7.4
requests>=2.31 requests>=2.31
pywin32>=306; sys_platform == "win32"

View File

@@ -31,6 +31,7 @@ class FakeWorker:
def stats(self): return {"camera_id": self.cam_cfg.id, "connected": True, def stats(self): return {"camera_id": self.cam_cfg.id, "connected": True,
"url": self.cam_cfg.safe_url()} "url": self.cam_cfg.safe_url()}
def latest_jpeg(self): return None def latest_jpeg(self): return None
def latest_jpeg_since(self, known_ts): return None, known_ts
@pytest.fixture @pytest.fixture

144
tests/test_live_picture.py Normal file
View File

@@ -0,0 +1,144 @@
"""The live picture is the camera's latest frame, not the pipeline's.
Until this, the MJPEG stream served the frame the recognition pipeline had
most recently *finished* - so on a shop PC where detection plus identification
ran a few times a second, the live view ran a few times a second too, and every
picture it showed was already as old as that processing. It read as lag
because it was. These tests pin the decoupling: the picture comes from the
capture thread at its own rate, the boxes come from the pipeline at theirs,
and nothing is encoded for a camera nobody is watching.
"""
import time
import numpy as np
from behavision.config import CameraConfig, Config
from behavision.engine import CameraWorker
from behavision.tracking import Track
class StillSource:
"""A capture thread stand-in that hands out whatever frame it is given."""
def __init__(self):
self.frame = None
self.ts = 0.0
def set(self, frame):
self.frame, self.ts = frame, time.time()
def latest(self):
return self.frame, self.ts
def latest_since(self, known_ts):
if self.frame is None or self.ts <= known_ts:
return None, known_ts
return self.frame, self.ts
def stop(self):
pass
def make_worker(tmp_path):
cfg = Config()
cfg.app.data_dir = tmp_path
cam = CameraConfig(id="cam1", host="127.0.0.1", port=1, path="/none")
w = CameraWorker(cam, cfg, detector=None, encoder=None, gallery=None,
bus=None, attrs=None)
w.source = StillSource()
return w
def grey(v=90):
return np.full((120, 160, 3), v, dtype=np.uint8)
def resolved_track(box=(30, 30, 90, 100), label="Priya"):
t = Track(id=1, box=box, kps=np.zeros((5, 2), dtype=np.float32), score=0.9)
t.state, t.label, t.similarity = "resolved", label, 0.71
return t
def test_the_picture_advances_with_the_camera_not_the_pipeline(tmp_path):
w = make_worker(tmp_path)
# Nothing captured yet: nothing to show, and no encode happened.
assert w.latest_jpeg_since(0.0) == (None, 0.0)
# A frame arrives from the camera. The pipeline has not touched it - and
# the live view must not wait for it to.
w.source.set(grey(80))
jpeg1, ts1 = w.latest_jpeg_since(0.0)
assert jpeg1 is not None and ts1 > 0
# Same frame again: the stream asks "anything newer than ts1?" and the
# answer is no. This is what stops duplicates going down the wire.
assert w.latest_jpeg_since(ts1) == (None, ts1)
# The camera produces a new frame; the pipeline still has not run.
time.sleep(0.002)
w.source.set(grey(160))
jpeg2, ts2 = w.latest_jpeg_since(ts1)
assert jpeg2 is not None and ts2 > ts1 and jpeg2 != jpeg1
def test_boxes_from_the_last_processed_frame_are_drawn_on_the_fresh_one(tmp_path):
w = make_worker(tmp_path)
w.source.set(grey())
plain, _ = w.latest_jpeg_since(0.0)
# The pipeline finishes a frame with one recognised person in it.
w._remember_tracks([resolved_track()])
# The NEXT camera frame - which the pipeline has not seen - still carries
# the box, because a person does not vanish between two frames.
time.sleep(0.002)
w.source.set(grey())
boxed, _ = w.latest_jpeg_since(0.0)
assert boxed != plain, "a resolved track should be drawn on the live picture"
def test_a_stale_overlay_is_not_drawn(tmp_path):
"""A stalled pipeline must not leave a box floating over an empty spot."""
w = make_worker(tmp_path)
w.source.set(grey())
plain, _ = w.latest_jpeg_since(0.0)
w._remember_tracks([resolved_track()])
# Pretend the pipeline last ran a while ago.
with w._lock:
w._overlay_ts = time.time() - 2.0
time.sleep(0.002)
w.source.set(grey())
fresh, _ = w.latest_jpeg_since(0.0)
assert fresh == plain, "boxes older than a second should not be drawn"
def test_only_tracks_matched_in_the_frame_are_drawn(tmp_path):
"""A track being coasted on misses has no face under it right now."""
w = make_worker(tmp_path)
missed = resolved_track()
missed.misses = 3
w._remember_tracks([missed])
assert w._overlay == []
seen = resolved_track()
w._remember_tracks([seen])
assert len(w._overlay) == 1
(box, _color, text) = w._overlay[0]
assert box == (30, 30, 90, 100) and text.startswith("Priya (")
def test_remembering_tracks_does_not_encode_or_copy(tmp_path):
"""The whole CPU argument: recording what to draw is a few tuples, and the
frame is never touched. A shop PC with no viewer pays nothing."""
w = make_worker(tmp_path)
tracks = [resolved_track() for _ in range(5)]
t0 = time.perf_counter()
for _ in range(1000):
w._remember_tracks(tracks)
per_call_us = (time.perf_counter() - t0) / 1000 * 1e6
# A JPEG encode of even a small frame is hundreds of microseconds; this
# should be an order of magnitude under that.
assert per_call_us < 100, f"remembering tracks took {per_call_us:.0f}us"