Files
Behavision/demo/console.py
Suriyakumarvijayanayagam f97ffc913a demo console: polled frames and the real stats field names
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
2026-09-24 13:50:08 +05:30

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()