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"
|
"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
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
|
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(),
|
||||||
|
|||||||
@@ -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",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
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]
|
[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]
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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
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