The add-camera form asked for an IP address, and a shop owner does not know their camera's IP address - it is on a sticker under the camera or in a menu that differs by make. That field is where onboarding stopped for anyone who was not an installer. behavision/discover.py: one ONVIF WS-Discovery multicast (names the camera and often its make) merged with a TCP sweep of port 554 across the local /24 (misses nothing that streams). Stdlib only, ~4 s on the office network, both cameras found. The add-camera sheet leads with 'Find cameras on this network'; picking a row fills the address and, when the make is recognisable, the stream path. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
417 lines
17 KiB
Python
417 lines
17 KiB
Python
"""HTTP API + minimal live dashboard (FastAPI).
|
|
|
|
Every endpoint is guarded by engine readiness; the server can start before
|
|
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
|
|
from typing import Optional
|
|
|
|
from fastapi import Depends, FastAPI, HTTPException
|
|
from fastapi.responses import HTMLResponse, Response, StreamingResponse
|
|
from fastapi.security import HTTPBasic, HTTPBasicCredentials
|
|
from pydantic import BaseModel, ValidationError
|
|
|
|
from .config import ApiSection, CameraConfig, CameraTuning
|
|
from .commission import CommissionRun
|
|
from .events import Event
|
|
from .engine import Engine
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
_STATIC = Path(__file__).parent / "static"
|
|
|
|
|
|
class RenamePayload(BaseModel):
|
|
label: str
|
|
|
|
|
|
class CommissionPayload(BaseModel):
|
|
seconds: float = 25.0
|
|
|
|
|
|
class MergePayload(BaseModel):
|
|
"""`into` is the identity that survives. `force` overrides the
|
|
similarity guard and is never the default: a wrong merge cannot be
|
|
undone, because nothing records which embedding came from whom."""
|
|
into: int
|
|
force: bool = False
|
|
|
|
|
|
class CameraPayload(BaseModel):
|
|
"""Camera as the UI submits it. Mirrors CameraConfig but every field is
|
|
optional so PATCH can send a subset.
|
|
|
|
Optional[...] rather than `X | None`: pydantic evaluates field annotations
|
|
at runtime, and CameraConfig already uses this form.
|
|
"""
|
|
id: Optional[str] = None
|
|
url: Optional[str] = None
|
|
host: Optional[str] = None
|
|
port: Optional[int] = None
|
|
path: Optional[str] = None
|
|
username: Optional[str] = None
|
|
password: Optional[str] = None
|
|
webcam: Optional[int] = None
|
|
max_width: Optional[int] = None
|
|
# Per-camera gate overrides. Without this the store could hold them but
|
|
# nothing could set them, so the commissioning advice ("loosen this
|
|
# camera's quality gate") had no way to be acted on.
|
|
tuning: Optional[CameraTuning] = None
|
|
|
|
|
|
def camera_public(cam: CameraConfig, worker=None) -> dict:
|
|
"""Camera as the API returns it.
|
|
|
|
The password is NEVER included — not masked, not empty-string-if-set,
|
|
absent. `safe_url()` already exists for exactly this and masks credentials
|
|
inside the URL form too.
|
|
"""
|
|
out = {
|
|
"id": cam.id, "host": cam.host, "port": cam.port, "path": cam.path,
|
|
"username": cam.username, "webcam": cam.webcam,
|
|
"max_width": cam.max_width, "has_password": bool(cam.password),
|
|
"url": cam.safe_url(),
|
|
"tuning": cam.tuning.model_dump(),
|
|
}
|
|
if worker is not None:
|
|
out.update(worker.stats())
|
|
out["url"] = cam.safe_url() # worker.stats() also carries a url key
|
|
return out
|
|
|
|
|
|
def _auth_dependencies(api_cfg: ApiSection) -> list:
|
|
"""HTTP Basic over every route when credentials are configured.
|
|
|
|
Applied at app level rather than per-route so a future endpoint cannot be
|
|
added unprotected by omission. Basic (not a token) because the dashboard
|
|
is a browser page: the browser prompts once and then attaches the header
|
|
to the MJPEG <img> subresource too, which a bearer token cannot do.
|
|
"""
|
|
if not api_cfg.auth_enabled:
|
|
return []
|
|
scheme = HTTPBasic()
|
|
|
|
def check(credentials: HTTPBasicCredentials = Depends(scheme)) -> None:
|
|
# compare_digest on both halves: no early exit, no timing signal.
|
|
ok_user = secrets.compare_digest(
|
|
credentials.username.encode("utf-8"),
|
|
api_cfg.username.encode("utf-8"))
|
|
ok_pass = secrets.compare_digest(
|
|
credentials.password.encode("utf-8"),
|
|
api_cfg.password.encode("utf-8"))
|
|
if not (ok_user and ok_pass):
|
|
raise HTTPException(401, "invalid credentials",
|
|
headers={"WWW-Authenticate": "Basic"})
|
|
|
|
return [Depends(check)]
|
|
|
|
|
|
def _reencode(jpeg: bytes, width: int, quality: int) -> "bytes | None":
|
|
"""Decode, scale and re-encode one frame. None on any failure.
|
|
|
|
None rather than an exception on purpose: the caller falls back to the
|
|
original frame, so a re-encode that fails costs bandwidth rather than the
|
|
picture. A live view that goes blank because a resize failed is a worse
|
|
outcome than one that is briefly larger than asked for.
|
|
"""
|
|
try:
|
|
import cv2
|
|
import numpy as np
|
|
|
|
img = cv2.imdecode(np.frombuffer(jpeg, np.uint8), cv2.IMREAD_COLOR)
|
|
if img is None:
|
|
return None
|
|
if 0 < width < img.shape[1]:
|
|
# Only ever DOWN. Upscaling a frame to a requested width would send
|
|
# more bytes than the original for no more detail.
|
|
scale = width / img.shape[1]
|
|
img = cv2.resize(img, (width, max(1, int(img.shape[0] * scale))),
|
|
interpolation=cv2.INTER_AREA)
|
|
q = quality if 1 <= quality <= 100 else 75
|
|
ok, buf = cv2.imencode(".jpg", img, [int(cv2.IMWRITE_JPEG_QUALITY), q])
|
|
return buf.tobytes() if ok else None
|
|
except Exception:
|
|
log.debug("frame re-encode failed", exc_info=True)
|
|
return None
|
|
|
|
|
|
def create_app(engine: Engine) -> FastAPI:
|
|
app = FastAPI(title="Behavision", version="1.0.0",
|
|
dependencies=_auth_dependencies(engine.cfg.api))
|
|
|
|
def worker_or_404(camera_id: str):
|
|
worker = engine.workers.get(camera_id)
|
|
if worker is None:
|
|
raise HTTPException(404, f"unknown camera '{camera_id}'")
|
|
return worker
|
|
|
|
@app.get("/", response_class=HTMLResponse)
|
|
def dashboard() -> str:
|
|
return (_STATIC / "dashboard.html").read_text(encoding="utf-8")
|
|
|
|
@app.get("/static/favicon.png")
|
|
def favicon() -> Response:
|
|
# The one static asset besides the page itself. Served explicitly
|
|
# rather than mounting the directory: nothing else in there is meant
|
|
# to be reachable, and a mount would make that a matter of what lands
|
|
# in the folder.
|
|
return Response((_STATIC / "favicon.png").read_bytes(), media_type="image/png",
|
|
headers={"cache-control": "public, max-age=86400"})
|
|
|
|
@app.get("/api/health")
|
|
def health() -> dict:
|
|
from .paths import describe
|
|
return {"status": "ok" if engine.started_at else "starting",
|
|
"recognition_model": engine.encoder.model_name,
|
|
# "where is my database" must be answerable from the API: the
|
|
# tray, the installer and support all need it, and installed
|
|
# it is not next to the code.
|
|
"paths": {**describe(), "data_dir": str(engine.cfg.app.data_dir),
|
|
"models_dir": str(engine.cfg.app.models_dir)},
|
|
"cameras": {cid: w.source.connected
|
|
for cid, w in engine.workers.items()}}
|
|
|
|
@app.get("/api/stats")
|
|
def stats() -> dict:
|
|
return engine.stats()
|
|
|
|
@app.get("/api/events")
|
|
def events(limit: int = 50) -> list:
|
|
return list(engine.bus.recent)[:limit]
|
|
|
|
@app.get("/api/identities")
|
|
def identities(limit: int = 200) -> list:
|
|
return engine.store.list_identities(limit)
|
|
|
|
@app.get("/api/sightings")
|
|
def sightings(limit: int = 100) -> list:
|
|
return engine.store.recent_sightings(limit)
|
|
|
|
@app.patch("/api/identities/{identity_id}")
|
|
def rename_identity(identity_id: int, payload: RenamePayload) -> dict:
|
|
if not engine.store.rename_identity(identity_id, payload.label.strip()):
|
|
raise HTTPException(404, "identity not found")
|
|
return engine.store.get_identity(identity_id)
|
|
|
|
@app.get("/api/identities/{identity_id}/embedding")
|
|
def identity_embedding(identity_id: int) -> dict:
|
|
"""One identity's best stored vector, for forwarding to the server.
|
|
|
|
This returns biometric personal data. It is on the authenticated local
|
|
API and bound to loopback in the product, and it exists because the
|
|
event bus deliberately does not carry embeddings — putting a 512-float
|
|
template on the bus would send it to the log sink and the email sink
|
|
too.
|
|
"""
|
|
best = engine.store.best_embedding(identity_id, engine.gallery.model_name)
|
|
if best is None:
|
|
raise HTTPException(404, "no embedding for this identity "
|
|
"(or it was made by a different model)")
|
|
vector, quality = best
|
|
return {"identity_id": identity_id,
|
|
"model": engine.gallery.model_name,
|
|
"quality": round(quality, 3),
|
|
"embedding": [round(float(x), 6) for x in vector]}
|
|
|
|
@app.get("/api/identities/duplicates")
|
|
def duplicate_identities(limit: int = 20) -> list:
|
|
"""Identity pairs that look like one person enrolled twice."""
|
|
return engine.gallery.duplicate_candidates(limit)
|
|
|
|
@app.post("/api/identities/{identity_id}/merge")
|
|
def merge_identity(identity_id: int, payload: MergePayload) -> dict:
|
|
result = engine.gallery.merge_identities(
|
|
identity_id, payload.into, force=payload.force)
|
|
if not result.get("ok"):
|
|
reason = result.get("reason", "merge refused")
|
|
# 409, not 400: the request is well formed, it conflicts with what
|
|
# the gallery believes. The body carries the measured similarity so
|
|
# the UI can show the operator what it is asking them to override.
|
|
status = 404 if "not found" in reason else 409
|
|
raise HTTPException(status, detail=result)
|
|
engine.bus.publish(Event(
|
|
type="identity.merged", camera_id="",
|
|
data={k: result[k] for k in
|
|
("source", "target", "label", "similarity", "forced",
|
|
"embeddings_moved", "sightings_moved")}))
|
|
return result
|
|
|
|
@app.delete("/api/identities/{identity_id}")
|
|
def delete_identity(identity_id: int) -> dict:
|
|
if not engine.gallery.delete_identity(identity_id):
|
|
raise HTTPException(404, "identity not found")
|
|
return {"deleted": identity_id}
|
|
|
|
@app.get("/api/cameras")
|
|
def cameras() -> list:
|
|
out = []
|
|
for cam in engine.camera_store.list():
|
|
out.append(camera_public(cam, engine.workers.get(cam.id)))
|
|
return out
|
|
|
|
@app.post("/api/cameras", status_code=201)
|
|
def add_camera(payload: CameraPayload) -> dict:
|
|
data = payload.model_dump(exclude_none=True)
|
|
if not data.get("id"):
|
|
raise HTTPException(400, "id is required")
|
|
try:
|
|
cam = CameraConfig.model_validate(data)
|
|
cam.source() # reject "no url, no host, no webcam" before storing
|
|
# Resolve the per-camera gates here too. Without this an inverted
|
|
# enroll/match pair was only caught when the worker was built,
|
|
# which surfaced as a 500 "stored but failed to start" instead of
|
|
# telling the user what was wrong with what they typed.
|
|
engine.cfg.recognition.merged(cam.tuning)
|
|
except (ValidationError, ValueError) as exc:
|
|
raise HTTPException(400, str(exc))
|
|
try:
|
|
engine.camera_store.add(cam)
|
|
except ValueError as exc:
|
|
raise HTTPException(409, str(exc))
|
|
try:
|
|
engine.add_camera(cam)
|
|
except Exception as exc:
|
|
# Never leave the store describing a camera the engine refused —
|
|
# the two would disagree until the next restart.
|
|
engine.camera_store.delete(cam.id)
|
|
raise HTTPException(500, f"camera stored but failed to start: {exc}")
|
|
return camera_public(cam, engine.workers.get(cam.id))
|
|
|
|
@app.patch("/api/cameras/{camera_id}")
|
|
def edit_camera(camera_id: str, payload: CameraPayload) -> dict:
|
|
fields = payload.model_dump(exclude_none=True)
|
|
fields.pop("id", None)
|
|
try:
|
|
if "tuning" in fields:
|
|
engine.cfg.recognition.merged(
|
|
CameraTuning.model_validate(fields["tuning"]))
|
|
cam = engine.camera_store.update(camera_id, fields)
|
|
except (ValidationError, ValueError) as exc:
|
|
raise HTTPException(400, str(exc))
|
|
if cam is None:
|
|
raise HTTPException(404, f"unknown camera '{camera_id}'")
|
|
engine.restart_camera(cam) # a changed URL needs a fresh connection
|
|
return camera_public(cam, engine.workers.get(cam.id))
|
|
|
|
@app.delete("/api/cameras/{camera_id}")
|
|
def delete_camera(camera_id: str) -> dict:
|
|
if not engine.camera_store.delete(camera_id):
|
|
raise HTTPException(404, f"unknown camera '{camera_id}'")
|
|
engine.remove_camera(camera_id)
|
|
return {"deleted": camera_id}
|
|
|
|
@app.post("/api/cameras/{camera_id}/commission")
|
|
def start_commission(camera_id: str,
|
|
payload: CommissionPayload) -> dict:
|
|
"""Begin a placement check: watch this camera for N seconds and judge
|
|
whether faces here are good enough to enrol."""
|
|
worker = worker_or_404(camera_id)
|
|
# The camera's own gate, not the global one - the whole point is to
|
|
# judge this view against the threshold it will actually run under.
|
|
worker.commission = CommissionRun(
|
|
camera_id, worker.rcfg.min_enroll_quality, payload.seconds)
|
|
return worker.commission.report()
|
|
|
|
@app.get("/api/cameras/{camera_id}/commission")
|
|
def commission_result(camera_id: str) -> dict:
|
|
worker = worker_or_404(camera_id)
|
|
if worker.commission is None:
|
|
raise HTTPException(404, "no placement check has been run")
|
|
return worker.commission.report()
|
|
|
|
@app.delete("/api/cameras/{camera_id}/commission")
|
|
def cancel_commission(camera_id: str) -> dict:
|
|
worker = worker_or_404(camera_id)
|
|
if worker.commission is not None:
|
|
worker.commission.cancel()
|
|
return {"cancelled": camera_id}
|
|
|
|
@app.get("/api/cameras/discover")
|
|
def discover_cameras() -> dict:
|
|
"""Cameras on this PC's network, for the add-camera form to pick from.
|
|
|
|
A sync def so FastAPI runs it in the threadpool: it holds a socket
|
|
open for a couple of seconds and sweeps a /24, and the event loop
|
|
must keep serving the live picture meanwhile.
|
|
"""
|
|
from .discover import discover
|
|
return discover()
|
|
|
|
@app.post("/api/cameras/test")
|
|
def test_camera(payload: CameraPayload) -> dict:
|
|
"""Try a camera WITHOUT saving it - the UI's Test button.
|
|
|
|
Deliberately a sync def so FastAPI runs it in the threadpool:
|
|
cv2.VideoCapture blocks hard and a wrong host can hang for the full
|
|
FFmpeg timeout, which would stall the whole event loop.
|
|
"""
|
|
from .capture import probe_source
|
|
|
|
data = payload.model_dump(exclude_none=True)
|
|
data.setdefault("id", "__test__")
|
|
try:
|
|
cam = CameraConfig.model_validate(data)
|
|
source = cam.source()
|
|
except (ValidationError, ValueError) as exc:
|
|
return {"ok": False, "error": str(exc)}
|
|
return probe_source(source, cam.max_width)
|
|
|
|
@app.get("/api/cameras/{camera_id}/frame.jpg")
|
|
def frame(camera_id: str, width: int = 0, quality: int = 0) -> Response:
|
|
"""The latest frame, optionally re-encoded smaller.
|
|
|
|
`width`/`quality` exist for the live relay, which sends several frames
|
|
a second up a shop's uplink and cannot afford the full-size picture the
|
|
dashboard uses. The re-encode happens here rather than in the agent
|
|
because this process already has OpenCV open and the frame decoded;
|
|
shipping a scaler into the agent would be the same work done twice.
|
|
|
|
Done on demand, not on every frame: a camera nobody is watching must
|
|
not pay for a second encode it will never use.
|
|
"""
|
|
jpeg = worker_or_404(camera_id).latest_jpeg()
|
|
if jpeg is None:
|
|
raise HTTPException(503, "no frame yet")
|
|
if width > 0 or quality > 0:
|
|
jpeg = _reencode(jpeg, width, quality) or jpeg
|
|
return Response(jpeg, media_type="image/jpeg")
|
|
|
|
@app.get("/api/cameras/{camera_id}/stream.mjpeg")
|
|
async def stream(camera_id: str) -> StreamingResponse:
|
|
worker = worker_or_404(camera_id)
|
|
|
|
async def generate():
|
|
boundary = b"--frame\r\nContent-Type: image/jpeg\r\n\r\n"
|
|
# 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():
|
|
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(),
|
|
media_type="multipart/x-mixed-replace; boundary=frame")
|
|
|
|
return app
|