Files
livekit-cameras/capture_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

177 lines
5.6 KiB
Python

"""Shared C920 capture: small gst-launch pool, latest I420 frames in mmap."""
from __future__ import annotations
import argparse
import json
import logging
import os
import select
import signal
import subprocess
import sys
import time
from capture_va import (
GST_GROUP_SIZE,
grouped_cameras,
shared_capture_cmd,
)
from frame_shm import LatestFrameWriter, frame_bytes
from uvc_exposure import pin_devices
log = logging.getLogger("cameras.capture_daemon")
READY_NAME = "ready"
def _open_pipes(n: int) -> tuple[list[int], list[int]]:
reads: list[int] = []
writes: list[int] = []
for _ in range(n):
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)
reads.append(r)
writes.append(w)
return reads, writes
def main(argv: list[str] | None = None) -> int:
ap = argparse.ArgumentParser(description="Shared JPEG capture daemon")
ap.add_argument("--cameras-json", required=True,
help="JSON list of {identity, video_device}")
ap.add_argument("--width", type=int, default=640)
ap.add_argument("--height", type=int, default=360)
ap.add_argument("--fps", type=int, default=30)
ap.add_argument("--shm-dir", default="/run/livekit-cameras/raw")
ap.add_argument("--vaapi-device", default="/dev/dri/renderD129")
ap.add_argument("--group-size", type=int, default=GST_GROUP_SIZE)
args = ap.parse_args(argv)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [capture-daemon] %(levelname)s %(message)s")
cams = json.loads(args.cameras_json)
pairs = [(c["identity"], c["video_device"]) for c in cams]
if not pairs:
log.error("no cameras")
return 2
os.environ["LIBVA_DRIVER_NAME"] = "iHD"
need = frame_bytes(args.width, args.height)
os.makedirs(args.shm_dir, exist_ok=True)
writers = [
LatestFrameWriter(os.path.join(args.shm_dir, f"{ident}.i420"),
args.width, args.height)
for ident, _dev in pairs
]
reads, writes = _open_pipes(len(pairs))
env = dict(os.environ)
env["LIBVA_DRM_DEVICE"] = args.vaapi_device
env["LIBVA_DRIVER_NAME"] = "iHD"
env["GST_REGISTRY_UPDATE"] = "no"
groups = grouped_cameras(pairs, writes, args.group_size)
log.info("starting %d gst-launch for %d cameras %dx%d@%d (group=%d)",
len(groups), len(pairs), args.width, args.height, args.fps,
args.group_size)
procs: list[subprocess.Popen] = []
for g_cams, g_fds in groups:
cmd = shared_capture_cmd(g_cams, args.width, args.height, args.fps, g_fds)
procs.append(subprocess.Popen(
cmd,
stdin=subprocess.DEVNULL,
stdout=subprocess.DEVNULL,
stderr=None,
pass_fds=tuple(g_fds),
env=env,
close_fds=True,
))
for w in writes:
os.close(w)
for r in reads:
os.set_blocking(r, False)
pin_devices([dev for _ident, dev in pairs])
bufs = [bytearray() for _ in pairs]
seqs = [0] * len(pairs)
stop = False
def _stop(_signum=None, _frame=None):
nonlocal stop
stop = True
signal.signal(signal.SIGTERM, _stop)
signal.signal(signal.SIGINT, _stop)
ready_path = os.path.join(args.shm_dir, READY_NAME)
if os.path.exists(ready_path):
os.unlink(ready_path)
marked = False
last_log = time.monotonic()
fd_index = {r: i for i, r in enumerate(reads)}
try:
while not stop:
dead = [p for p in procs if p.poll() is not None]
if dead:
log.error("gst exited rc=%s", [p.returncode for p in dead])
break
ready, _, _ = select.select(reads, [], [], 0.2)
now_ns = time.time_ns()
for r in ready:
i = fd_index[r]
while True:
try:
chunk = os.read(r, max(need, need - len(bufs[i])))
except BlockingIOError:
break
if not chunk:
break
bufs[i].extend(chunk)
while len(bufs[i]) >= need:
frame = bytearray(bufs[i][:need])
del bufs[i][:need]
seqs[i] = writers[i].write(frame, now_ns)
if not marked and any(s > 0 for s in seqs):
pin_devices([dev for _ident, dev in pairs])
open(ready_path, "w").write("ok\n")
marked = True
log.info("ready first frames %s",
",".join(f"{pairs[i][0]}={seqs[i]}" for i in range(len(pairs))))
now = time.monotonic()
if now - last_log >= 5.0:
last_log = now
log.info("seqs %s",
" ".join(f"{pairs[i][0]}:{seqs[i]}" for i in range(len(pairs))))
finally:
for proc in procs:
if proc.poll() is None:
proc.terminate()
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
proc.kill()
for r in reads:
try:
os.close(r)
except OSError:
pass
for w in writers:
w.close()
if os.path.exists(ready_path):
os.unlink(ready_path)
return 0
if __name__ == "__main__":
sys.exit(main())