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.
590 lines
21 KiB
Python
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())
|