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

403 lines
16 KiB
Python

"""One camera -> one LiveKit participant.
Run as a subprocess per camera: python -m cameras.run_worker --index 0 ...
"""
from __future__ import annotations
import argparse
import asyncio
import concurrent.futures
import contextlib
import json
import logging
import os
import signal
import subprocess
import sys
import time
from collections import deque
import numpy as np
log = logging.getLogger("cameras.worker")
def drain_held_frames(pending: deque, now: float, hold_s: float):
"""Yield frame payloads captured at least hold_s seconds ago."""
while pending and (now - pending[0][0]) >= hold_s:
yield pending.popleft()[1]
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--camera-json", required=True,
help="JSON of one discovered camera dict")
ap.add_argument("--room", required=True)
ap.add_argument("--url", required=True)
ap.add_argument("--token", required=True)
ap.add_argument("--width", type=int, default=320)
ap.add_argument("--height", type=int, default=180)
ap.add_argument("--fps", type=int, default=15)
ap.add_argument("--bitrate", type=int, default=400_000)
ap.add_argument("--video-codec", default="h264")
ap.add_argument("--video-encoder", default="vaapi")
ap.add_argument("--vaapi-device", default="/dev/dri/renderD129")
ap.add_argument("--audio-rate", type=int, default=16000)
ap.add_argument("--enhance-mode", default="auto")
ap.add_argument("--enhance-model",
default="speechbrain/metricgan-plus-voicebank")
ap.add_argument("--enhance-chunk-s", type=float, default=1.0)
ap.add_argument("--enhance-hop-s", type=float, default=0.1)
ap.add_argument("--enhance-infer", default="auto")
ap.add_argument("--beamform-mode", default="auto")
ap.add_argument("--torch-num-threads", type=int, default=1)
ap.add_argument("--tdoa-every", type=int, default=5)
ap.add_argument("--audio-queue-ms", type=int, default=None)
ap.add_argument("--video-hold-s", type=float, default=None)
ap.add_argument("--split-executor", action="store_true")
ap.add_argument("--enhance-socket", default="")
ap.add_argument("--speaker-camera", default="rally")
ap.add_argument("--audio-gate", action="store_true")
ap.add_argument("--audio-gate-open-db", type=float, default=-28.0)
ap.add_argument("--audio-gate-close-db", type=float, default=-34.0)
ap.add_argument("--audio-gate-hold-s", type=float, default=0.3)
ap.add_argument("--video-shm", default="",
help="latest-frame I420 mmap from capture_daemon")
args = ap.parse_args()
cam = json.loads(args.camera_json)
logging.basicConfig(
level=logging.INFO,
format=f"%(asctime)s [{cam['identity']}] %(levelname)s %(message)s",
)
from config import Config # noqa: F401 (ensures .env is loaded)
from audio_cleanup import AudioCleaner, apply_torch_thread_limits
from audio_gate import make_gate
from latency import audio_queue_size_ms, video_hold_seconds
identity = cam["identity"]
os.environ["LIBVA_DRIVER_NAME"] = "iHD"
os.environ.pop("LIBVA_DRM_DEVICE", None)
if not args.enhance_socket:
apply_torch_thread_limits(args.torch_num_threads)
import livekit.rtc as rtc
video_dev = cam["video_device"]
audio_card = cam.get("audio_card")
async def run() -> int:
cleaner = AudioCleaner(
mode=args.enhance_mode,
model_source=args.enhance_model,
sample_rate=args.audio_rate,
chunk_s=args.enhance_chunk_s,
hop_s=args.enhance_hop_s,
beamform=args.beamform_mode,
torch_num_threads=args.torch_num_threads,
tdoa_every=args.tdoa_every,
enhance_infer=args.enhance_infer,
enhance_socket=args.enhance_socket,
identity=identity,
gate=make_gate(
identity,
args.speaker_camera,
bool(args.audio_gate),
args.audio_gate_open_db,
args.audio_gate_close_db,
args.audio_gate_hold_s,
args.enhance_hop_s,
),
)
log.info("audio backend: %s gate=%s", cleaner.backend,
"on" if cleaner._gate is not None else "off")
io_ex = None
enhance_ex = None
split_owned: list[concurrent.futures.Executor] = []
if args.split_executor:
io_ex = concurrent.futures.ThreadPoolExecutor(
max_workers=2, thread_name_prefix=f"{identity}-io")
enhance_ex = concurrent.futures.ThreadPoolExecutor(
max_workers=1, thread_name_prefix=f"{identity}-enh")
split_owned.extend([io_ex, enhance_ex])
log.info("split executor: ffmpeg io vs enhance")
room = rtc.Room()
await room.connect(
args.url, args.token, rtc.RoomOptions(auto_subscribe=False))
log.info("connected to room %r as %r (auto_subscribe=False)",
room.name, identity)
# ---------------- video ----------------
video_source = rtc.VideoSource(args.width, args.height)
video_track = rtc.LocalVideoTrack.create_video_track(
f"{identity}-video", video_source)
frame_bytes = args.width * args.height * 3 // 2
video_frames = 0
vproc: subprocess.Popen | None = None
shm_path = (args.video_shm or "").strip()
video_hold_s = video_hold_seconds(
args.enhance_hop_s,
audio_card is not None,
0.0 if shm_path else args.video_hold_s,
)
pending: deque[tuple[float, bytes]] = deque()
max_pending = max(8, int(args.fps * (video_hold_s + 0.15)) + 4)
if shm_path:
from frame_shm import LatestFrameReader
shm_reader = LatestFrameReader(shm_path)
log.info("video from shared capture shm=%s hold=0 (latest frame)",
shm_path)
async def video_loop():
nonlocal video_frames
period = 1.0 / max(1, int(args.fps))
last_seq = -1
next_t = time.monotonic()
scratch = bytearray(frame_bytes)
while True:
seq = shm_reader.copy_into(scratch)
if seq is not None and seq != last_seq:
last_seq = seq
try:
video_source.capture_frame(
rtc.VideoFrame(
args.width, args.height,
rtc.VideoBufferType.I420, scratch),
timestamp_us=int(time.monotonic() * 1_000_000),
)
video_frames += 1
except Exception: # noqa: BLE001
log.exception("video capture_frame failed")
next_t += period
delay = next_t - time.monotonic()
if delay > 0:
await asyncio.sleep(delay)
else:
next_t = time.monotonic()
else:
from capture_va import video_capture_cmd
vproc = subprocess.Popen(
video_capture_cmd(
video_dev, args.width, args.height, args.fps, fd=1),
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)
log.info("capture gst jpegdec Arc (no sidecar h264)")
if video_hold_s > 0:
log.info("av sync: holding video %.0f ms to match audio hop",
video_hold_s * 1000.0)
async def video_loop():
nonlocal video_frames
assert vproc is not None
loop = asyncio.get_running_loop()
stdout = vproc.stdout
assert stdout is not None
while vproc.poll() is None:
data = await loop.run_in_executor(
io_ex, stdout.read, frame_bytes)
if not data or len(data) < frame_bytes:
if not data:
break
keep = b""
while len(keep) < frame_bytes:
more = await loop.run_in_executor(
io_ex, stdout.read,
frame_bytes - len(keep))
if not more:
return
keep += more
data = keep
now = time.monotonic()
pending.append((now, data))
while len(pending) > max_pending:
pending.popleft()
for payload in drain_held_frames(pending, now, video_hold_s):
try:
video_source.capture_frame(
rtc.VideoFrame(
args.width, args.height,
rtc.VideoBufferType.I420, payload),
timestamp_us=int(time.monotonic() * 1_000_000),
)
video_frames += 1
except Exception: # noqa: BLE001
log.exception("video capture_frame failed")
err = b""
if vproc.stderr:
try:
err = vproc.stderr.read() or b""
except Exception: # noqa: BLE001
pass
log.error("capture gst exited rc=%s stderr=%s",
vproc.returncode,
err.decode("utf-8", "replace")[-800:])
# ---------------- audio ----------------
aproc: subprocess.Popen | None = None
atrack: rtc.LocalAudioTrack | None = None
if audio_card is not None:
audio_source = rtc.AudioSource(
sample_rate=args.audio_rate, num_channels=1,
queue_size_ms=audio_queue_size_ms(
args.enhance_hop_s, args.audio_queue_ms))
atrack = rtc.LocalAudioTrack.create_audio_track(
f"{identity}-audio", audio_source)
log.info("audio queue_size_ms=%d", audio_queue_size_ms(
args.enhance_hop_s, args.audio_queue_ms))
hop_s = args.enhance_hop_s
chunk = int(hop_s * args.audio_rate) # samples per hop
async def audio_loop():
assert aproc is not None and atrack is not None
loop = asyncio.get_running_loop()
# mic delivers stereo s16le: bytes per chunk =
# chunk * 2 ch * 2 bytes
need = chunk * 2 * 2
buf = b""
hops = 0
while aproc.poll() is None:
data = await loop.run_in_executor(
io_ex, aproc.stdout.read, 4096)
if not data:
break
buf += data
while len(buf) >= need:
piece, buf = buf[:need], buf[need:]
pcm_stereo = np.frombuffer(
piece, dtype=np.int16).reshape(-1, 2)
n_frames = pcm_stereo.shape[0]
usable = (n_frames // chunk) * chunk
if usable:
cleaned = await loop.run_in_executor(
enhance_ex, cleaner.process,
pcm_stereo[:usable])
t_cap = time.perf_counter()
await audio_source.capture_frame(
rtc.AudioFrame(
cleaned.tobytes(),
args.audio_rate,
1,
len(cleaned)))
hops += 1
if hops % 50 == 0:
log.info(
"audio hop process_ms=%.1f capture_ms=%.1f",
cleaner.last_dt_ms,
(time.perf_counter() - t_cap) * 1000.0,
)
aproc = subprocess.Popen(
["ffmpeg", "-hide_banner", "-loglevel", "error",
"-fflags", "nobuffer", "-flags", "low_delay",
"-thread_queue_size", "8",
"-f", "alsa", "-i", f"hw:{audio_card},0",
"-ar", str(args.audio_rate), "-ac", "2",
"-f", "s16le", "pipe:1"],
stdout=subprocess.PIPE,
)
# ---------------- publish ----------------
from livekit.rtc import TrackSource, TrackPublishOptions, VideoEncoding, VideoCodec
from livekit.rtc._proto.room_pb2 import (
ENCODER_BACKEND_AUTO,
ENCODER_BACKEND_SOFTWARE,
ENCODER_BACKEND_HARDWARE,
ENCODER_BACKEND_NVENC,
ENCODER_BACKEND_VAAPI,
ENCODER_BACKEND_VIDEOTOOLBOX,
DEGRADATION_PREFERENCE_MAINTAIN_RESOLUTION,
)
encoder_map = {
"auto": ENCODER_BACKEND_AUTO,
"software": ENCODER_BACKEND_SOFTWARE,
"hardware": ENCODER_BACKEND_HARDWARE,
"nvenc": ENCODER_BACKEND_NVENC,
"vaapi": ENCODER_BACKEND_VAAPI,
"videotoolbox": ENCODER_BACKEND_VIDEOTOOLBOX,
}
codec_map = {
"vp8": VideoCodec.VP8,
"h264": VideoCodec.H264,
"av1": VideoCodec.AV1,
"vp9": VideoCodec.VP9,
"h265": VideoCodec.H265,
}
encoder = encoder_map.get(args.video_encoder, ENCODER_BACKEND_VAAPI)
codec = codec_map.get(args.video_codec, VideoCodec.H264)
pub_opts = TrackPublishOptions(
source=TrackSource.SOURCE_CAMERA,
simulcast=False,
video_codec=codec,
video_encoding=VideoEncoding(
max_bitrate=int(args.bitrate),
max_framerate=int(args.fps),
),
video_encoder=encoder,
degradation_preference=DEGRADATION_PREFERENCE_MAINTAIN_RESOLUTION,
)
log.info(
"publish video %dx%d@%d %s encoder=%s bitrate=%d vaapi=%s",
args.width, args.height, args.fps, args.video_codec,
args.video_encoder, args.bitrate, args.vaapi_device,
)
await room.local_participant.publish_track(video_track, pub_opts)
if atrack is not None:
await room.local_participant.publish_track(
atrack,
TrackPublishOptions(
source=TrackSource.SOURCE_MICROPHONE,
dtx=False))
log.info("published video+audio (audio backend: %s)",
cleaner.backend)
else:
log.info("published video only (no audio card matched)")
log.info("status: UP identity=%s video=%s audio=%s",
identity, video_dev,
f"hw:{audio_card},0" if audio_card is not None else "none")
# ---------------- run ----------------
tasks = [asyncio.create_task(video_loop(), name="video")]
if aproc is not None:
tasks.append(asyncio.create_task(audio_loop(), name="audio"))
stop = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGTERM, signal.SIGINT):
with contextlib.suppress(NotImplementedError):
loop.add_signal_handler(sig, stop.set)
await stop.wait()
for t in tasks:
t.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
for ex in split_owned:
with contextlib.suppress(Exception):
ex.shutdown(wait=False, cancel_futures=True)
for p in (vproc, aproc):
if p is not None and p.poll() is None:
with contextlib.suppress(Exception):
p.terminate()
with contextlib.suppress(Exception):
p.wait(timeout=3)
log.info("stopped; delivered %d video frames", video_frames)
with contextlib.suppress(Exception):
await room.disconnect()
return 0
try:
return asyncio.run(run())
except Exception: # noqa: BLE001
log.exception("worker fatal")
return 1
if __name__ == "__main__":
sys.exit(main())