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.
749 lines
29 KiB
Python
749 lines
29 KiB
Python
#!/usr/bin/env python3
|
|
"""Headless 3-display LiveKit wall via GStreamer kmssink.
|
|
|
|
1. grid — 5x4 mosaic of cam-01..cam-20
|
|
2. speaker — designated speaker camera (DISPLAY_SPEAKER_CAMERA)
|
|
3. screenshare — STUB until the production LiveKit server is online
|
|
|
|
.venv/bin/python run_displays.py --list
|
|
.venv/bin/python run_displays.py --test-pattern
|
|
.venv/bin/python run_displays.py
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import asyncio
|
|
import logging
|
|
import os
|
|
import signal
|
|
import sys
|
|
import time
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from typing import Optional
|
|
from urllib.parse import urlparse
|
|
|
|
import numpy as np
|
|
|
|
from config import load_config
|
|
from display_grid import (
|
|
CameraGrid,
|
|
ROLES,
|
|
active_speaker_ids,
|
|
bgra_from_livekit,
|
|
camera_ids,
|
|
collapse_bleed,
|
|
confirm_speakers,
|
|
describe_livekit_error,
|
|
format_livekit_status,
|
|
hold_speakers,
|
|
placeholder,
|
|
select_grid_speakers,
|
|
update_noise_floor,
|
|
)
|
|
from drm_outputs import DrmOutput, format_outputs, list_outputs, select_outputs
|
|
from gst_sink import KmsPipeline, VaGridPipeline
|
|
from hand_raise import HandRaiseMonitor
|
|
from latency import video_stream_kwargs
|
|
from participant_tags import want_display_video
|
|
from tokens import make_viewer_token
|
|
|
|
log = logging.getLogger("cameras.display")
|
|
|
|
H264_DIR = "/run/livekit-cameras/h264"
|
|
|
|
|
|
def _h264_fifo_paths(identities: list[str]) -> list[str] | None:
|
|
# Sidecar H.264 encode removed: grid tiles come from LiveKit I420.
|
|
del identities
|
|
return None
|
|
|
|
_PATTERNS = ("smpte", "ball", "snow")
|
|
_PLACE_W, _PLACE_H = 1920, 1080
|
|
|
|
|
|
def _public_url(cfg) -> str:
|
|
if cfg.livekit_public_url:
|
|
return cfg.livekit_public_url.rstrip("/")
|
|
parsed = urlparse(cfg.url)
|
|
host = parsed.hostname or "127.0.0.1"
|
|
if host in ("127.0.0.1", "localhost"):
|
|
host = "10.200.0.20"
|
|
scheme = parsed.scheme or "ws"
|
|
port = parsed.port
|
|
if port:
|
|
return f"{scheme}://{host}:{port}"
|
|
return f"{scheme}://{host}"
|
|
|
|
|
|
def _bind_roles(roles: list[str], connectors: list[str]) -> list[tuple[str, DrmOutput]]:
|
|
roles = [r.strip().lower() for r in roles if r.strip()]
|
|
for r in roles:
|
|
if r not in ROLES:
|
|
raise KeyError(f"unknown display role {r!r}; expected {ROLES}")
|
|
names = connectors if connectors else None
|
|
outs = select_outputs(n=len(roles), names=names)
|
|
if len(outs) < len(roles):
|
|
log.warning("wanted %d displays, found %d DRM connectors", len(roles), len(outs))
|
|
slots = list(zip(roles, outs))
|
|
for role, out in slots:
|
|
log.info("role %s -> %s (%s connector-id=%s %s)",
|
|
role, out.name, out.status, out.connector_id, out.driver)
|
|
return slots
|
|
|
|
|
|
class HopLevels:
|
|
"""Latest hop RMS + idle floor per grid identity (does not mute publish)."""
|
|
|
|
def __init__(self, shm_dir: str, identities: list[str], rise_db: float = 6.0) -> None:
|
|
from hop_shm import HopReader
|
|
from audio_gate import hop_mono_int16, rms_dbfs
|
|
|
|
self._mono = hop_mono_int16
|
|
self._rms = rms_dbfs
|
|
self.rise_db = float(rise_db)
|
|
self._db = {i: -120.0 for i in identities}
|
|
self._wave: dict[str, np.ndarray | None] = {i: None for i in identities}
|
|
self._floor: dict[str, float | None] = {i: None for i in identities}
|
|
self._n = {i: 0 for i in identities}
|
|
self._readers = {
|
|
i: HopReader(os.path.join(shm_dir, f"{i}.pcm")) for i in identities
|
|
}
|
|
|
|
def poll(self, speaking: set[str] | None = None) -> None:
|
|
hot = speaking or set()
|
|
for ident, r in self._readers.items():
|
|
blob = None
|
|
while True:
|
|
nxt = r.read()
|
|
if nxt is None:
|
|
break
|
|
blob = nxt
|
|
if blob is None:
|
|
continue
|
|
pcm = np.frombuffer(blob, dtype=np.int16)
|
|
if pcm.size % 2 == 0:
|
|
pcm = pcm.reshape(-1, 2)
|
|
db = self._rms(self._mono(pcm))
|
|
self._db[ident] = db
|
|
self._wave[ident] = np.asarray(self._mono(pcm), dtype=np.float32)
|
|
fl = self._floor[ident]
|
|
self._n[ident] += 1
|
|
if fl is None:
|
|
self._floor[ident] = db
|
|
else:
|
|
self._floor[ident] = update_noise_floor(
|
|
fl, db, speaking=ident in hot, rise_db=self.rise_db)
|
|
|
|
def levels(self, identities: list[str]) -> list[tuple[str, float]]:
|
|
return [
|
|
(i, self._db.get(i, -120.0))
|
|
for i in identities
|
|
if i in self._readers and self._n.get(i, 0) >= 3
|
|
]
|
|
|
|
def waves(self, identities: list[str]) -> dict[str, np.ndarray]:
|
|
out: dict[str, np.ndarray] = {}
|
|
for i in identities:
|
|
w = self._wave.get(i)
|
|
if w is not None:
|
|
out[i] = w
|
|
return out
|
|
|
|
@property
|
|
def floors(self) -> dict[str, float]:
|
|
return {i: fl for i, fl in self._floor.items() if fl is not None}
|
|
|
|
|
|
class Wall:
|
|
def __init__(self, cfg, slots: list[tuple[str, DrmOutput]]):
|
|
self.cfg = cfg
|
|
self.slots = {role: KmsPipeline(out, role) for role, out in slots}
|
|
prefix = cfg.participant_prefix or "cam"
|
|
n = cfg.display_grid_cols * cfg.display_grid_rows
|
|
self.grid = CameraGrid(
|
|
width=cfg.display_grid_width,
|
|
height=cfg.display_grid_height,
|
|
cols=cfg.display_grid_cols,
|
|
rows=cfg.display_grid_rows,
|
|
identities=[],
|
|
)
|
|
self.grid.set_status(format_livekit_status(url=cfg.url, state="connecting"))
|
|
self._retry_s = 5.0
|
|
# LiveKit FFI retries ~3 times (~12s) on ENETUNREACH; stay above that
|
|
# so the overlay can show "host unreachable" instead of a bare timeout.
|
|
self._connect_timeout_s = 20.0
|
|
self.speaker_camera = (cfg.display_speaker_camera or "rally").strip()
|
|
self.active_speaker: Optional[str] = None
|
|
self._last_speakers: tuple[str, ...] = ()
|
|
self._speaker_frame: Optional[np.ndarray] = None
|
|
self._tasks: dict[str, asyncio.Task] = {}
|
|
self._stop = asyncio.Event()
|
|
self._pump_pool = ThreadPoolExecutor(
|
|
max_workers=max(1, int(cfg.display_pump_workers)),
|
|
thread_name_prefix="disp-pump",
|
|
)
|
|
# Screenshare is a stub until production LiveKit is online.
|
|
self._share_placeholder = placeholder(_PLACE_W, _PLACE_H, [
|
|
"Screen share",
|
|
"stub — waiting for production LiveKit",
|
|
_public_url(cfg),
|
|
])
|
|
cam = self.speaker_camera
|
|
self._speaker_placeholder = placeholder(_PLACE_W, _PLACE_H, [
|
|
"Speaker camera",
|
|
"Logitech Rally",
|
|
cam if cam.lower() != "active" else "(following active speaker)",
|
|
"not connected",
|
|
])
|
|
self.va_grid: Optional[VaGridPipeline] = None
|
|
self.hand_raise = HandRaiseMonitor(
|
|
enabled=cfg.display_hand_raise,
|
|
hz=cfg.display_hand_raise_hz,
|
|
hold_s=cfg.display_hand_raise_hold_s,
|
|
on_change=self.grid.set_highlight,
|
|
model=cfg.display_hand_raise_model,
|
|
)
|
|
self._hop_levels = HopLevels(
|
|
cfg.audio_shm_dir, self.grid.identities,
|
|
rise_db=cfg.display_speak_rise_db)
|
|
self._speak_cur: list[str] = []
|
|
self._speak_peak = -120.0
|
|
self._speak_hold_until = 0.0
|
|
self._speak_counts: dict[str, int] = {}
|
|
|
|
def _refresh_speakers(self) -> None:
|
|
"""Hop RMS is the detector; LiveKit is a hint, not a gate."""
|
|
tiles = list(self.grid.identities)
|
|
ranked = self._hop_levels.levels(tiles)
|
|
picked = select_grid_speakers(
|
|
ranked,
|
|
min_db=self.cfg.display_speak_min_db,
|
|
margin_db=self.cfg.display_speak_margin_db,
|
|
floors=self._hop_levels.floors,
|
|
rise_db=self.cfg.display_speak_rise_db,
|
|
)
|
|
picked = collapse_bleed(
|
|
picked,
|
|
self._hop_levels.waves(picked),
|
|
dict(ranked),
|
|
corr_min=self.cfg.display_speak_corr,
|
|
)
|
|
picked, self._speak_counts = confirm_speakers(
|
|
self._speak_counts, picked,
|
|
need=self.cfg.display_speak_confirm,
|
|
already=self._speak_cur,
|
|
)
|
|
new_peak = max((db for i, db in ranked if i in picked), default=-120.0)
|
|
now = time.monotonic()
|
|
chrome, self._speak_hold_until = hold_speakers(
|
|
self._speak_cur, picked,
|
|
prev_peak=self._speak_peak, new_peak=new_peak,
|
|
now=now, hold_until=self._speak_hold_until,
|
|
hold_s=self.cfg.display_speak_hold_s,
|
|
)
|
|
if chrome == picked:
|
|
self._speak_peak = new_peak
|
|
self._speak_cur = chrome
|
|
self.grid.set_speakers(chrome)
|
|
key = tuple(chrome)
|
|
if key != self._last_speakers:
|
|
self._last_speakers = key
|
|
log.info("active-speakers %s", ",".join(chrome) if chrome else "-")
|
|
|
|
def speaker_identity(self) -> Optional[str]:
|
|
if self.speaker_camera.lower() == "active":
|
|
return self.active_speaker
|
|
return self.speaker_camera
|
|
|
|
def _set_lk_status(self, state: str, detail: str = "", retry_s: float | None = None) -> None:
|
|
lines = format_livekit_status(
|
|
url=self.cfg.url, state=state, detail=detail, retry_s=retry_s,
|
|
room=self.cfg.room)
|
|
if state == "connected":
|
|
if self.grid.identities:
|
|
self.grid.set_status(None)
|
|
else:
|
|
self.grid.set_status(["NO LIVE CAMERAS", f"room {self.cfg.room}"])
|
|
return
|
|
self.grid.set_status(lines)
|
|
|
|
def _seat_participant(self, participant) -> None:
|
|
ident = getattr(participant, "identity", "") or ""
|
|
if not ident or not self._want_video(participant):
|
|
return
|
|
if ident == self.speaker_identity():
|
|
return
|
|
if not self.grid.add_identity(ident):
|
|
log.warning("grid full, skip %s", ident)
|
|
return
|
|
if self.grid.identities:
|
|
self.grid.set_status(None)
|
|
|
|
def _unseat_participant(self, identity: str) -> None:
|
|
ident = (identity or "").strip()
|
|
if not ident:
|
|
return
|
|
self.grid.remove_identity(ident)
|
|
self.grid.clear(ident)
|
|
if ident == self.speaker_identity():
|
|
self._speaker_frame = None
|
|
# Keep mosaic if other tiles remain; otherwise show empty-room text
|
|
# only while the SFU session is still up (status overlay owns errors).
|
|
if not self.grid.identities and not self.grid._status_lines:
|
|
self.grid.set_status(["NO LIVE CAMERAS", f"room {self.cfg.room}"])
|
|
|
|
def _tile_index(self, identity: str) -> Optional[int]:
|
|
try:
|
|
return self.grid.identities.index(identity)
|
|
except ValueError:
|
|
return None
|
|
|
|
def _start_va_grid(self) -> None:
|
|
# vacompositor prerolls one HDMI frame then never flips (crtc CRC unique=1).
|
|
# CPU mosaic through appsrc+kmssink is the path that actually presents.
|
|
if os.environ.get("DISPLAY_VA_GRID", "0").strip() != "1":
|
|
log.info("va-grid skipped; CPU mosaic -> kmssink")
|
|
return
|
|
sink = self.slots.get("grid")
|
|
if sink is None or not sink.out.connected:
|
|
return
|
|
self.va_grid = VaGridPipeline(sink.out)
|
|
self.va_grid.start(
|
|
self.grid,
|
|
ingest_w=max(16, int(self.cfg.width)),
|
|
ingest_h=max(16, int(self.cfg.height)),
|
|
fps=max(1, int(self.cfg.fps) or 15),
|
|
queue_buffers=self.cfg.display_gst_queue_buffers,
|
|
h264_paths=_h264_fifo_paths(self.grid.identities),
|
|
)
|
|
if not self.va_grid.alive:
|
|
log.warning("va-grid failed to start; falling back to CPU mosaic")
|
|
self.va_grid = None
|
|
|
|
def _want_video(self, participant) -> bool:
|
|
return want_display_video(
|
|
identity=getattr(participant, "identity", "") or "",
|
|
attributes=getattr(participant, "attributes", None) or {},
|
|
metadata=getattr(participant, "metadata", "") or "",
|
|
hide_tags=self.cfg.display_hide_tags,
|
|
participant_prefix=self.cfg.participant_prefix or "cam",
|
|
display_identity=self.cfg.display_identity,
|
|
speaker_identity=self.speaker_identity() or "",
|
|
)
|
|
|
|
def _push_sink(self, role: str, data: bytes, width: int, height: int, fps: int,
|
|
pixel_format: str = "bgra") -> None:
|
|
sink = self.slots.get(role)
|
|
if sink is None or not sink.out.connected:
|
|
return
|
|
if (not sink.alive or sink.width != width or sink.height != height
|
|
or getattr(sink, "pixel_format", "bgra") != pixel_format):
|
|
sink.start_appsrc(
|
|
width, height, fps,
|
|
queue_buffers=self.cfg.display_gst_queue_buffers,
|
|
pixel_format=pixel_format)
|
|
if not sink.push(data):
|
|
sink.start_appsrc(
|
|
width, height, fps,
|
|
queue_buffers=self.cfg.display_gst_queue_buffers,
|
|
pixel_format=pixel_format)
|
|
sink.push(data)
|
|
err = sink.gst_errors()
|
|
if err:
|
|
for line in err.strip().splitlines()[-6:]:
|
|
log.warning("%s gst: %s", role, line)
|
|
|
|
async def _output_loops(self) -> None:
|
|
self._start_va_grid()
|
|
grid_fps = max(1, int(self.cfg.display_grid_fps))
|
|
grid_period = 1.0 / grid_fps
|
|
spk_period = 1.0 / max(1, int(self.cfg.fps))
|
|
next_grid = time.monotonic()
|
|
next_spk = time.monotonic()
|
|
next_hop = time.monotonic()
|
|
while not self._stop.is_set():
|
|
now = time.monotonic()
|
|
if now >= next_hop:
|
|
self._hop_levels.poll(speaking=set(self._speak_cur))
|
|
self._refresh_speakers()
|
|
next_hop = now + 0.2
|
|
if now >= next_grid:
|
|
if self.va_grid is not None:
|
|
if not self.va_grid.alive:
|
|
log.warning("va-grid died; falling back to CPU mosaic")
|
|
self.va_grid.stop()
|
|
self.va_grid = None
|
|
else:
|
|
err = self.va_grid.gst_errors()
|
|
if err:
|
|
for line in err.strip().splitlines()[-6:]:
|
|
log.warning("va-grid gst: %s", line)
|
|
if self.va_grid is None:
|
|
i420 = self.cfg.display_video_format in ("i420", "yuv420p")
|
|
if i420:
|
|
self._push_sink(
|
|
"grid", self.grid.render_i420(),
|
|
self.grid.width, self.grid.height, grid_fps,
|
|
pixel_format="i420",
|
|
)
|
|
else:
|
|
self._push_sink(
|
|
"grid", self.grid.render(),
|
|
self.grid.width, self.grid.height, grid_fps,
|
|
)
|
|
next_grid = now + grid_period
|
|
if now >= next_spk:
|
|
spk = self._speaker_frame
|
|
arrival = self.cfg.display_speaker_push == "arrival"
|
|
if spk is None:
|
|
img = self._speaker_placeholder
|
|
self._push_sink(
|
|
"speaker", img.tobytes(), img.shape[1], img.shape[0],
|
|
self.cfg.fps)
|
|
elif not arrival:
|
|
self._push_sink(
|
|
"speaker", spk.tobytes(), spk.shape[1], spk.shape[0],
|
|
self.cfg.fps)
|
|
share = self._share_placeholder
|
|
self._push_sink("screenshare", share.tobytes(), share.shape[1], share.shape[0], self.cfg.fps)
|
|
next_spk = now + spk_period
|
|
await asyncio.sleep(0.005)
|
|
|
|
async def _pump_camera(self, track, identity: str) -> None:
|
|
import livekit.rtc as rtc
|
|
|
|
stream_kw = video_stream_kwargs(
|
|
self.cfg.display_video_capacity, self.cfg.display_video_format)
|
|
if stream_kw.get("format") == "bgra":
|
|
stream_kw["format"] = rtc.VideoBufferType.BGRA
|
|
elif stream_kw.get("format") in ("i420", "yuv420p"):
|
|
stream_kw["format"] = rtc.VideoBufferType.I420
|
|
stream = rtc.VideoStream(track, **stream_kw)
|
|
want_bgra = self.cfg.display_video_format == "bgra"
|
|
use_i420 = self.cfg.display_video_format in ("i420", "yuv420p")
|
|
loop = asyncio.get_running_loop()
|
|
try:
|
|
async for event in stream:
|
|
frame = event.frame
|
|
if frame.width <= 0 or frame.height <= 0:
|
|
continue
|
|
if (self.va_grid is not None
|
|
and getattr(self.va_grid, "codec", "") == "i420"):
|
|
idx = self._tile_index(identity)
|
|
if (idx is not None
|
|
and frame.width == self.va_grid.ingest_w
|
|
and frame.height == self.va_grid.ingest_h):
|
|
self.va_grid.push_tile(idx, frame.data)
|
|
is_spk = identity == self.speaker_identity()
|
|
cpu_grid = self.va_grid is None
|
|
if cpu_grid and use_i420 and not is_spk:
|
|
blob = bytes(frame.data)
|
|
await loop.run_in_executor(
|
|
self._pump_pool,
|
|
self.grid.set_frame_i420,
|
|
identity, blob, frame.width, frame.height)
|
|
self.hand_raise.offer(identity, blob, frame.width, frame.height)
|
|
continue
|
|
if cpu_grid or is_spk:
|
|
try:
|
|
if frame.type == rtc.VideoBufferType.BGRA:
|
|
converted = frame
|
|
else:
|
|
converted = frame.convert(rtc.VideoBufferType.BGRA)
|
|
except Exception as exc:
|
|
log.warning("%s: convert failed: %s", identity, exc)
|
|
continue
|
|
bgra = bgra_from_livekit(
|
|
converted.data, converted.width, converted.height)
|
|
if cpu_grid:
|
|
self.grid.set_frame(identity, bgra)
|
|
if is_spk:
|
|
self._speaker_frame = bgra
|
|
if self.cfg.display_speaker_push == "arrival":
|
|
self._push_sink(
|
|
"speaker", bgra.tobytes(),
|
|
bgra.shape[1], bgra.shape[0], self.cfg.fps)
|
|
finally:
|
|
await stream.aclose()
|
|
self.hand_raise.forget(identity)
|
|
self.grid.clear(identity)
|
|
if identity == self.speaker_identity():
|
|
self._speaker_frame = None
|
|
|
|
def _track_key(self, participant_identity: str, track) -> str:
|
|
return f"{participant_identity}:{getattr(track, 'sid', id(track))}"
|
|
|
|
def attach(self, track, publication, participant, rtc) -> None:
|
|
ident = participant.identity
|
|
if not self._want_video(participant):
|
|
return
|
|
is_spk = ident == self.speaker_identity()
|
|
if not is_spk:
|
|
self._seat_participant(participant)
|
|
if ident not in self.grid.identities:
|
|
return
|
|
key = self._track_key(ident, track)
|
|
old = self._tasks.pop(key, None)
|
|
if old:
|
|
old.cancel()
|
|
# Screenshare subscribe/render is stubbed until production LiveKit.
|
|
self._tasks[key] = asyncio.create_task(
|
|
self._pump_camera(track, ident), name=f"cam-{ident}")
|
|
|
|
def subscribe_if_wanted(self, participant, rtc) -> None:
|
|
ident = participant.identity
|
|
if not self._want_video(participant):
|
|
return
|
|
for pub in participant.track_publications.values():
|
|
if pub.kind != rtc.TrackKind.KIND_VIDEO:
|
|
continue
|
|
if not pub.subscribed:
|
|
log.info("subscribe %s from %s (source=%s)", pub.sid, ident, pub.source)
|
|
pub.set_subscribed(True)
|
|
|
|
def detach(self, track, participant) -> None:
|
|
key = self._track_key(participant.identity, track)
|
|
task = self._tasks.pop(key, None)
|
|
if task:
|
|
task.cancel()
|
|
|
|
async def run(self) -> int:
|
|
loops = asyncio.create_task(self._output_loops(), name="output-loops")
|
|
loop = asyncio.get_running_loop()
|
|
for sig in (signal.SIGINT, signal.SIGTERM):
|
|
try:
|
|
loop.add_signal_handler(sig, self._stop.set)
|
|
except NotImplementedError:
|
|
pass
|
|
self.hand_raise.start()
|
|
self._set_lk_status("connecting")
|
|
log.info(
|
|
"wall up (livekit optional): grid=%dx%d speaker_camera=%s "
|
|
"video_capacity=%d format=%s gst_queue=%d speaker_push=%s "
|
|
"pump_workers=%d hide_tags=%s hand_raise=%s url=%s",
|
|
self.cfg.display_grid_cols, self.cfg.display_grid_rows, self.speaker_camera,
|
|
self.cfg.display_video_capacity,
|
|
self.cfg.display_video_format or "native",
|
|
self.cfg.display_gst_queue_buffers,
|
|
self.cfg.display_speaker_push,
|
|
self.cfg.display_pump_workers,
|
|
",".join(self.cfg.display_hide_tags) or "off",
|
|
"on" if self.hand_raise.enabled else "off",
|
|
self.cfg.url,
|
|
)
|
|
try:
|
|
while not self._stop.is_set():
|
|
try:
|
|
await self._session()
|
|
except asyncio.CancelledError:
|
|
raise
|
|
except Exception as exc:
|
|
detail = describe_livekit_error(exc)
|
|
log.error("livekit %s: %s", detail, exc or type(exc).__name__)
|
|
self.grid.set_identities([])
|
|
self._speaker_frame = None
|
|
self._set_lk_status(
|
|
"disconnected", detail=detail, retry_s=self._retry_s)
|
|
if self._stop.is_set():
|
|
break
|
|
try:
|
|
await asyncio.wait_for(self._stop.wait(), timeout=self._retry_s)
|
|
except asyncio.TimeoutError:
|
|
pass
|
|
if not self._stop.is_set():
|
|
self._set_lk_status("connecting")
|
|
finally:
|
|
loops.cancel()
|
|
for t in self._tasks.values():
|
|
t.cancel()
|
|
self._tasks.clear()
|
|
for sink in self.slots.values():
|
|
sink.stop()
|
|
if self.va_grid is not None:
|
|
self.va_grid.stop()
|
|
self.va_grid = None
|
|
self._pump_pool.shutdown(wait=False, cancel_futures=True)
|
|
self.hand_raise.stop()
|
|
return 0
|
|
|
|
async def _session(self) -> None:
|
|
import livekit.rtc as rtc
|
|
|
|
room = rtc.Room()
|
|
session_done = asyncio.Event()
|
|
|
|
@room.on("track_subscribed")
|
|
def _on_sub(track, publication, participant):
|
|
if track.kind == rtc.TrackKind.KIND_VIDEO:
|
|
log.info("subscribed %s from %s", publication.sid, participant.identity)
|
|
self.attach(track, publication, participant, rtc)
|
|
|
|
@room.on("track_unsubscribed")
|
|
def _on_unsub(track, publication, participant):
|
|
self.detach(track, participant)
|
|
|
|
@room.on("track_published")
|
|
def _on_pub(publication, participant):
|
|
self.subscribe_if_wanted(participant, rtc)
|
|
|
|
@room.on("participant_connected")
|
|
def _on_pc(participant):
|
|
log.info("participant joined: %s", participant.identity)
|
|
self._seat_participant(participant)
|
|
self.subscribe_if_wanted(participant, rtc)
|
|
|
|
@room.on("participant_disconnected")
|
|
def _on_pd(participant):
|
|
ident = getattr(participant, "identity", "") or ""
|
|
log.info("participant left: %s", ident)
|
|
self._unseat_participant(ident)
|
|
|
|
@room.on("disconnected")
|
|
def _on_disc(*args):
|
|
reason = ""
|
|
if args:
|
|
reason = str(args[0] or "")
|
|
detail = describe_livekit_error(RuntimeError(reason)) if reason else "server closed the session"
|
|
if reason:
|
|
low = reason.lower()
|
|
if "disconnect" in low or "close" in low:
|
|
detail = reason[:96]
|
|
log.warning("livekit disconnected: %s", reason or "no reason")
|
|
self.grid.set_identities([])
|
|
self._speaker_frame = None
|
|
self._set_lk_status(
|
|
"disconnected", detail=detail, retry_s=self._retry_s)
|
|
session_done.set()
|
|
|
|
@room.on("active_speakers_changed")
|
|
def _on_spk(speakers):
|
|
cams = active_speaker_ids(speakers, want=self._want_video)
|
|
self.active_speaker = cams[0] if cams else None
|
|
self._hop_levels.poll(speaking=set(self._speak_cur))
|
|
self._refresh_speakers()
|
|
|
|
token = make_viewer_token(
|
|
self.cfg.api_key, self.cfg.api_secret, self.cfg.room, self.cfg.display_identity)
|
|
log.info("connecting to %s room=%s as %s", self.cfg.url, self.cfg.room, self.cfg.display_identity)
|
|
self._set_lk_status("connecting")
|
|
try:
|
|
await asyncio.wait_for(
|
|
room.connect(self.cfg.url, token, rtc.RoomOptions(auto_subscribe=False)),
|
|
timeout=self._connect_timeout_s,
|
|
)
|
|
except Exception:
|
|
try:
|
|
await room.disconnect()
|
|
except Exception:
|
|
pass
|
|
raise
|
|
|
|
self._set_lk_status("connected")
|
|
for rp in room.remote_participants.values():
|
|
self._seat_participant(rp)
|
|
self.subscribe_if_wanted(rp, rtc)
|
|
for pub in rp.track_publications.values():
|
|
if pub.track and pub.track.kind == rtc.TrackKind.KIND_VIDEO:
|
|
self.attach(pub.track, pub, rp, rtc)
|
|
log.info(
|
|
"livekit connected: tiles=%d speaker_camera=%s",
|
|
len(self.grid.identities), self.speaker_camera,
|
|
)
|
|
stop_task = asyncio.create_task(self._stop.wait())
|
|
disc_task = asyncio.create_task(session_done.wait())
|
|
try:
|
|
await asyncio.wait({stop_task, disc_task}, return_when=asyncio.FIRST_COMPLETED)
|
|
finally:
|
|
stop_task.cancel()
|
|
disc_task.cancel()
|
|
for t in self._tasks.values():
|
|
t.cancel()
|
|
self._tasks.clear()
|
|
try:
|
|
await room.disconnect()
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def run_test_pattern(slots: list[tuple[str, DrmOutput]]) -> int:
|
|
sinks = []
|
|
for i, (role, out) in enumerate(slots):
|
|
sink = KmsPipeline(out, role)
|
|
sink.start_test_pattern(_PATTERNS[i % len(_PATTERNS)])
|
|
sinks.append(sink)
|
|
log.info("test patterns on %d sink(s); Ctrl-C to stop", sum(1 for s in sinks if s.alive))
|
|
stop = False
|
|
|
|
def _stop(*_):
|
|
nonlocal stop
|
|
stop = True
|
|
|
|
signal.signal(signal.SIGINT, _stop)
|
|
signal.signal(signal.SIGTERM, _stop)
|
|
try:
|
|
while not stop:
|
|
for s in sinks:
|
|
err = s.gst_errors()
|
|
if err:
|
|
for line in err.strip().splitlines()[-8:]:
|
|
log.warning("%s gst: %s", s.identity, line)
|
|
if any(s.out.connected and not s.alive and s.mode == "dead" for s in sinks):
|
|
return 1
|
|
time.sleep(0.5)
|
|
finally:
|
|
for s in sinks:
|
|
s.stop()
|
|
return 0
|
|
|
|
|
|
def main(argv: Optional[list[str]] = None) -> int:
|
|
ap = argparse.ArgumentParser(description=__doc__)
|
|
ap.add_argument("--list", action="store_true", help="list DRM connectors and exit")
|
|
ap.add_argument("--test-pattern", action="store_true",
|
|
help="drive displays with GStreamer videotestsrc (no LiveKit)")
|
|
ap.add_argument("--connectors", default="",
|
|
help="comma-separated DRM names, e.g. HDMI-A-7,HDMI-A-5,HDMI-A-6")
|
|
ap.add_argument("--roles", default="",
|
|
help="comma-separated roles (default grid,speaker,screenshare)")
|
|
ap.add_argument("--speaker-camera", default="",
|
|
help="cam-NN or 'active' (overrides DISPLAY_SPEAKER_CAMERA)")
|
|
args = ap.parse_args(argv)
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [display] %(levelname)s %(message)s",
|
|
)
|
|
|
|
if args.list:
|
|
print(format_outputs(list_outputs()) or "(no DRM connectors)")
|
|
return 0
|
|
|
|
cfg = load_config()
|
|
logging.getLogger().setLevel(getattr(logging, cfg.log_level, logging.INFO))
|
|
if cfg.url or cfg.room:
|
|
log.info("livekit destination url=%s room=%s", cfg.url, cfg.room)
|
|
if args.speaker_camera:
|
|
cfg.display_speaker_camera = args.speaker_camera
|
|
roles = [p.strip() for p in args.roles.split(",") if p.strip()] or cfg.display_roles
|
|
connectors = [p.strip() for p in args.connectors.split(",") if p.strip()] or cfg.display_connectors
|
|
try:
|
|
slots = _bind_roles(roles, connectors)
|
|
except KeyError as exc:
|
|
log.error("%s", exc)
|
|
return 2
|
|
if not slots:
|
|
log.error("no DRM outputs to drive")
|
|
return 2
|
|
|
|
if args.test_pattern:
|
|
return run_test_pattern(slots)
|
|
if not (cfg.url or "").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)
|
|
return asyncio.run(Wall(cfg, slots).run())
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|