Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e262fc8482 | |||
| b59e667a68 | |||
| 70c447873d | |||
| 719ba2c7f5 |
80
agent/cmd/behavision-demo-pack/main.go
Normal file
80
agent/cmd/behavision-demo-pack/main.go
Normal 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)
|
||||
}
|
||||
@@ -32,7 +32,11 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
|
||||
"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/paths"
|
||||
)
|
||||
@@ -69,6 +73,17 @@ func run() error {
|
||||
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()
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -113,10 +128,19 @@ func run() error {
|
||||
// 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
|
||||
// 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)
|
||||
}
|
||||
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(" 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
|
||||
// answer. Any reply counts, including 401: the engine invents its own
|
||||
// credential when none is configured, and a refusal proves it is serving.
|
||||
func smokeTest(vpy string) error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second)
|
||||
func smokeTest(vpy string, demoCams []demo.Camera) error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
|
||||
defer cancel()
|
||||
|
||||
cmd := exec.CommandContext(ctx, vpy, "-m", "behavision", "run")
|
||||
@@ -370,7 +394,15 @@ func smokeTest(vpy string) error {
|
||||
if err == nil {
|
||||
_, _ = io.Copy(io.Discard, resp.Body)
|
||||
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() {
|
||||
break
|
||||
@@ -410,3 +442,102 @@ func pause() {
|
||||
fmt.Print(" Press Enter to close. ")
|
||||
_, _ = 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
114
agent/pkg/demo/bundle.go
Normal 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
|
||||
}
|
||||
87
agent/pkg/demo/bundle_test.go
Normal file
87
agent/pkg/demo/bundle_test.go
Normal 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")
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ models finish loading without a single unguarded None dereference.
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import time
|
||||
import logging
|
||||
import secrets
|
||||
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
|
||||
# removed camera leaves this generator running for the life of the
|
||||
# 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():
|
||||
jpeg = worker.latest_jpeg()
|
||||
if jpeg is not None:
|
||||
yield boundary + jpeg + b"\r\n"
|
||||
await asyncio.sleep(0.1) # ~10 fps to the browser
|
||||
now = time.time()
|
||||
if now - sent_at < min_gap:
|
||||
await asyncio.sleep(min_gap - (now - sent_at))
|
||||
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(
|
||||
generate(),
|
||||
|
||||
@@ -17,10 +17,20 @@ import numpy as np
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# Force TCP transport and a 5s socket timeout for RTSP before OpenCV loads
|
||||
# ffmpeg. UDP is the default and silently drops frames on lossy Wi-Fi.
|
||||
# Set before OpenCV loads ffmpeg, which reads this once.
|
||||
#
|
||||
# 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(
|
||||
"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",
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -162,7 +162,11 @@ class CameraWorker(threading.Thread):
|
||||
# every test using a stubbed worker passed.
|
||||
self._stopping = threading.Event()
|
||||
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._was_connected = False
|
||||
self.frames_processed = 0
|
||||
@@ -187,8 +191,46 @@ class CameraWorker(threading.Thread):
|
||||
self.source.stop()
|
||||
|
||||
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:
|
||||
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:
|
||||
return {
|
||||
@@ -236,7 +278,7 @@ class CameraWorker(threading.Thread):
|
||||
for track in ended:
|
||||
self._finish_track(track, ts)
|
||||
|
||||
self._publish_annotated(frame, active)
|
||||
self._remember_tracks(active)
|
||||
self.frames_processed += 1
|
||||
except Exception:
|
||||
log.exception("[%s] frame processing failed", self.cam_cfg.id)
|
||||
@@ -426,12 +468,14 @@ class CameraWorker(threading.Thread):
|
||||
track.quality, rcfg=self.rcfg):
|
||||
track.reinforcements += 1
|
||||
|
||||
def _publish_annotated(self, frame: np.ndarray, tracks: "list[Track]") -> None:
|
||||
canvas = frame.copy()
|
||||
def _remember_tracks(self, tracks: "list[Track]") -> None:
|
||||
"""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:
|
||||
if t.misses > 0:
|
||||
continue # only draw tracks matched in this frame
|
||||
x1, y1, x2, y2 = t.box
|
||||
if t.state == "resolved":
|
||||
color = _COLORS["known"] if t.label and not str(t.label).startswith(
|
||||
"Visitor") else _COLORS["new"]
|
||||
@@ -440,15 +484,10 @@ class CameraWorker(threading.Thread):
|
||||
color, text = _COLORS["ambiguous"], "?"
|
||||
else:
|
||||
color, text = _COLORS["pending"], ""
|
||||
cv2.rectangle(canvas, (x1, y1), (x2, y2), color, 2)
|
||||
if text:
|
||||
cv2.putText(canvas, text, (x1, max(20, y1 - 8)),
|
||||
cv2.FONT_HERSHEY_SIMPLEX, 0.55, color, 2)
|
||||
ok, buf = cv2.imencode(".jpg", canvas,
|
||||
[int(cv2.IMWRITE_JPEG_QUALITY), 80])
|
||||
if ok:
|
||||
with self._lock:
|
||||
self._annotated_jpeg = buf.tobytes()
|
||||
overlay.append((tuple(t.box), color, text))
|
||||
with self._lock:
|
||||
self._overlay = overlay
|
||||
self._overlay_ts = time.time()
|
||||
|
||||
|
||||
class Engine:
|
||||
|
||||
@@ -124,6 +124,23 @@ func (a *App) startup(ctx context.Context) {
|
||||
})
|
||||
|
||||
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
|
||||
|
||||
@@ -180,9 +180,16 @@ func (c *Client) send(ctx context.Context, method, path string, raw []byte, out
|
||||
switch {
|
||||
case resp.StatusCode == http.StatusUnauthorized && e.Error == "token_expired":
|
||||
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
|
||||
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
|
||||
if msg == "" {
|
||||
msg = fmt.Sprintf("%s %s: %s", method, path, resp.Status)
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
@@ -137,3 +138,43 @@ func TestVisitorIDIsPathEscaped(t *testing.T) {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -38,6 +38,13 @@ SETTING UP
|
||||
This takes several minutes. Leave the window open until it says Done.
|
||||
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
|
||||
|
||||
The window opens and an icon appears in the system tray, next to the
|
||||
|
||||
21
installer/run-with-lan-head-office.cmd
Normal file
21
installer/run-with-lan-head-office.cmd
Normal 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"
|
||||
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "behavision"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
description = "Production face recognition over RTSP"
|
||||
requires-python = ">=3.10"
|
||||
dependencies = [
|
||||
@@ -14,6 +14,10 @@ dependencies = [
|
||||
"python-dotenv>=1.0",
|
||||
"faiss-cpu>=1.7.4",
|
||||
"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]
|
||||
|
||||
@@ -8,3 +8,4 @@ PyYAML>=6.0
|
||||
python-dotenv>=1.0
|
||||
faiss-cpu>=1.7.4
|
||||
requests>=2.31
|
||||
pywin32>=306; sys_platform == "win32"
|
||||
|
||||
@@ -31,6 +31,7 @@ class FakeWorker:
|
||||
def stats(self): return {"camera_id": self.cam_cfg.id, "connected": True,
|
||||
"url": self.cam_cfg.safe_url()}
|
||||
def latest_jpeg(self): return None
|
||||
def latest_jpeg_since(self, known_ts): return None, known_ts
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
|
||||
144
tests/test_live_picture.py
Normal file
144
tests/test_live_picture.py
Normal 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"
|
||||
Reference in New Issue
Block a user