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.
357 lines
14 KiB
Python
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())
|