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.
137 lines
4.2 KiB
Python
137 lines
4.2 KiB
Python
"""Pure helpers for env-flagged capture/display latency knobs."""
|
|
from __future__ import annotations
|
|
|
|
|
|
def audio_queue_size_ms(hop_s: float, audio_queue_ms: int | None = None) -> int:
|
|
"""LiveKit AudioSource queue. Unset AUDIO_QUEUE_MS keeps the old 150 ms floor."""
|
|
hop_ms = int(round(float(hop_s) * 1000.0))
|
|
if audio_queue_ms is not None:
|
|
return max(1, int(audio_queue_ms))
|
|
return max(150, hop_ms + 50)
|
|
|
|
|
|
def video_hold_seconds(
|
|
hop_s: float,
|
|
has_audio: bool,
|
|
video_hold_s: float | None = None,
|
|
) -> float:
|
|
"""Delay published video to match audio hops. VIDEO_HOLD_S=0 disables."""
|
|
if not has_audio:
|
|
return 0.0
|
|
if video_hold_s is not None:
|
|
return max(0.0, float(video_hold_s))
|
|
return float(hop_s)
|
|
|
|
|
|
def effective_enhance_mode(
|
|
mode: str,
|
|
identity: str,
|
|
speaker_only: bool,
|
|
speaker_camera: str,
|
|
) -> str:
|
|
"""ENHANCE_SPEAKER_ONLY=1 keeps enhance on DISPLAY_SPEAKER_CAMERA only.
|
|
|
|
``speaker_camera=active`` has no single identity, so every camera keeps
|
|
the configured mode.
|
|
"""
|
|
if not speaker_only:
|
|
return mode
|
|
target = (speaker_camera or "").strip()
|
|
if not target or target.lower() == "active":
|
|
return mode
|
|
if identity != target:
|
|
return "off"
|
|
return mode
|
|
|
|
|
|
def next_video_capture(now: float, due: float, period: float) -> tuple[bool, float]:
|
|
"""Pace VideoSource.capture_frame.
|
|
|
|
When the loop is late, capture once and jump the next due time to
|
|
``now + period`` so we never catch up by pushing extra frames into
|
|
libwebrtc's EncoderQueue (that path leaked ~2.6 MiB/s native RSS).
|
|
"""
|
|
period = float(period)
|
|
if period <= 0:
|
|
period = 1.0 / 30.0
|
|
if now < due:
|
|
return False, due
|
|
return True, now + period
|
|
|
|
|
|
def should_drop_stale_vaapi_frame(
|
|
now_us: int,
|
|
capture_us: int,
|
|
fps: int,
|
|
max_frames_behind: int = 3,
|
|
is_keyframe: bool = False,
|
|
) -> bool:
|
|
"""True when VAAPI Encode should drop before ToI420 (encoder behind).
|
|
|
|
LiveKit's VAAPIH264EncoderWrapper::Encode never dropped; EncoderQueue
|
|
kept I420 copies (~0.37 extra frames/s/cam → OOM). Keyframes always
|
|
pass. Missing timestamps are not treated as stale.
|
|
"""
|
|
if is_keyframe:
|
|
return False
|
|
if int(capture_us) <= 0:
|
|
return False
|
|
fps = max(1, int(fps))
|
|
budget = max(1, int(max_frames_behind))
|
|
max_delay_us = budget * 1_000_000 // fps
|
|
return int(now_us) - int(capture_us) > max_delay_us
|
|
|
|
|
|
def should_drop_encoder_queue_depth(queued: int, max_queued: int = 1) -> bool:
|
|
"""True when VideoTrackSource must not push another frame (latest-wins)."""
|
|
return int(queued) >= max(1, int(max_queued))
|
|
|
|
|
|
def video_capture_gate_acquire(gate: dict) -> bool:
|
|
"""Admit a capture_frame if in-flight slots remain (AudioSource-like).
|
|
|
|
``gate`` is ``{in_flight, dropped, max_in_flight}``. Caller must
|
|
``video_capture_gate_release`` when the encoder slot is free (or after
|
|
one frame period as a stand-in when FFI has no completion callback).
|
|
"""
|
|
max_in_flight = max(1, int(gate.get("max_in_flight") or 3))
|
|
if int(gate.get("in_flight") or 0) >= max_in_flight:
|
|
gate["dropped"] = int(gate.get("dropped") or 0) + 1
|
|
return False
|
|
gate["in_flight"] = int(gate.get("in_flight") or 0) + 1
|
|
return True
|
|
|
|
|
|
def video_capture_gate_release(gate: dict) -> None:
|
|
"""Free one in-flight capture slot."""
|
|
n = int(gate.get("in_flight") or 0)
|
|
gate["in_flight"] = n - 1 if n > 0 else 0
|
|
|
|
|
|
def should_recycle_encoders(
|
|
rss_anon_kb: int,
|
|
limit_kb: int,
|
|
now: float,
|
|
last_recycle_at: float,
|
|
min_interval_s: float = 60.0,
|
|
) -> bool:
|
|
"""True when native RSS is over the cap and a recycle is allowed."""
|
|
if int(limit_kb) <= 0:
|
|
return False
|
|
if int(rss_anon_kb) < int(limit_kb):
|
|
return False
|
|
if last_recycle_at > 0 and (now - last_recycle_at) < float(min_interval_s):
|
|
return False
|
|
return True
|
|
|
|
|
|
def video_stream_kwargs(capacity: int = 0, format_name: str = "") -> dict:
|
|
"""LiveKit VideoStream options. capacity=0 is unbounded (SDK default)."""
|
|
out: dict = {}
|
|
if int(capacity) > 0:
|
|
out["capacity"] = int(capacity)
|
|
fmt = (format_name or "").strip().lower()
|
|
if fmt:
|
|
out["format"] = fmt
|
|
return out
|