Files
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

581 lines
20 KiB
Python

"""Rally speaker follow: capture /dev/rally to I420 SHM + YuNet PTZ.
No X11. Local HDMI preview is a LiveKit subscriber (run_displays speaker
role on i915 HDMI-A-5). This process does not kmssink unless --connector
is set. Person tracking is host software. V-R0010 has no native VISCA.
"""
from __future__ import annotations
import argparse
import logging
import os
import select
import signal
import subprocess
import sys
import time
from pathlib import Path
import numpy as np
from drm_outputs import DrmOutput, list_outputs
from frame_shm import LatestFrameWriter, frame_bytes
from gst_sink import GST_LAUNCH, _kms_props
from hop_shm import HopWriter, hop_bytes
from visca_xu_bridge import RallyV4L2
log = logging.getLogger("cameras.rally_follow")
def pick_preview(name: str) -> DrmOutput:
for o in list_outputs():
if o.name == name:
return o
have = [f"{o.name}/{o.driver}" for o in list_outputs() if o.connected]
raise SystemExit(f"DRM connector {name!r} not found. connected={have}")
def video_cmd(device: str, width: int, height: int, fps: int,
fd: int, preview: DrmOutput | None = None) -> list[str]:
caps = f"image/jpeg,width={int(width)},height={int(height)},framerate={int(fps)}/1"
head = [
GST_LAUNCH, "-q",
"v4l2src", f"device={device}", "do-timestamp=false",
"!", caps,
"!", "jpegparse",
"!", "jpegdec", "qos=false",
]
i420 = [
"!", "videoconvert", "qos=false", "n-threads=2",
"!", f"video/x-raw,format=I420,width={int(width)},height={int(height)}",
"!", "fdsink", f"fd={int(fd)}", "sync=false",
]
if preview is None:
return head + i420
return head + [
"!", "tee", "name=t",
"t.", "!", "queue", "max-size-buffers=1", "max-size-bytes=0",
"max-size-time=0", "leaky=downstream",
"!", "videoconvert", "qos=false",
"!", "videoscale", "method=0",
"!", *_kms_props(preview),
"t.", "!", "queue", "max-size-buffers=1", "max-size-bytes=0",
"max-size-time=0", "leaky=downstream",
] + i420
def audio_cmd(card: int, rate: int, fd: int) -> list[str]:
return [
GST_LAUNCH, "-q",
"alsasrc", f"device=hw:{int(card)},0", "do-timestamp=false",
"!", "audioconvert",
"!", "audioresample",
"!", f"audio/x-raw,format=S16LE,rate={int(rate)},channels=2",
"!", "queue", "max-size-buffers=2", "max-size-bytes=0",
"max-size-time=0", "leaky=downstream",
"!", "fdsink", f"fd={int(fd)}", "sync=false",
]
_YUNET_2026 = Path(__file__).parent / "models" / "face_detection_yunet_2026may.onnx"
_YUNET_2023 = Path(__file__).parent / "models" / "face_detection_yunet_2023mar.onnx"
_YUNET = _YUNET_2026 if _YUNET_2026.is_file() else _YUNET_2023
# YuNet row: x,y,w,h, re, le, nose, rmouth, lmouth, score (15).
_SCORE_MIN = 0.80
_MIN_FRAC = 0.06
_MAX_FRAC = 0.50
_TOP_BAND = 0.10
_BOT_BAND = 0.90
_ASPECT_LO = 0.45
_ASPECT_HI = 1.40
_CONFIRM = 3
_HOLD_S = 2.0
_MOTION_MAD = 8.0
def resolve_yunet(models_dir: Path | None = None) -> Path:
root = Path(models_dir) if models_dir is not None else Path(__file__).parent / "models"
newer = root / "face_detection_yunet_2026may.onnx"
old = root / "face_detection_yunet_2023mar.onnx"
return newer if newer.is_file() else old
def _box_ok(x: float, y: float, w: float, h: float, score: float,
frame_w: int, frame_h: int) -> bool:
if score < _SCORE_MIN or w <= 1 or h <= 1:
return False
fw = float(frame_w)
fh = float(frame_h)
short = min(fw, fh)
if min(w, h) < _MIN_FRAC * short:
return False
if max(w, h) > _MAX_FRAC * max(fw, fh):
return False
ar = w / h
if ar < _ASPECT_LO or ar > _ASPECT_HI:
return False
cx = (x + w / 2.0) / fw
cy = (y + h / 2.0) / fh
if cy < _TOP_BAND or cy > _BOT_BAND:
return False
if cx < 0.02 or cx > 0.98:
return False
return True
def _frontal_landmarks(row, x: float, y: float, w: float, h: float) -> bool:
re_x, re_y = float(row[4]), float(row[5])
le_x, le_y = float(row[6]), float(row[7])
ns_x, ns_y = float(row[8]), float(row[9])
rm_x, rm_y = float(row[10]), float(row[11])
lm_x, lm_y = float(row[12]), float(row[13])
eye_y = 0.5 * (re_y + le_y)
mouth_y = 0.5 * (rm_y + lm_y)
if not (eye_y + 0.04 * h < ns_y < mouth_y - 0.04 * h):
return False
eye_lo = min(re_x, le_x)
eye_hi = max(re_x, le_x)
if not (eye_lo + 0.08 * w <= ns_x <= eye_hi - 0.08 * w):
return False
eye_dist = ((le_x - re_x) ** 2 + (le_y - re_y) ** 2) ** 0.5
if eye_dist < 0.22 * w or eye_dist > 0.85 * w:
return False
if abs(le_y - re_y) > 0.35 * h:
return False
return True
def _profile_landmarks(row, x: float, y: float, w: float, h: float) -> bool:
"""Side / turned head: landmarks clustered, nose not between the eyes."""
pts = [(float(row[i]), float(row[i + 1])) for i in range(4, 14, 2)]
pad_x, pad_y = 0.2 * w, 0.2 * h
inside = 0
for px, py in pts:
if (x - pad_x) <= px <= (x + w + pad_x) and (y - pad_y) <= py <= (y + h + pad_y):
inside += 1
if inside < 3:
return False
re_x, re_y = float(row[4]), float(row[5])
le_x, le_y = float(row[6]), float(row[7])
eye_dist = ((le_x - re_x) ** 2 + (le_y - re_y) ** 2) ** 0.5
if eye_dist >= 0.22 * w:
return False
ns_y = float(row[9])
rm_y, lm_y = float(row[11]), float(row[13])
eye_y = 0.5 * (re_y + le_y)
mouth_y = 0.5 * (rm_y + lm_y)
if mouth_y + 0.02 * h < eye_y:
return False
if ns_y + 0.12 * h < eye_y:
return False
return True
def accept_yunet_face(row, frame_w: int, frame_h: int) -> bool:
"""True if this YuNet hit looks like a real face (frontal or turned), not ceiling/stand."""
if row is None or len(row) < 15:
return False
x, y, w, h = (float(row[0]), float(row[1]), float(row[2]), float(row[3]))
score = float(row[14])
if not _box_ok(x, y, w, h, score, frame_w, frame_h):
return False
return _frontal_landmarks(row, x, y, w, h) or _profile_landmarks(row, x, y, w, h)
def select_yunet_face(faces, frame_w: int, frame_h: int):
"""Highest-score accepted face, or None."""
if faces is None or len(faces) == 0:
return None
best = None
best_score = -1.0
for row in faces:
if not accept_yunet_face(row, frame_w, frame_h):
continue
sc = float(row[14])
if sc > best_score:
best_score = sc
best = row
return best
class FaceGate:
"""Need `need` similar hits before the box is trusted for PTZ."""
def __init__(self, need: int = _CONFIRM, max_jump: float = 0.18) -> None:
self.need = max(1, int(need))
self.max_jump = float(max_jump)
self.hits = 0
self.cx = 0.0
self.cy = 0.0
def update(self, box: tuple[int, int, int, int] | None,
frame_w: int, frame_h: int) -> tuple[int, int, int, int] | None:
if box is None:
self.hits = 0
return None
x, y, w, h = box
cx = (x + w / 2.0) / float(frame_w)
cy = (y + h / 2.0) / float(frame_h)
if self.hits > 0:
jump = ((cx - self.cx) ** 2 + (cy - self.cy) ** 2) ** 0.5
if jump > self.max_jump:
self.hits = 1
else:
self.hits += 1
else:
self.hits = 1
self.cx, self.cy = cx, cy
if self.hits >= self.need:
return box
return None
def roi_motion(
prev: np.ndarray | None,
gray: np.ndarray | None,
box: tuple[int, int, int, int] | None,
*,
expand: float = 1.4,
min_mad: float = _MOTION_MAD,
) -> bool:
"""True if the last head ROI moved; motion elsewhere is ignored."""
if prev is None or gray is None or box is None:
return False
if prev.shape != gray.shape:
return False
h_img, w_img = gray.shape[:2]
x, y, w, h = (int(box[0]), int(box[1]), int(box[2]), int(box[3]))
extra_w = int((expand - 1.0) * w / 2.0)
extra_h = int((expand - 1.0) * h / 2.0)
x0 = max(0, x - extra_w)
y0 = max(0, y - extra_h)
x1 = min(w_img, x + w + extra_w)
y1 = min(h_img, y + h + extra_h)
if x1 - x0 < 4 or y1 - y0 < 4:
return False
a = prev[y0:y1, x0:x1].astype(np.float32)
b = gray[y0:y1, x0:x1].astype(np.float32)
return float(np.mean(np.abs(a - b))) >= float(min_mad)
class HeadHold:
"""Keep PTZ on the last confirmed face while the head ROI is still moving."""
def __init__(self, hold_s: float = _HOLD_S) -> None:
self.hold_s = max(0.0, float(hold_s))
self.box: tuple[int, int, int, int] | None = None
self.until = 0.0
def update(
self,
box: tuple[int, int, int, int] | None,
*,
moving: bool,
now: float,
) -> tuple[int, int, int, int] | None:
if box is not None:
self.box = box
self.until = now + self.hold_s
return box
if self.box is None:
return None
if moving:
self.until = now + self.hold_s
return self.box
if now < self.until:
return self.box
self.box = None
return None
def steer_speeds(
box: tuple[int, int, int, int],
frame_w: int,
frame_h: int,
deadband: float = 0.14,
) -> tuple[int, int]:
x, y, w, h = box
err_x = ((x + w / 2.0) - frame_w / 2.0) / (frame_w / 2.0)
err_y = ((y + h / 2.0) - frame_h / 2.0) / (frame_h / 2.0)
pan = 1 if err_x > deadband else (-1 if err_x < -deadband else 0)
tilt = -1 if err_y > deadband else (1 if err_y < -deadband else 0)
return pan, tilt
def lock_action(
confirmed: tuple[int, int, int, int] | None,
held: tuple[int, int, int, int] | None,
) -> str:
"""Steer only on a live face. A stale hold freezes PTZ (no ghost pan)."""
if confirmed is not None:
return "steer"
if held is not None:
return "freeze"
return "lost"
class Follower:
"""Keep a confirmed YuNet face near frame center. Hold on loss.
OpenCV 5 dropped Haar/HOG; YuNet FaceDetectorYN is the replacement.
Tiny / ceiling / stand / landmark-invalid hits are dropped before PTZ.
"""
def __init__(self, ptz: RallyV4L2, width: int, height: int,
deadband: float = 0.14, detect_every: int = 5):
self.ptz = ptz
self.w = int(width)
self.h = int(height)
self.deadband = float(deadband)
self.detect_every = max(1, int(detect_every))
self._n = 0
self._lost = 0
self._det = None
self._cv2 = None
self._gate = FaceGate()
self._hold = HeadHold()
self._frozen = False
self._gray: np.ndarray | None = None
self._prev: np.ndarray | None = None
self._dw = max(160, self.w // 2)
self._dh = max(90, self.h // 2)
model = resolve_yunet()
try:
import cv2
self._cv2 = cv2
if not model.is_file():
raise FileNotFoundError(model)
self._det = cv2.FaceDetectorYN_create(
str(model), "", (self._dw, self._dh), _SCORE_MIN, 0.3, 5000)
log.info("yunet face detector %s input=%dx%d score>=%.2f",
model.name, self._dw, self._dh, _SCORE_MIN)
except Exception as exc: # noqa: BLE001
log.warning("opencv tracker unavailable: %s", exc)
def _boxes(self, payload: bytes) -> list[tuple[int, int, int, int]]:
if self._cv2 is None or self._det is None:
return []
yuv = np.frombuffer(payload, dtype=np.uint8).reshape((self.h * 3 // 2, self.w))
bgr = self._cv2.cvtColor(yuv, self._cv2.COLOR_YUV2BGR_I420)
small = self._cv2.resize(bgr, (self._dw, self._dh))
self._gray = self._cv2.cvtColor(small, self._cv2.COLOR_BGR2GRAY)
self._det.setInputSize((self._dw, self._dh))
_ok, faces = self._det.detect(small)
picked = select_yunet_face(faces, self._dw, self._dh)
if picked is None:
return []
sx = self.w / float(self._dw)
sy = self.h / float(self._dh)
x, yy, ww, hh = (float(picked[0]), float(picked[1]),
float(picked[2]), float(picked[3]))
return [(int(x * sx), int(yy * sy), int(ww * sx), int(hh * sy))]
def on_i420(self, payload: bytes) -> None:
self._n += 1
if self._n % self.detect_every != 0:
return
boxes = self._boxes(payload)
now = time.monotonic()
hold_box = self._hold.box
if hold_box is not None and self._gray is not None:
sx = self._dw / float(self.w)
sy = self._dh / float(self.h)
small_box = (
int(hold_box[0] * sx), int(hold_box[1] * sy),
int(hold_box[2] * sx), int(hold_box[3] * sy),
)
else:
small_box = None
moving = roi_motion(self._prev, self._gray, small_box)
self._prev = None if self._gray is None else self._gray.copy()
face = boxes[0] if boxes else None
confirmed = self._gate.update(face, self.w, self.h)
held = self._hold.update(confirmed, moving=moving, now=now)
action = lock_action(confirmed, held)
if action == "steer" and confirmed is not None:
self._lost = 0
self._frozen = False
pan, tilt = steer_speeds(confirmed, self.w, self.h, self.deadband)
self.ptz.set_speed(pan, tilt)
return
if action == "freeze":
self._lost = 0
if not self._frozen:
self.ptz.stop()
self._frozen = True
return
self._lost += 1
if self._lost == 1 or self._lost == 6:
self.ptz.stop()
self._frozen = True
if self._lost == 6:
log.info("target lost; hold")
def _pipe() -> tuple[int, int]:
r, w = os.pipe()
try:
import fcntl
fcntl.fcntl(r, fcntl.F_SETPIPE_SZ, 1048576)
fcntl.fcntl(w, fcntl.F_SETPIPE_SZ, 1048576)
except OSError:
pass
os.set_inheritable(w, True)
os.set_inheritable(r, False)
return r, w
def main(argv: list[str] | None = None) -> int:
ap = argparse.ArgumentParser(description=__doc__)
ap.add_argument("--device", default="/dev/rally")
ap.add_argument("--connector", default="",
help="optional local kmssink; empty = LiveKit hairpin only")
ap.add_argument("--width", type=int, default=1920)
ap.add_argument("--height", type=int, default=1080)
ap.add_argument("--fps", type=int, default=30)
ap.add_argument("--shm-dir", default="/run/livekit-cameras/raw")
ap.add_argument("--audio-shm-dir", default="/run/livekit-cameras/pcm")
ap.add_argument("--audio-card", type=int, default=-1)
ap.add_argument("--audio-rate", type=int, default=16000)
ap.add_argument("--hop-s", type=float, default=0.1)
ap.add_argument("--identity", default="rally")
ap.add_argument("--no-follow", action="store_true")
args = ap.parse_args(argv)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [rally-follow] %(levelname)s %(message)s")
preview = pick_preview(args.connector) if args.connector else None
if preview is not None and preview.driver == "xe":
log.error("refusing to kmssink on Arc (%s); wall stays there", preview.name)
return 2
if preview is None:
log.info("preview via LiveKit hairpin (no local kmssink)")
else:
log.info("preview %s id=%s %s %s", preview.name, preview.connector_id,
preview.driver, preview.node)
need = frame_bytes(args.width, args.height)
os.makedirs(args.shm_dir, exist_ok=True)
video_w = LatestFrameWriter(
os.path.join(args.shm_dir, f"{args.identity}.i420"),
args.width, args.height)
hop_n = hop_bytes(args.audio_rate, 2, args.hop_s)
audio_w = None
audio_r = audio_wr = None
audio_proc = None
audio_buf = bytearray()
if args.audio_card >= 0:
os.makedirs(args.audio_shm_dir, exist_ok=True)
audio_w = HopWriter(
os.path.join(args.audio_shm_dir, f"{args.identity}.pcm"), hop_n)
audio_r, audio_wr = _pipe()
vr, vw = _pipe()
vcmd = video_cmd(args.device, args.width, args.height, args.fps, vw, preview)
log.info("gst video %s", " ".join(vcmd))
vproc = subprocess.Popen(
vcmd, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL,
stderr=None, pass_fds=(vw,))
os.close(vw)
if audio_r is not None and audio_wr is not None:
acmd = audio_cmd(args.audio_card, args.audio_rate, audio_wr)
log.info("gst audio %s", " ".join(acmd))
audio_proc = subprocess.Popen(
acmd, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL,
stderr=None, pass_fds=(audio_wr,))
os.close(audio_wr)
os.set_blocking(audio_r, False)
os.set_blocking(vr, False)
follower = None if args.no_follow else Follower(
RallyV4L2(args.device), args.width, args.height)
stop = False
def _stop(_signum=None, _frame=None):
nonlocal stop
stop = True
signal.signal(signal.SIGTERM, _stop)
signal.signal(signal.SIGINT, _stop)
vbuf = bytearray()
seq = 0
last = time.monotonic()
try:
while not stop:
if vproc.poll() is not None:
log.error("video gst exited rc=%s", vproc.returncode)
break
if audio_proc is not None and audio_proc.poll() is not None:
log.error("audio gst exited rc=%s", audio_proc.returncode)
break
fds = [vr]
if audio_r is not None:
fds.append(audio_r)
ready, _, _ = select.select(fds, [], [], 0.2)
if vr in ready:
while True:
try:
chunk = os.read(vr, max(need, need - len(vbuf)))
except BlockingIOError:
break
if not chunk:
break
vbuf.extend(chunk)
while len(vbuf) >= need:
frame = bytes(vbuf[:need])
del vbuf[:need]
seq = video_w.write(frame)
if follower is not None:
follower.on_i420(frame)
if audio_r is not None and audio_r in ready and audio_w is not None:
while True:
try:
chunk = os.read(audio_r, max(hop_n, hop_n - len(audio_buf)))
except BlockingIOError:
break
if not chunk:
break
audio_buf.extend(chunk)
while len(audio_buf) >= hop_n:
hop = bytes(audio_buf[:hop_n])
del audio_buf[:hop_n]
audio_w.write(hop)
now = time.monotonic()
if now - last >= 5.0:
last = now
log.info("video seq=%d preview=%s follow=%s",
seq, preview.name if preview else "livekit",
follower is not None)
finally:
if follower is not None:
follower.ptz.stop()
for proc in (vproc, audio_proc):
if proc is None or proc.poll() is not None:
continue
proc.terminate()
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
proc.kill()
video_w.close()
if audio_w is not None:
audio_w.close()
for fd in (vr, audio_r):
if fd is None:
continue
try:
os.close(fd)
except OSError:
pass
return 0
if __name__ == "__main__":
sys.exit(main())