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

590 lines
21 KiB
Python

"""Orchestrator: discover cameras and run one publisher worker each.
Usage:
python run.py # use .env (or environment variables)
python run.py --cameras 1,5 # publish only cameras 1 and 5
python run.py --once # run, print a status report, exit (for CI)
"""
from __future__ import annotations
import argparse
import asyncio
import contextlib
import json
import logging
import os
import signal
import subprocess
import sys
import time
from pathlib import Path
from config import load_config, _load_dotenv
from discovery import Camera, discover_cameras, select_cameras
from latency import effective_enhance_mode
from participant_tags import token_attributes
from rally import discover_rally
from tokens import make_token
log = logging.getLogger("cameras.main")
def spawn_enhance_daemon(cfg) -> subprocess.Popen:
from enhance_daemon import wait_for_socket
cmd = [
sys.executable, str(Path(__file__).parent / "enhance_daemon.py"),
"--socket", cfg.enhance_socket,
"--enhance-mode", "off" if cfg.enhance_mode == "off" else "force",
"--enhance-model", cfg.enhance_model,
"--enhance-chunk-s", str(cfg.enhance_chunk_s),
"--enhance-hop-s", str(cfg.enhance_hop_s),
"--enhance-infer", cfg.enhance_infer,
"--beamform-mode", cfg.beamform_mode,
"--torch-num-threads", str(max(2, cfg.torch_num_threads)),
"--tdoa-every", str(cfg.beamform_tdoa_every),
"--audio-rate", str(cfg.audio_rate),
]
log.info("starting shared enhance daemon at %s", cfg.enhance_socket)
proc = subprocess.Popen(cmd)
if not wait_for_socket(cfg.enhance_socket, timeout=90):
proc.kill()
raise RuntimeError(
f"enhance daemon did not listen on {cfg.enhance_socket}")
log.info("enhance daemon ready pid=%s", proc.pid)
return proc
def spawn_capture_daemon(cfg, cameras) -> subprocess.Popen:
from frame_shm import wait_for_ready
payload = [{"identity": c.identity, "video_device": c.video_device}
for c in cameras]
cmd = [
sys.executable, str(Path(__file__).parent / "capture_daemon.py"),
"--cameras-json", json.dumps(payload),
"--width", str(cfg.width),
"--height", str(cfg.height),
"--fps", str(cfg.fps),
"--shm-dir", cfg.capture_shm_dir,
"--vaapi-device", cfg.vaapi_device,
]
log.info("starting shared capture daemon shm=%s cams=%d",
cfg.capture_shm_dir, len(payload))
env = dict(os.environ)
env["GST_REGISTRY_UPDATE"] = "no"
env.setdefault("LIBVA_DRM_DEVICE", cfg.vaapi_device)
env.setdefault("LIBVA_DRIVER_NAME", "iHD")
proc = subprocess.Popen(cmd, env=env)
ready = str(Path(cfg.capture_shm_dir) / "ready")
if not wait_for_ready(ready, timeout=30):
proc.kill()
raise RuntimeError(f"capture daemon did not become ready at {ready}")
log.info("capture daemon ready pid=%s", proc.pid)
return proc
def spawn_audio_daemon(cfg, cameras) -> subprocess.Popen:
from frame_shm import wait_for_ready
payload = [{"identity": c.identity, "audio_card": c.audio_card}
for c in cameras if c.audio_card is not None]
if not payload:
raise RuntimeError("no cameras with audio cards")
cmd = [
sys.executable, str(Path(__file__).parent / "audio_daemon.py"),
"--cameras-json", json.dumps(payload),
"--rate", str(cfg.audio_rate),
"--hop-s", str(cfg.enhance_hop_s),
"--shm-dir", cfg.audio_shm_dir,
]
log.info("starting shared audio daemon shm=%s mics=%d",
cfg.audio_shm_dir, len(payload))
proc = subprocess.Popen(cmd)
ready = str(Path(cfg.audio_shm_dir) / "ready")
if not wait_for_ready(ready, timeout=30):
proc.kill()
raise RuntimeError(f"audio daemon did not become ready at {ready}")
log.info("audio daemon ready pid=%s", proc.pid)
return proc
def spawn_encode_daemon(cfg, cameras, portrait_dir: str = "") -> subprocess.Popen:
from frame_shm import wait_for_ready
payload = []
for c in cameras:
payload.append({
"identity": c.identity,
"width": int(c.width or cfg.width),
"height": int(c.height or cfg.height),
})
cmd = [
sys.executable, str(Path(__file__).parent / "encode_daemon.py"),
"--cameras-json", json.dumps(payload),
"--i420-dir", cfg.capture_shm_dir,
"--h264-dir", cfg.encode_shm_dir,
"--fps", str(cfg.fps),
"--bitrate", str(cfg.video_bitrate),
"--hi-bitrate", str(cfg.rally_video_bitrate),
"--vaapi-device", cfg.vaapi_device,
"--codec", cfg.video_codec,
]
if portrait_dir:
cmd.extend(["--portrait-dir", portrait_dir])
log.info("starting encode daemon h264=%s cams=%d codec=%s portrait=%s",
cfg.encode_shm_dir, len(payload), cfg.video_codec,
portrait_dir or "off")
env = dict(os.environ)
env["GST_REGISTRY_UPDATE"] = "no"
env.setdefault("LIBVA_DRM_DEVICE", cfg.vaapi_device)
env.setdefault("LIBVA_DRIVER_NAME", "iHD")
proc = subprocess.Popen(cmd, env=env)
ready = str(Path(cfg.encode_shm_dir) / "ready")
if not wait_for_ready(ready, timeout=60):
proc.kill()
raise RuntimeError(f"encode daemon did not become ready at {ready}")
log.info("encode daemon ready pid=%s", proc.pid)
return proc
def spawn_portrait_daemon(cfg, cameras) -> subprocess.Popen:
from frame_shm import wait_for_ready
payload = []
for c in cameras:
if c.identity == "rally":
continue
payload.append({
"identity": c.identity,
"width": int(c.width or cfg.width),
"height": int(c.height or cfg.height),
})
cmd = [
sys.executable, str(Path(__file__).parent / "portrait_daemon.py"),
"--cameras-json", json.dumps(payload),
"--in-dir", cfg.capture_shm_dir,
"--out-dir", cfg.portrait_shm_dir,
"--hz", str(cfg.portrait_hz),
"--blur-px", str(cfg.portrait_blur_px),
"--hold-s", str(cfg.portrait_hold_s),
]
log.info("starting portrait daemon shm=%s cams=%d hz=%.1f",
cfg.portrait_shm_dir, len(payload), cfg.portrait_hz)
env = dict(os.environ)
env["OMP_NUM_THREADS"] = "1"
env["OPENBLAS_NUM_THREADS"] = "1"
env["MKL_NUM_THREADS"] = "1"
proc = subprocess.Popen(cmd, env=env)
ready = str(Path(cfg.portrait_shm_dir) / "ready")
if not wait_for_ready(ready, timeout=45):
proc.kill()
raise RuntimeError(f"portrait daemon did not become ready at {ready}")
log.info("portrait daemon ready pid=%s", proc.pid)
return proc
def publisher_args_for(cfg, cameras) -> list[str]:
payload = []
for cam in cameras:
d = dict(cam.__dict__)
d["token"] = make_token(
cfg.url, cfg.api_key, cfg.api_secret,
cfg.room, cam.identity,
attributes=token_attributes(cfg.participant_tags))
payload.append(d)
args = [
"--cameras-json", json.dumps(payload),
"--url", cfg.url,
"--room", cfg.room,
"--width", str(cfg.width),
"--height", str(cfg.height),
"--fps", str(cfg.fps),
"--bitrate", str(cfg.video_bitrate),
"--video-codec", cfg.video_codec,
"--video-encoder", cfg.video_encoder,
"--vaapi-device", cfg.vaapi_device,
"--audio-rate", str(cfg.audio_rate),
"--enhance-mode", cfg.enhance_mode,
"--enhance-model", cfg.enhance_model,
"--enhance-chunk-s", str(cfg.enhance_chunk_s),
"--enhance-hop-s", str(cfg.enhance_hop_s),
"--enhance-infer", cfg.enhance_infer,
"--beamform-mode", cfg.beamform_mode,
"--torch-num-threads", str(cfg.torch_num_threads),
"--tdoa-every", str(cfg.beamform_tdoa_every),
"--video-shm-dir", cfg.capture_shm_dir,
"--audio-shm-dir", cfg.audio_shm_dir,
"--h264-shm-dir", cfg.encode_shm_dir,
"--speaker-camera", cfg.display_speaker_camera,
"--audio-gate-open-db", str(cfg.audio_gate_open_db),
"--audio-gate-close-db", str(cfg.audio_gate_close_db),
"--audio-gate-hold-s", str(cfg.audio_gate_hold_s),
]
if cfg.audio_queue_ms is not None:
args.extend(["--audio-queue-ms", str(cfg.audio_queue_ms)])
if cfg.enhance_daemon and cfg.enhance_socket:
args.extend(["--enhance-socket", cfg.enhance_socket])
if cfg.enhance_speaker_only:
args.append("--enhance-speaker-only")
if cfg.audio_gate:
args.append("--audio-gate")
return args
def fleet_cameras(cameras: list[Camera]) -> list[Camera]:
return [c for c in cameras if c.identity != "rally"]
def spawn_rally_follow(cfg, cam: Camera) -> subprocess.Popen:
cmd = [
sys.executable, str(Path(__file__).parent / "rally_follow.py"),
"--device", cam.video_device,
"--width", str(cam.width or 1920),
"--height", str(cam.height or 1080),
"--fps", "30",
"--shm-dir", cfg.capture_shm_dir,
"--audio-shm-dir", cfg.audio_shm_dir,
"--identity", cam.identity,
"--audio-rate", str(cfg.audio_rate),
"--hop-s", str(cfg.enhance_hop_s),
]
if cam.audio_card is not None:
cmd.extend(["--audio-card", str(cam.audio_card)])
log.info("starting rally follow video=%s audio=%s (preview=LiveKit hairpin)",
cam.video_device,
f"hw:{cam.audio_card},0" if cam.audio_card is not None else "none")
return subprocess.Popen(cmd)
def spawn_publisher(cfg, cameras) -> subprocess.Popen:
args = publisher_args_for(cfg, cameras)
env = dict(os.environ)
env.setdefault("LIBVA_DRIVER_NAME", "iHD")
env.setdefault("LIBVA_DRM_DEVICE", cfg.vaapi_device)
log.info("starting single publisher cams=%d", len(cameras))
return subprocess.Popen(
[sys.executable, str(Path(__file__).parent / "run_publisher.py"), *args],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
text=True, bufsize=1, env=env,
)
def stop_proc(proc: subprocess.Popen | None) -> None:
if proc is None or proc.poll() is not None:
return
proc.send_signal(signal.SIGTERM)
with contextlib.suppress(Exception):
proc.wait(timeout=10)
if proc.poll() is None:
proc.kill()
def worker_args_for(cfg, cam) -> list[str]:
enhance_mode = effective_enhance_mode(
cfg.enhance_mode,
cam.identity,
cfg.enhance_speaker_only,
cfg.display_speaker_camera,
)
beamform_mode = cfg.beamform_mode
if cfg.enhance_speaker_only and enhance_mode == "off":
beamform_mode = "off"
args = [
"--camera-json", json.dumps(cam.__dict__),
"--room", cfg.room,
"--url", cfg.url,
"--token", make_token(
cfg.url, cfg.api_key, cfg.api_secret,
cfg.room, cam.identity,
attributes=token_attributes(cfg.participant_tags)),
"--width", str(cfg.width),
"--height", str(cfg.height),
"--fps", str(cfg.fps),
"--bitrate", str(cfg.video_bitrate),
"--video-codec", cfg.video_codec,
"--video-encoder", cfg.video_encoder,
"--vaapi-device", cfg.vaapi_device,
"--audio-rate", str(cfg.audio_rate),
"--enhance-mode", enhance_mode,
"--enhance-model", cfg.enhance_model,
"--enhance-chunk-s", str(cfg.enhance_chunk_s),
"--enhance-hop-s", str(cfg.enhance_hop_s),
"--enhance-infer", cfg.enhance_infer,
"--beamform-mode", beamform_mode,
"--torch-num-threads", str(cfg.torch_num_threads),
"--tdoa-every", str(cfg.beamform_tdoa_every),
"--speaker-camera", cfg.display_speaker_camera,
"--audio-gate-open-db", str(cfg.audio_gate_open_db),
"--audio-gate-close-db", str(cfg.audio_gate_close_db),
"--audio-gate-hold-s", str(cfg.audio_gate_hold_s),
]
if cfg.audio_queue_ms is not None:
args.extend(["--audio-queue-ms", str(cfg.audio_queue_ms)])
if cfg.video_hold_s is not None:
args.extend(["--video-hold-s", str(cfg.video_hold_s)])
if cfg.worker_split_executor:
args.append("--split-executor")
if cfg.enhance_daemon and cfg.enhance_socket:
args.extend(["--enhance-socket", cfg.enhance_socket])
if cfg.audio_gate:
args.append("--audio-gate")
if cfg.capture_daemon:
args.extend([
"--video-shm",
f"{cfg.capture_shm_dir.rstrip('/')}/{cam.identity}.i420",
])
return args
class Supervisor:
def __init__(self, cfg, cameras):
self.cfg = cfg
self.cameras = cameras
self.procs: dict[str, subprocess.Popen] = {}
self.last_exit: dict[str, int | None] = {}
self.restarts: dict[str, int] = {}
def spawn(self, cam) -> None:
args = worker_args_for(self.cfg, cam)
env = dict(os.environ)
env.setdefault("LIBVA_DRIVER_NAME", "iHD")
env.setdefault("LIBVA_DRM_DEVICE", self.cfg.vaapi_device)
p = subprocess.Popen(
[sys.executable, str(Path(__file__).parent / "run_worker.py"),
*args],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
text=True, bufsize=1, env=env,
)
self.procs[cam.identity] = p
self.restarts[cam.identity] = self.restarts.get(cam.identity, 0) + 1
def pump(self) -> None:
"""Read one line from each running worker; print them."""
for ident, p in list(self.procs.items()):
line = p.stdout.readline() if p.stdout else ""
if line:
sys.stdout.write(f"[{ident}] {line}")
sys.stdout.flush()
if p.poll() is not None and not line:
pass # EOF handled by poll below
def reap(self) -> None:
for ident, p in list(self.procs.items()):
rc = p.poll()
if rc is None:
continue
self.last_exit[ident] = rc
if rc != 0:
log.warning("%s exited rc=%s (restart %d/%d)", ident, rc,
self.restarts.get(ident, 0), 10)
if self.restarts.get(ident, 0) < 10:
time.sleep(1.0)
# find the camera object again
cam = next(c for c in self.cameras
if c.identity == ident)
self.spawn(cam)
else:
log.error("%s gave up after 10 restarts", ident)
else:
log.info("%s exited cleanly", ident)
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--cameras", default=None,
help="subset spec, e.g. 'all', '1,5', '0-19'")
ap.add_argument("--once", action="store_true",
help="spawn workers, report status after "
"PUBLISH_TIMEOUT_S, then exit")
args = ap.parse_args()
cfg = load_config()
logging.basicConfig(
level=getattr(logging, cfg.log_level, logging.INFO),
format="%(asctime)s [main] %(levelname)s %(message)s")
if args.cameras:
from config import _parse_cameras
cfg.cameras = _parse_cameras(args.cameras)
if not cfg.url or not cfg.api_key or not cfg.api_secret:
log.error("LIVEKIT_URL / LIVEKIT_API_KEY / LIVEKIT_API_SECRET "
"not set (create .env from .env.example)")
return 2
if not cfg.url.startswith("wss://"):
log.error("LIVEKIT_URL must be wss:// (refusing %s)", cfg.url)
return 2
log.info("livekit destination url=%s room=%s", cfg.url, cfg.room)
cameras = discover_cameras()
log.info("discovered %d cameras:", len(cameras))
for c in cameras:
log.info(" %s video=%s audio=hw:%s usb=%s",
c.identity, c.video_device,
c.audio_card if c.audio_card is not None else "-",
c.usb_port or "-")
cameras = select_cameras(cameras, cfg.cameras)
if cfg.cameras == ["all"]:
rally = discover_rally()
if rally is not None:
cameras = list(cameras) + [rally]
log.info(" %s video=%s audio=hw:%s usb=%s",
rally.identity, rally.video_device,
rally.audio_card if rally.audio_card is not None else "-",
rally.usb_port or "-")
if not cameras:
log.error("no cameras match --cameras %s", args.cameras)
return 2
fleet = fleet_cameras(cameras)
daemon = None
capture = None
audio = None
encode = None
rally_proc = None
portrait = None
if cfg.enhance_daemon and cfg.enhance_mode != "off":
try:
daemon = spawn_enhance_daemon(cfg)
except Exception as exc:
log.error("enhance daemon failed: %s", exc)
return 2
if cfg.capture_daemon:
try:
capture = spawn_capture_daemon(cfg, fleet)
except Exception as exc:
log.error("capture daemon failed: %s", exc)
stop_proc(daemon)
return 2
if cfg.audio_daemon:
try:
audio = spawn_audio_daemon(cfg, fleet)
except Exception as exc:
log.error("audio daemon failed: %s", exc)
stop_proc(capture)
stop_proc(daemon)
return 2
rally_cam = next((c for c in cameras if c.identity == "rally"), None)
if rally_cam is not None:
try:
rally_proc = spawn_rally_follow(cfg, rally_cam)
shm = Path(cfg.capture_shm_dir) / "rally.i420"
from frame_shm import wait_for_ready
if not wait_for_ready(str(shm), timeout=20):
raise RuntimeError(f"rally follow did not create {shm}")
except Exception as exc:
log.error("rally follow failed: %s", exc)
stop_proc(rally_proc)
stop_proc(audio)
stop_proc(capture)
stop_proc(daemon)
return 2
portrait_dir = ""
if cfg.portrait_blur and cfg.capture_daemon:
try:
portrait = spawn_portrait_daemon(cfg, fleet)
portrait_dir = cfg.portrait_shm_dir
except Exception as exc:
log.error("portrait daemon failed, encode uses raw: %s", exc)
stop_proc(portrait)
portrait = None
portrait_dir = ""
if cfg.encode_daemon and cfg.capture_daemon:
try:
encode = spawn_encode_daemon(cfg, cameras, portrait_dir=portrait_dir)
except Exception as exc:
log.error("encode daemon failed: %s", exc)
stop_proc(portrait)
stop_proc(rally_proc)
stop_proc(audio)
stop_proc(capture)
stop_proc(daemon)
return 2
use_pub = cfg.single_publisher and cfg.capture_daemon
sup = Supervisor(cfg, cameras)
pub = None
if use_pub:
pub = spawn_publisher(cfg, cameras)
log.info("spawned single publisher for %d cameras", len(cameras))
else:
for cam in cameras:
sup.spawn(cam)
log.info("spawned %d workers", len(cameras))
def _shutdown():
if pub is not None and pub.poll() is None:
pub.send_signal(signal.SIGTERM)
with contextlib.suppress(Exception):
pub.wait(timeout=10)
if pub.poll() is None:
pub.kill()
for p in sup.procs.values():
if p.poll() is None:
p.send_signal(signal.SIGTERM)
for p in sup.procs.values():
with contextlib.suppress(Exception):
p.wait(timeout=10)
for p in sup.procs.values():
if p.poll() is None:
p.kill()
stop_proc(daemon)
stop_proc(capture)
stop_proc(audio)
stop_proc(encode)
stop_proc(portrait)
stop_proc(rally_proc)
def _pump_pub():
nonlocal pub
if pub is None:
return
line = pub.stdout.readline() if pub.stdout else ""
if line:
sys.stdout.write(line)
sys.stdout.flush()
rc = pub.poll()
if rc is not None:
log.warning("publisher exited rc=%s; restarting", rc)
time.sleep(1.0)
pub = spawn_publisher(cfg, cameras)
if args.once:
deadline = time.monotonic() + cfg.publish_timeout_s
while time.monotonic() < deadline:
if use_pub:
_pump_pub()
else:
sup.pump()
sup.reap()
time.sleep(0.5)
if use_pub:
alive = pub is not None and pub.poll() is None
log.info("STATUS: publisher %s", "UP" if alive else "DOWN")
_shutdown()
return 0 if alive else 1
up = [i for i, p in sup.procs.items() if p.poll() is None]
log.info("STATUS: %d/%d workers alive", len(up), len(sup.procs))
_shutdown()
return 0 if len(up) == len(sup.procs) else 1
try:
while True:
if use_pub:
_pump_pub()
else:
sup.pump()
sup.reap()
time.sleep(0.2)
except KeyboardInterrupt:
log.info("shutting down ...")
_shutdown()
return 0
if __name__ == "__main__":
sys.exit(main())