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.
403 lines
16 KiB
Python
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())
|