Files
livekit-cameras/portrait_daemon.py
root 95bb7c50ee feat: point the fleet at wss LiveKit room uwh-telhai
Publish and subscribe over wss://livekit.uni-wh.de:7800 and refuse
cleartext ws://. Conference room is uwh-telhai. Includes the uncommitted
encoded H.264 publish path, Rally hairpin, and KMS wall overlay.
2026-10-11 00:48:50 +00:00

234 lines
7.0 KiB
Python

"""C920 I420 SHM -> person mask -> blurred background SHM.
OpenVINO CPU selfie seg. Does not open DRM. Rally is not processed.
"""
from __future__ import annotations
import os
os.environ.setdefault("OMP_NUM_THREADS", "1")
os.environ.setdefault("OPENBLAS_NUM_THREADS", "1")
os.environ.setdefault("MKL_NUM_THREADS", "1")
import argparse
import json
import logging
import signal
import threading
import time
from dataclasses import dataclass, field
from pathlib import Path
from typing import Callable
import numpy as np
from frame_shm import LatestFrameReader, LatestFrameWriter, wait_for_ready
from portrait import (
MaskHold,
PersonSeg,
apply_i420,
i420_nbytes,
i420_to_bgr,
person_present,
)
log = logging.getLogger("cameras.portrait-daemon")
READY_NAME = "ready"
InferFn = Callable[[np.ndarray], np.ndarray]
@dataclass
class CamSlot:
identity: str
width: int
height: int
reader: LatestFrameReader
writer: LatestFrameWriter
infer_hz: float = 8.0
hold_s: float = 0.8
blur_px: int = 7
last_seq: int = -1
last_infer: float = 0.0
hold: MaskHold = field(default_factory=MaskHold)
scratch: bytearray = field(default_factory=bytearray)
def __post_init__(self) -> None:
self.hold = MaskHold(hold_s=self.hold_s, ema=0.4)
n = i420_nbytes(self.width, self.height)
self.scratch = bytearray(n)
self.infer_scratch = bytearray(n)
self.infer_reader = LatestFrameReader(self.reader.path)
def tick_once(
slots: list[CamSlot],
infer: InferFn,
now: float,
infer_budget_s: float = 0.008,
) -> int:
"""Infer a few due cameras, composite every new SHM frame. Returns writes."""
wrote = 0
deadline = time.monotonic() + max(0.001, float(infer_budget_s))
pending: list[tuple[CamSlot, int]] = []
for s in sorted(slots, key=lambda c: c.last_infer):
seq = s.reader.copy_into(s.scratch)
if seq is None or seq == s.last_seq:
continue
due = (now - s.last_infer) >= (1.0 / max(0.5, s.infer_hz))
if due and time.monotonic() < deadline:
bgr = i420_to_bgr(s.scratch, s.width, s.height)
raw = infer(bgr)
present = person_present(raw)
s.hold.update(raw if present else None, present, now)
s.last_infer = now
pending.append((s, seq))
for s, seq in pending:
out = apply_i420(
s.scratch, s.width, s.height, s.hold.mask, radius=s.blur_px)
s.writer.write(out)
s.last_seq = seq
wrote += 1
return wrote
def composite_once(slots: list[CamSlot]) -> int:
"""30 fps path: no OpenVINO."""
wrote = 0
for s in slots:
seq = s.reader.copy_into(s.scratch)
if seq is None or seq == s.last_seq:
continue
out = apply_i420(
s.scratch, s.width, s.height, s.hold.mask, radius=s.blur_px)
s.writer.write(out)
s.last_seq = seq
wrote += 1
return wrote
def infer_once(slots: list[CamSlot], infer: InferFn, now: float, budget_s: float = 0.012) -> int:
n = 0
deadline = time.monotonic() + max(0.001, float(budget_s))
for s in sorted(slots, key=lambda c: c.last_infer):
if time.monotonic() >= deadline:
break
if (now - s.last_infer) < (1.0 / max(0.5, s.infer_hz)):
continue
seq = s.infer_reader.copy_into(s.infer_scratch)
if seq is None:
continue
bgr = i420_to_bgr(s.infer_scratch, s.width, s.height)
raw = infer(bgr)
present = person_present(raw)
s.hold.update(raw if present else None, present, now)
s.last_infer = now
n += 1
return n
def main(argv: list[str] | None = None) -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--cameras-json", required=True)
ap.add_argument("--in-dir", default="/run/livekit-cameras/raw")
ap.add_argument("--out-dir", default="/run/livekit-cameras/portrait")
ap.add_argument("--hz", type=float, default=8.0)
ap.add_argument("--blur-px", type=int, default=7)
ap.add_argument("--hold-s", type=float, default=0.8)
ap.add_argument("--model", default="")
args = ap.parse_args(argv)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [portrait-daemon] %(levelname)s %(message)s")
import cv2
cv2.setNumThreads(1)
os.environ.setdefault("OMP_NUM_THREADS", "1")
os.environ.setdefault("OPENBLAS_NUM_THREADS", "1")
cams = json.loads(args.cameras_json)
cams = [c for c in cams if c.get("identity") != "rally"]
if not cams:
log.error("no C920 cameras")
return 2
model_path = Path(args.model) if args.model else None
seg = PersonSeg(model_path=model_path, num_threads=2)
log.info(
"portrait OpenVINO %s model=%s cams=%d hz=%.1f blur=%d hold=%.2fs",
seg.device,
(model_path or Path("models/selfie_segmentation.onnx")).name,
len(cams), args.hz, args.blur_px, args.hold_s,
)
os.makedirs(args.out_dir, exist_ok=True)
slots: list[CamSlot] = []
for c in cams:
ident = c["identity"]
w, h = int(c["width"]), int(c["height"])
src = os.path.join(args.in_dir, f"{ident}.i420")
dst = os.path.join(args.out_dir, f"{ident}.i420")
if not wait_for_ready(src, timeout=20):
log.warning("missing capture SHM %s", src)
continue
slots.append(CamSlot(
identity=ident, width=w, height=h,
reader=LatestFrameReader(src),
writer=LatestFrameWriter(dst, w, h),
infer_hz=float(args.hz),
hold_s=float(args.hold_s),
blur_px=int(args.blur_px),
))
if not slots:
log.error("no portrait slots")
return 2
ready = os.path.join(args.out_dir, READY_NAME)
Path(ready).write_text("ok\n")
stop = False
t0 = time.monotonic()
writes = 0
infers = 0
stop_ev = threading.Event()
def _stop(*_a):
nonlocal stop
stop = True
stop_ev.set()
signal.signal(signal.SIGTERM, _stop)
signal.signal(signal.SIGINT, _stop)
def _infer_loop() -> None:
nonlocal infers
while not stop_ev.is_set():
infers += infer_once(slots, seg.infer_bgr, time.monotonic(), 0.012)
time.sleep(0.004)
th = threading.Thread(target=_infer_loop, name="portrait-infer", daemon=True)
th.start()
try:
while not stop:
n = composite_once(slots)
writes += n
if time.monotonic() - t0 >= 10:
log.info(
"writes=%d infers=%d last_seq %s",
writes, infers,
" ".join(f"{s.identity}:{s.last_seq}" for s in slots),
)
writes = 0
infers = 0
t0 = time.monotonic()
if n == 0:
time.sleep(0.003)
finally:
stop_ev.set()
try:
os.unlink(ready)
except OSError:
pass
return 0
if __name__ == "__main__":
raise SystemExit(main())