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.
177 lines
5.6 KiB
Python
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())
|