The agent read the engine's generated credential file once, at startup. On a brand new install that file does not exist yet: the agent starts the engine, and the engine writes its credential seconds later. So the agent held an empty credential for the life of the process and every call it makes - health, stats, camera sync, the embedding for a visit - came back 401, with a tray showing a red engine that was running perfectly. Measured on a fresh state directory today: three 401s, no camera ever reconciled, and the engine left running the YAML-seeded main stream instead of the sub-stream head office holds. The install script hid this on Windows because setup runs the engine once before the app starts. config.Creds resolves lazily and re-reads on a rejection; the camera client, the supervisor and the desktop app's engine client all retry once when it changes. A configured BEHAVISION_API_USER is never re-read - an operator who set one means it. Tests pin the actual first-run ordering. Also adds demo/, a one-screen live console for showing the whole chain: camera, the six steps with a measured camera-to-cloud latency, the customer editable in place, and the raw JSON a phone and a dashboard receive from production side by side. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
309 lines
13 KiB
Python
309 lines
13 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})
|
|
if st == 200 and isinstance(d, dict):
|
|
for e in d.get("events", []):
|
|
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"), "faces": c.get("faces")} 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.mjpeg":
|
|
cam = parse_qs(u.query).get("id", [""])[0]
|
|
return self._proxy_stream(f"{ENGINE}/api/cameras/{cam}/stream.mjpeg")
|
|
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()
|