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.
234 lines
7.0 KiB
Python
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())
|