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

357 lines
14 KiB
Python

"""One process: N LiveKit identities, shared video SHM + audio hop rings."""
from __future__ import annotations
import argparse
import asyncio
import json
import logging
import os
import signal
import sys
import time
import numpy as np
log = logging.getLogger("cameras.publisher")
def _rss_anon_kb() -> int:
try:
with open("/proc/self/status", encoding="utf-8") as f:
for line in f:
if line.startswith("RssAnon:"):
return int(line.split()[1])
except OSError:
return 0
return 0
def _pub_opts(rtc, args, height: int = 0):
from livekit.rtc import TrackPublishOptions, VideoEncoding, VideoCodec, TrackSource
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_AUTO)
codec = codec_map.get(args.video_codec, VideoCodec.H264)
br = int(args.bitrate)
if int(height) >= 720:
br = max(br, int(os.environ.get("VIDEO_RALLY_BITRATE", "20000000")))
video = TrackPublishOptions(
source=TrackSource.SOURCE_CAMERA,
simulcast=False,
video_codec=codec,
video_encoding=VideoEncoding(
max_bitrate=br,
max_framerate=int(args.fps),
),
video_encoder=encoder,
degradation_preference=DEGRADATION_PREFERENCE_MAINTAIN_RESOLUTION,
)
audio = TrackPublishOptions(source=TrackSource.SOURCE_MICROPHONE, dtx=False)
return video, audio
async def publish_one(rtc, cam: dict, args, stop: asyncio.Event,
connect_gate: asyncio.Semaphore, enc_state: dict) -> None:
from audio_cleanup import AudioCleaner
from frame_shm import LatestFrameReader, frame_bytes
from hop_shm import HopReader
from audio_gate import make_gate
from latency import (
audio_queue_size_ms,
effective_enhance_mode,
next_video_capture,
should_recycle_encoders,
video_capture_gate_acquire,
video_capture_gate_release,
)
identity = cam["identity"]
clog = logging.getLogger(f"cameras.publisher.{identity}")
w = int(cam.get("width") or args.width)
h = int(cam.get("height") or args.height)
if args.video_codec in ("av1", "h264"):
h = h - (h % 16)
w = w - (w % 16)
enhance_mode = effective_enhance_mode(
args.enhance_mode, identity, args.enhance_speaker_only,
args.speaker_camera)
beamform = args.beamform_mode
if args.enhance_speaker_only and enhance_mode == "off":
beamform = "off"
cleaner = AudioCleaner(
mode=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=beamform,
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,
),
)
clog.info("audio backend=%s enhance=%s beamform=%s gate=%s",
cleaner.backend, enhance_mode, beamform,
"on" if cleaner._gate is not None else "off")
room = rtc.Room()
nbytes = frame_bytes(w, h)
shm_path = os.path.join(args.video_shm_dir, f"{identity}.i420")
shm = LatestFrameReader(shm_path)
scratch = bytearray(nbytes)
video_frames = 0
video_gate = {"in_flight": 0, "dropped": 0, "max_in_flight": 3}
from encoded_pub import capture_encoded
from nal_shm import NalReader
if args.video_codec == "av1":
from av1_obu import is_keyframe
else:
from annexb import is_keyframe
nal = NalReader(os.path.join(args.h264_shm_dir, f"{identity}.h264"))
clog.info("video path=encoded-preencoded identity=%s codec=%s h264=%s",
identity, args.video_codec, nal.path)
audio_card = cam.get("audio_card")
hops = HopReader(os.path.join(args.audio_shm_dir, f"{identity}.pcm")) if audio_card is not None else None
atrack = None
audio_source = None
v_opts, a_opts = _pub_opts(rtc, args, h)
async with connect_gate:
await room.connect(args.url, cam["token"], rtc.RoomOptions(auto_subscribe=False))
clog.info("connected room=%r", room.name)
video_source = rtc.VideoSource(w, h)
video_track = rtc.LocalVideoTrack.create_video_track(
f"{identity}-video", video_source)
if hops 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)
await room.local_participant.publish_track(video_track, v_opts)
if atrack is not None:
await room.local_participant.publish_track(atrack, a_opts)
clog.info("status: UP identity=%s video_shm=%s audio=%s",
identity, shm_path,
f"hw:{audio_card},0" if audio_card is not None else "none")
async def recycle_video():
nonlocal video_source, video_track
sid = getattr(video_track, "sid", None) or ""
clog.warning("encoder recycle rss_anon_kb=%d sid=%s",
_rss_anon_kb(), sid)
if sid:
try:
await room.local_participant.unpublish_track(sid)
except Exception: # noqa: BLE001
clog.exception("unpublish for encoder recycle failed")
stagger = max(0, int(cam.get("index") or 0)) * 0.05
if stagger:
await asyncio.sleep(stagger)
video_source = rtc.VideoSource(w, h)
video_track = rtc.LocalVideoTrack.create_video_track(
f"{identity}-video", video_source)
await room.local_participant.publish_track(video_track, v_opts)
async def video_loop():
nonlocal video_frames, video_source, video_track
period = 1.0 / max(1, int(args.fps))
last_seq = -1
due = time.monotonic()
local_gen = enc_state["gen"]
loop = asyncio.get_running_loop()
while not stop.is_set():
now = time.monotonic()
capture, due = next_video_capture(now, due, period)
if capture:
au = nal.read()
if au is not None:
last_seq, payload = au
try:
rc = capture_encoded(
video_source,
payload,
w,
h,
is_keyframe(payload),
timestamp_us=int(time.monotonic() * 1_000_000),
codec=args.video_codec,
)
if rc == 0:
video_frames += 1
elif video_frames % 90 == 1:
clog.warning("capture_encoded rc=%s seq=%s", rc, last_seq)
except Exception: # noqa: BLE001
clog.exception("capture_encoded failed")
if args.encoder_rss_limit_kb and video_frames % 90 == 0:
rss = _rss_anon_kb()
async with enc_state["lock"]:
if should_recycle_encoders(
rss, args.encoder_rss_limit_kb, now,
enc_state["last"],
args.encoder_recycle_interval_s):
enc_state["last"] = now
enc_state["gen"] += 1
if local_gen < enc_state["gen"]:
await recycle_video()
local_gen = enc_state["gen"]
delay = due - time.monotonic()
if delay > 0:
await asyncio.sleep(delay)
async def audio_loop():
assert hops is not None and audio_source is not None
hop_n = int(round(args.enhance_hop_s * args.audio_rate))
n = 0
loop = asyncio.get_running_loop()
while not stop.is_set():
blob = hops.read()
if blob is None:
await asyncio.sleep(0.005)
continue
try:
pcm = np.frombuffer(blob, dtype=np.int16)
if pcm.size % 2 == 0:
pcm = pcm.reshape(-1, 2)
cleaned = await loop.run_in_executor(None, cleaner.process, pcm)
await audio_source.capture_frame(
rtc.AudioFrame(
cleaned.tobytes(), args.audio_rate, 1, len(cleaned)))
n += 1
if n % 50 == 0:
clog.info("audio hop process_ms=%.1f", cleaner.last_dt_ms)
except Exception: # noqa: BLE001
clog.exception("audio hop failed")
tasks = [asyncio.create_task(video_loop(), name=f"{identity}-video")]
if hops is not None:
tasks.append(asyncio.create_task(audio_loop(), name=f"{identity}-audio"))
try:
await stop.wait()
finally:
for t in tasks:
t.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
clog.info("stopped; delivered %d video frames", video_frames)
try:
await room.disconnect()
except Exception: # noqa: BLE001
pass
async def amain(args) -> int:
os.environ["LIBVA_DRIVER_NAME"] = "iHD"
os.environ.pop("LIBVA_DRM_DEVICE", None)
import livekit.rtc as rtc
cams = json.loads(args.cameras_json)
stop = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGTERM, signal.SIGINT):
try:
loop.add_signal_handler(sig, stop.set)
except NotImplementedError:
pass
log.info("publisher start cams=%d shm=%s pcm=%s speaker_only=%s "
"encoder_rss_limit_kb=%d",
len(cams), args.video_shm_dir, args.audio_shm_dir,
args.enhance_speaker_only, args.encoder_rss_limit_kb)
gate = asyncio.Semaphore(1)
enc_state = {"lock": asyncio.Lock(), "last": 0.0, "gen": 0}
tasks = [
asyncio.create_task(
publish_one(rtc, cam, args, stop, gate, enc_state),
name=cam["identity"])
for cam in cams
]
await stop.wait()
for t in tasks:
t.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
return 0
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--cameras-json", required=True)
ap.add_argument("--url", required=True)
ap.add_argument("--room", required=True)
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("--bitrate", type=int, default=1_200_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("--enhance-socket", default="")
ap.add_argument("--video-shm-dir", default="/run/livekit-cameras/raw")
ap.add_argument("--audio-shm-dir", default="/run/livekit-cameras/pcm")
ap.add_argument("--h264-shm-dir", default="/run/livekit-cameras/h264")
ap.add_argument("--enhance-speaker-only", action="store_true")
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("--encoder-rss-limit-kb", type=int, default=0,
help="Recycle VAAPI VideoSource tracks when RssAnon exceeds this (0=off)")
ap.add_argument("--encoder-recycle-interval-s", type=float, default=120.0)
args = ap.parse_args()
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(name)s] %(levelname)s %(message)s")
try:
return asyncio.run(amain(args))
except Exception: # noqa: BLE001
log.exception("publisher fatal")
return 1
if __name__ == "__main__":
sys.exit(main())