The camera used an MJPEG stream through the proxy and the engine's stats under names it does not use (frames/faces rather than frames_processed/faces_seen), so the picture was blank and the counter read zero. The engine re-serves its latest frame until the pipeline produces a new one, so a polled still is the same picture with none of the multipart fragility - which matters when the audience is in the room. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
331 lines
14 KiB
Python
331 lines
14 KiB
Python
"""A one-screen live demo of the whole Behavision chain, for showing someone.
|
|
|
|
Run it on the shop PC (here, this Mac) while the engine and agent are running.
|
|
It holds every credential itself and the browser holds none, so the page can be
|
|
put on a projector without putting a token on it.
|
|
|
|
What it shows, and why each part is there:
|
|
|
|
- the live camera, so the person walking past sees themselves;
|
|
- the CHAIN, measured rather than described: the engine recognised a face at
|
|
this instant, the same visit appeared in the cloud API this many seconds
|
|
later. That number is the product's claim, and it is computed here from two
|
|
independent sources rather than asserted;
|
|
- the customer, editable - type a name, walk past again, watch the name come
|
|
back through the cloud instead of "Visitor 5";
|
|
- the raw JSON a phone and a dashboard receive, side by side, because a
|
|
colleague's real question is "is this actually wired up or is it a mock".
|
|
|
|
.venv/bin/python demo/console.py # http://127.0.0.1:8099
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import base64
|
|
import json
|
|
import os
|
|
import threading
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
from pathlib import Path
|
|
|
|
ROOT = Path(__file__).resolve().parent.parent
|
|
STATE = Path(os.environ.get("BEHAVISION_DATA_DIR", ROOT / ".demo"))
|
|
CLOUD = os.environ.get("BEHAVISION_CLOUD", "https://mcp.loyaly.ai")
|
|
ENGINE = "http://127.0.0.1:8010"
|
|
PORT = int(os.environ.get("DEMO_PORT", "8099"))
|
|
|
|
# Whoever the demo signs in as. Staff on purpose: it is the weakest role that
|
|
# can do everything the shop floor does, so nothing here is only possible
|
|
# because we used an owner.
|
|
EMAIL = os.environ.get("DEMO_EMAIL", "staff.demo@tenext.in")
|
|
PASSWORD = os.environ.get("DEMO_PASSWORD", "admin@123")
|
|
|
|
|
|
def engine_auth() -> str:
|
|
"""The engine invents a Basic credential when none is configured, and
|
|
writes it here. Read it rather than keeping a second copy."""
|
|
f = STATE / "data" / "api_credentials.txt"
|
|
if not f.exists():
|
|
return ""
|
|
user = pw = ""
|
|
for line in f.read_text().splitlines():
|
|
# `key=value`, and `key: value` too - the engine writes one and people
|
|
# read the other, and which is which is not worth a support call.
|
|
if "=" in line or ":" in line:
|
|
k, v = line.split("=", 1) if "=" in line else line.split(":", 1)
|
|
if k.strip().lower() == "username":
|
|
user = v.strip()
|
|
elif k.strip().lower() == "password":
|
|
pw = v.strip()
|
|
if not user:
|
|
return ""
|
|
return "Basic " + base64.b64encode(f"{user}:{pw}".encode()).decode()
|
|
|
|
|
|
def fetch(url: str, *, headers=None, body=None, method="GET", timeout=20):
|
|
req = urllib.request.Request(url, method=method,
|
|
data=json.dumps(body).encode() if body is not None else None,
|
|
headers={k: v for k, v in (headers or {}).items() if v})
|
|
if body is not None:
|
|
req.add_header("content-type", "application/json")
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=timeout) as r:
|
|
raw = r.read()
|
|
return r.status, (json.loads(raw) if raw and r.headers.get("content-type", "").startswith("application/json") else raw)
|
|
except urllib.error.HTTPError as e:
|
|
raw = e.read()
|
|
try:
|
|
return e.code, json.loads(raw or b"{}")
|
|
except Exception:
|
|
return e.code, raw[:400]
|
|
except Exception as e:
|
|
return 0, {"error": str(e)}
|
|
|
|
|
|
class Cloud:
|
|
"""The signed-in session, refreshed when it expires."""
|
|
|
|
def __init__(self):
|
|
self.token = ""
|
|
self.lock = threading.Lock()
|
|
|
|
def sign_in(self) -> bool:
|
|
st, d = fetch(f"{CLOUD}/api/auth/login", method="POST",
|
|
body={"email": EMAIL, "password": PASSWORD, "device": "Demo console"})
|
|
if st == 200 and isinstance(d, dict):
|
|
self.token = d.get("access_token", "")
|
|
return True
|
|
return False
|
|
|
|
def call(self, path, method="GET", body=None, retry=True):
|
|
with self.lock:
|
|
if not self.token and not self.sign_in():
|
|
return 0, {"error": "cannot sign in to the platform"}
|
|
tok = self.token
|
|
st, d = fetch(f"{CLOUD}{path}", method=method, body=body,
|
|
headers={"authorization": f"Bearer {tok}"})
|
|
if st == 401 and retry:
|
|
with self.lock:
|
|
self.sign_in()
|
|
return self.call(path, method, body, retry=False)
|
|
return st, d
|
|
|
|
|
|
cloud = Cloud()
|
|
|
|
# The chain, as the watcher builds it. One dict per recognition, newest first.
|
|
events: list[dict] = []
|
|
events_lock = threading.Lock()
|
|
|
|
|
|
def watch():
|
|
"""Poll the engine's own event log and the cloud feed, and join them.
|
|
|
|
They are joined on the identity and the second, not on a shared id,
|
|
because the engine numbers identities locally and the server numbers them
|
|
per tenant - the two are deliberately different (see CLAUDE.md). What
|
|
matters for the demo is the LATENCY between one seeing a person and the
|
|
other, and that only needs the same person and the same moment.
|
|
"""
|
|
seen_local: set[str] = set()
|
|
while True:
|
|
try:
|
|
auth = engine_auth()
|
|
st, d = fetch(f"{ENGINE}/api/events?limit=25", headers={"authorization": auth})
|
|
# The engine returns a bare list; a dict with "events" is accepted
|
|
# too so this survives either shape.
|
|
evs = d if isinstance(d, list) else (d or {}).get("events", []) if isinstance(d, dict) else []
|
|
if st == 200:
|
|
for e in evs:
|
|
if e.get("type") not in ("person.new", "person.seen"):
|
|
continue
|
|
key = f"{e.get('ts')}|{e.get('camera_id')}|{(e.get('data') or {}).get('identity_id')}"
|
|
if key in seen_local:
|
|
continue
|
|
seen_local.add(key)
|
|
data = e.get("data") or {}
|
|
with events_lock:
|
|
events.insert(0, {
|
|
"key": key,
|
|
"at": time.time(),
|
|
"kind": "new" if e["type"] == "person.new" else "seen",
|
|
"engine": {
|
|
"label": data.get("label"),
|
|
"identity_id": data.get("identity_id"),
|
|
"similarity": data.get("similarity"),
|
|
"quality": data.get("quality"),
|
|
"gender": data.get("gender"),
|
|
"age": data.get("age"),
|
|
"camera": e.get("camera_id"),
|
|
"ts": e.get("ts"),
|
|
},
|
|
"cloud": None,
|
|
"latency": None,
|
|
})
|
|
del events[40:]
|
|
except Exception:
|
|
pass
|
|
|
|
# the other half: has the cloud got it yet?
|
|
try:
|
|
with events_lock:
|
|
pending = [e for e in events if e["cloud"] is None][:6]
|
|
if pending:
|
|
st, d = cloud.call("/api/visits?limit=12")
|
|
arrivals = (d or {}).get("arrivals", []) if isinstance(d, dict) else []
|
|
for e in pending:
|
|
for a in arrivals:
|
|
# same camera, and the cloud's visit is not older than
|
|
# the engine's sighting
|
|
if a.get("camera_id") != e["engine"]["camera"]:
|
|
continue
|
|
if a.get("visit_id") in [x["cloud"].get("visit_id") for x in events if x["cloud"]]:
|
|
continue
|
|
with events_lock:
|
|
e["cloud"] = {
|
|
"visit_id": a.get("visit_id"),
|
|
"visitor_id": a.get("visitor_id"),
|
|
"label": a.get("label"),
|
|
"ref": a.get("customer_ref") or a.get("ref"),
|
|
"is_new": a.get("is_new_visitor"),
|
|
"similarity": a.get("similarity"),
|
|
"site": a.get("site"),
|
|
"occurred_at": a.get("occurred_at"),
|
|
"image": a.get("image"),
|
|
}
|
|
e["latency"] = round(time.time() - e["at"], 1)
|
|
break
|
|
except Exception:
|
|
pass
|
|
time.sleep(1.0)
|
|
|
|
|
|
def snapshot() -> dict:
|
|
"""Everything the page draws, in one reply."""
|
|
auth = engine_auth()
|
|
_, health = fetch(f"{ENGINE}/api/health", headers={"authorization": auth})
|
|
_, stats = fetch(f"{ENGINE}/api/stats", headers={"authorization": auth})
|
|
st_v, visits = cloud.call("/api/visits?limit=3")
|
|
st_f, foot = cloud.call("/api/reports/footfall?from=%s&to=%s"
|
|
% (time.strftime("%Y-%m-%d", time.localtime(time.time() - 7 * 86400)),
|
|
time.strftime("%Y-%m-%d")))
|
|
with events_lock:
|
|
chain = json.loads(json.dumps(events[:8]))
|
|
cams = (stats or {}).get("cameras", []) if isinstance(stats, dict) else []
|
|
return {
|
|
"engine": {
|
|
"up": isinstance(health, dict) and bool(health.get("status")),
|
|
"model": (health or {}).get("recognition_model") if isinstance(health, dict) else None,
|
|
"cameras": [{"id": c.get("camera_id"), "connected": c.get("connected"),
|
|
"frames": c.get("frames_processed") or c.get("frames_total"),
|
|
"faces": c.get("faces_seen"),
|
|
"tracks": c.get("active_tracks")} for c in cams],
|
|
},
|
|
"cloud_ok": st_v == 200,
|
|
"chain": chain,
|
|
"mobile": {"request": "GET /api/visits?limit=3", "status": st_v, "body": visits},
|
|
"dashboard": {"request": "GET /api/reports/footfall?from=…&to=…", "status": st_f, "body": foot},
|
|
}
|
|
|
|
|
|
def main():
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
|
from urllib.parse import urlparse, parse_qs
|
|
|
|
page = (Path(__file__).parent / "console.html").read_bytes()
|
|
|
|
class H(BaseHTTPRequestHandler):
|
|
def log_message(self, *a): # quiet
|
|
pass
|
|
|
|
def _send(self, code, body, ctype="application/json"):
|
|
self.send_response(code)
|
|
self.send_header("content-type", ctype)
|
|
self.send_header("content-length", str(len(body)))
|
|
self.end_headers()
|
|
try:
|
|
self.wfile.write(body)
|
|
except (BrokenPipeError, ConnectionResetError):
|
|
pass
|
|
|
|
def do_GET(self):
|
|
u = urlparse(self.path)
|
|
if u.path == "/":
|
|
return self._send(200, page, "text/html; charset=utf-8")
|
|
if u.path == "/api/snapshot":
|
|
return self._send(200, json.dumps(snapshot()).encode())
|
|
if u.path == "/api/customer":
|
|
vid = parse_qs(u.query).get("id", [""])[0]
|
|
if not vid:
|
|
return self._send(400, b'{"error":"no id"}')
|
|
st, d = cloud.call(f"/api/visitors/{vid}/history?limit=8")
|
|
return self._send(200, json.dumps({"status": st, "history": d}).encode())
|
|
if u.path == "/camera.jpg":
|
|
# A polled still, not the MJPEG stream. The engine re-serves its
|
|
# latest frame until the pipeline produces a new one, so polling
|
|
# shows the same picture - and a multipart stream through a
|
|
# proxy is one more thing to fail in front of an audience.
|
|
cam = parse_qs(u.query).get("id", [""])[0]
|
|
st, body = fetch(f"{ENGINE}/api/cameras/{cam}/frame.jpg",
|
|
headers={"authorization": engine_auth()}, timeout=15)
|
|
if st != 200 or not isinstance(body, (bytes, bytearray)):
|
|
return self._send(502, b'{"error":"no frame"}')
|
|
self.send_response(200)
|
|
self.send_header("content-type", "image/jpeg")
|
|
self.send_header("cache-control", "no-store")
|
|
self.send_header("content-length", str(len(body)))
|
|
self.end_headers()
|
|
try:
|
|
self.wfile.write(body)
|
|
except (BrokenPipeError, ConnectionResetError):
|
|
pass
|
|
return
|
|
self._send(404, b'{"error":"no"}')
|
|
|
|
def do_POST(self):
|
|
u = urlparse(self.path)
|
|
n = int(self.headers.get("content-length", 0))
|
|
body = json.loads(self.rfile.read(n) or b"{}")
|
|
if u.path == "/api/profile":
|
|
vid = body.pop("id", "")
|
|
st, d = cloud.call(f"/api/visitors/{vid}/profile", method="PUT", body=body)
|
|
return self._send(200, json.dumps({"status": st, "body": d}).encode())
|
|
self._send(404, b'{"error":"no"}')
|
|
|
|
def _proxy_stream(self, url):
|
|
"""The engine's MJPEG, relayed so the browser needs no credential.
|
|
|
|
The engine's API is Basic-authenticated with a credential it
|
|
generated locally; putting that in a page would hand the whole
|
|
biometric API to anyone who opened it.
|
|
"""
|
|
try:
|
|
req = urllib.request.Request(url, headers={"authorization": engine_auth()})
|
|
up = urllib.request.urlopen(req, timeout=20)
|
|
except Exception:
|
|
return self._send(502, b'{"error":"camera not available"}')
|
|
self.send_response(200)
|
|
self.send_header("content-type", up.headers.get("content-type", "multipart/x-mixed-replace"))
|
|
self.end_headers()
|
|
try:
|
|
while True:
|
|
chunk = up.read(8192)
|
|
if not chunk:
|
|
break
|
|
self.wfile.write(chunk)
|
|
self.wfile.flush()
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
up.close()
|
|
|
|
threading.Thread(target=watch, daemon=True).start()
|
|
print(f"\n Demo console → http://127.0.0.1:{PORT}\n")
|
|
print(f" engine {ENGINE} · cloud {CLOUD} · signed in as {EMAIL}\n")
|
|
ThreadingHTTPServer(("127.0.0.1", PORT), H).serve_forever()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|