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

451 lines
15 KiB
Python

"""GStreamer kmssink pipelines for one DRM connector."""
from __future__ import annotations
import logging
import os
import fcntl
import shutil
import subprocess
from typing import Optional
from drm_outputs import DrmOutput
log = logging.getLogger("cameras.display")
GST_LAUNCH = shutil.which("gst-launch-1.0") or "gst-launch-1.0"
def _kms_props(out: DrmOutput) -> list[str]:
props = [
"kmssink",
f"connector-id={out.connector_id}",
"force-modesetting=true",
"can-scale=true",
"sync=false",
"restore-crtc=true",
]
if out.driver:
props.append(f"driver-name={out.driver}")
return props
def test_pattern_cmd(out: DrmOutput, pattern: str = "smpte", width: int = 1280, height: int = 720) -> list[str]:
return [
GST_LAUNCH, "-e",
"videotestsrc", "is-live=true", f"pattern={pattern}",
"!", f"video/x-raw,width={width},height={height},framerate=30/1",
"!", "videoconvert",
"!", "videoscale",
"!", *_kms_props(out),
]
def appsrc_cmd(
out: DrmOutput,
width: int,
height: int,
fps: int = 30,
queue_buffers: int = 2,
pixel_format: str = "bgra",
) -> list[str]:
"""Raw BGRA or I420 frames on stdin -> kmssink."""
fmt = (pixel_format or "bgra").strip().lower()
if fmt in ("i420", "yuv420p"):
fmt = "i420"
block = width * height * 3 // 2
else:
fmt = "bgra"
block = width * height * 4
q = max(1, int(queue_buffers))
return [
GST_LAUNCH, "-e",
"fdsrc", "fd=0", f"blocksize={block}",
"!", "videoparse",
f"width={width}", f"height={height}", f"format={fmt}", f"framerate={fps}/1",
"!", "queue", f"max-size-buffers={q}", "leaky=downstream",
"!", "videoconvert",
"!", "videoscale", "method=0",
"!", *_kms_props(out),
]
def va_grid_cmd(
out: DrmOutput,
grid,
ingest_w: int,
ingest_h: int,
fps: int,
tile_fds: list[int] | None = None,
chrome_fd: int | None = None,
chrome_path: str | None = None,
queue_buffers: int = 1,
tile_h264_paths: list[str] | None = None,
) -> list[str]:
"""Chrome + N tiles -> Arc vacompositor -> kmssink.
Tiles are I420 fdsrcs (default) or Annex-B H.264 files/fifos decoded
with varenderD129h264dec. Chrome is still a frozen BGRA overlay.
"""
q = max(1, int(queue_buffers))
fps = max(1, int(fps))
if tile_h264_paths:
n = len(tile_h264_paths)
elif tile_fds:
n = len(tile_fds)
else:
raise ValueError("tile_fds or tile_h264_paths required")
cmd: list[str] = [GST_LAUNCH, "-e", "varenderD129compositor", "name=c"]
cmd += [
"sink_0::zorder=0", "sink_0::xpos=0", "sink_0::ypos=0",
f"sink_0::width={grid.width}", f"sink_0::height={grid.height}",
]
for i in range(n):
x, y, w, h = grid.inner_rect(i)
s = i + 1
cmd += [
f"sink_{s}::zorder=1",
f"sink_{s}::xpos={x}",
f"sink_{s}::ypos={y}",
f"sink_{s}::width={w}",
f"sink_{s}::height={h}",
]
cmd += [
"!", f"video/x-raw,width={grid.width},height={grid.height},framerate={fps}/1",
"!", "varenderD129postproc",
"!", f"video/x-raw,format=NV12,width={grid.width},height={grid.height}",
"!", "queue", f"max-size-buffers={q}", "leaky=downstream",
"!", *_kms_props(out),
]
chrome_block = grid.width * grid.height * 4
if chrome_path:
cmd += [
"filesrc", f"location={chrome_path}", f"blocksize={chrome_block}",
"!", "videoparse",
f"width={grid.width}", f"height={grid.height}",
"format=bgra", f"framerate={fps}/1",
"!", "imagefreeze", "is-live=true",
"!", "varenderD129postproc",
"!", "c.sink_0",
]
else:
if chrome_fd is None:
raise ValueError("chrome_fd or chrome_path required")
cmd += [
"fdsrc", f"fd={chrome_fd}", f"blocksize={chrome_block}",
"!", "videoparse",
f"width={grid.width}", f"height={grid.height}",
"format=bgra", f"framerate={fps}/1",
"!", "queue", f"max-size-buffers={q}", "leaky=downstream",
"!", "varenderD129postproc",
"!", "c.sink_0",
]
if tile_h264_paths:
for i, path in enumerate(tile_h264_paths):
s = i + 1
cmd += [
"filesrc", f"location={path}",
"!", "queue", f"max-size-buffers={q}", "leaky=downstream",
"!", "h264parse",
"!", "varenderD129h264dec",
"!", f"c.sink_{s}",
]
return cmd
tile_block = ingest_w * ingest_h * 3 // 2
assert tile_fds is not None
for i, fd in enumerate(tile_fds):
s = i + 1
cmd += [
"fdsrc", f"fd={fd}", f"blocksize={tile_block}",
"!", "videoparse",
f"width={ingest_w}", f"height={ingest_h}",
"format=i420", f"framerate={fps}/1",
"!", "queue", f"max-size-buffers={q}", "leaky=downstream",
"!", "varenderD129postproc", "add-borders=true",
"!", f"c.sink_{s}",
]
return cmd
class KmsPipeline:
"""One gst-launch kmssink process, fed raw BGRA on stdin (or self-running testsrc)."""
def __init__(self, out: DrmOutput, identity: str):
self.out = out
self.identity = identity
self.proc: Optional[subprocess.Popen] = None
self.width = 0
self.height = 0
self.mode = "idle"
self.queue_buffers = 2
self.pixel_format = "bgra"
@property
def alive(self) -> bool:
return self.proc is not None and self.proc.poll() is None
def stop(self) -> None:
if self.proc is None:
return
if self.proc.stdin:
try:
self.proc.stdin.close()
except OSError:
pass
if self.proc.poll() is None:
self.proc.terminate()
try:
self.proc.wait(timeout=3)
except subprocess.TimeoutExpired:
self.proc.kill()
self.proc.wait(timeout=2)
self.proc = None
self.mode = "idle"
def start_test_pattern(self, pattern: str = "smpte") -> None:
self.stop()
if not self.out.connected:
log.warning("%s: %s not connected, skip test pattern", self.identity, self.out.name)
return
cmd = test_pattern_cmd(self.out, pattern=pattern)
log.info("%s: gst test pattern %s on %s id=%s", self.identity, pattern, self.out.name, self.out.connector_id)
self.proc = subprocess.Popen(
cmd,
stdin=subprocess.DEVNULL,
stdout=subprocess.DEVNULL,
stderr=subprocess.PIPE,
env=_gst_env(),
)
self.mode = "pattern"
def start_appsrc(self, width: int, height: int, fps: int = 30,
queue_buffers: int = 2, pixel_format: str = "bgra") -> None:
queue_buffers = max(1, int(queue_buffers))
fmt = (pixel_format or "bgra").strip().lower()
if fmt in ("i420", "yuv420p"):
fmt = "i420"
else:
fmt = "bgra"
if (self.alive and self.mode == "appsrc" and self.width == width
and self.height == height and self.queue_buffers == queue_buffers
and self.pixel_format == fmt):
return
self.stop()
if not self.out.connected:
log.warning("%s: %s not connected, cannot start kmssink", self.identity, self.out.name)
return
cmd = appsrc_cmd(self.out, width, height, fps, queue_buffers=queue_buffers,
pixel_format=fmt)
log.info(
"%s: gst appsrc %dx%d %s -> %s (connector %s %s)",
self.identity, width, height, fmt, self.out.name, self.out.connector_id, self.out.driver,
)
self.proc = subprocess.Popen(
cmd,
stdin=subprocess.PIPE,
stdout=subprocess.DEVNULL,
stderr=subprocess.PIPE,
env=_gst_env(),
bufsize=0,
)
self.width = width
self.height = height
self.mode = "appsrc"
self.queue_buffers = queue_buffers
self.pixel_format = fmt
def push(self, frame: bytes) -> bool:
if not self.alive or self.proc is None or self.proc.stdin is None:
return False
try:
self.proc.stdin.write(frame)
return True
except BrokenPipeError:
log.warning("%s: gst stdin closed (rc=%s)", self.identity, self.proc.poll())
self.proc = None
return False
def gst_errors(self) -> str:
if self.proc is None or self.proc.stderr is None:
return ""
try:
os.set_blocking(self.proc.stderr.fileno(), False)
except OSError:
pass
try:
chunk = self.proc.stderr.read() or b""
except OSError:
return ""
return chunk.decode("utf-8", "replace")
class VaGridPipeline:
"""Arc vacompositor mosaic: chrome + one BGRA pipe per tile."""
def __init__(self, out: DrmOutput):
self.out = out
self.proc: Optional[subprocess.Popen] = None
self._writes: list[int] = []
self._chrome_path: Optional[str] = None
self._chrome_w: Optional[int] = None
self.ingest_w = 0
self.ingest_h = 0
self.n_tiles = 0
self.codec = "i420"
self._fifo_holds: list[int] = []
@property
def alive(self) -> bool:
return self.proc is not None and self.proc.poll() is None
def start(self, grid, ingest_w: int, ingest_h: int, fps: int = 15,
queue_buffers: int = 1, h264_paths: list[str] | None = None) -> None:
self.stop()
if not self.out.connected:
log.warning("va-grid: %s not connected", self.out.name)
return
n = len(grid.identities)
import tempfile
from display_grid import placeholder, display_name, bgra_to_i420
chrome_path = tempfile.NamedTemporaryFile(
prefix="livekit-grid-chrome-", suffix=".bgra", delete=False).name
with open(chrome_path, "wb") as f:
f.write(grid.render())
self._chrome_path = chrome_path
self.codec = "h264" if h264_paths else "i420"
reads: list[int] = []
writes: list[int] = []
if h264_paths:
from h264_fifo import hold_fifo
self._fifo_holds = [hold_fifo(p) for p in h264_paths]
cmd = va_grid_cmd(
self.out, grid, ingest_w, ingest_h, fps,
tile_h264_paths=h264_paths, chrome_path=chrome_path,
queue_buffers=queue_buffers)
pass_fds: tuple[int, ...] = ()
else:
pipe_sz = 1048576
for _ in range(n):
r, w = os.pipe()
os.set_inheritable(r, True)
try:
fcntl.fcntl(w, fcntl.F_SETPIPE_SZ, pipe_sz)
fcntl.fcntl(r, fcntl.F_SETPIPE_SZ, pipe_sz)
except OSError:
pass
reads.append(r)
writes.append(w)
cmd = va_grid_cmd(
self.out, grid, ingest_w, ingest_h, fps,
tile_fds=reads, chrome_path=chrome_path,
queue_buffers=queue_buffers)
pass_fds = tuple(reads)
log.info(
"va-grid: Arc compositor %dx%d ingest %dx%d %d tiles codec=%s -> %s id=%s",
grid.width, grid.height, ingest_w, ingest_h, n, self.codec,
self.out.name, self.out.connector_id)
self.proc = subprocess.Popen(
cmd,
stdin=subprocess.DEVNULL,
stdout=subprocess.DEVNULL,
stderr=None,
pass_fds=pass_fds,
env=_gst_env(),
close_fds=True,
)
for r in reads:
try:
os.close(r)
except OSError:
pass
self._writes = writes
self.ingest_w = ingest_w
self.ingest_h = ingest_h
self.n_tiles = n
if self.codec == "i420":
for i, ident in enumerate(grid.identities):
ph = placeholder(ingest_w, ingest_h, [display_name(ident)])
os.set_blocking(writes[i], True)
self._write_fd(writes[i], bgra_to_i420(ph))
os.set_blocking(writes[i], False)
def push_tile(self, index: int, frame: bytes | memoryview) -> bool:
if not self.alive or index < 0 or index >= self.n_tiles:
return False
return self._write_fd(self._writes[index], frame)
def _write_fd(self, fd: Optional[int], data: bytes | memoryview) -> bool:
if fd is None:
return False
view = memoryview(data)
off = 0
try:
while off < len(view):
try:
n = os.write(fd, view[off:])
if n == 0:
return False
off += n
except BlockingIOError:
if off == 0:
return True
continue
return True
except BrokenPipeError:
log.warning("va-grid: pipe closed (rc=%s)",
None if self.proc is None else self.proc.poll())
return False
def gst_errors(self) -> str:
if self.proc is None or self.proc.stderr is None:
return ""
try:
os.set_blocking(self.proc.stderr.fileno(), False)
except OSError:
pass
try:
chunk = self.proc.stderr.read() or b""
except OSError:
return ""
return chunk.decode("utf-8", "replace")
def stop(self) -> None:
for w in self._writes:
try:
os.close(w)
except OSError:
pass
self._writes = []
for fd in getattr(self, "_fifo_holds", []):
try:
os.close(fd)
except OSError:
pass
self._fifo_holds = []
self._chrome_w = None
if self._chrome_path:
try:
os.unlink(self._chrome_path)
except OSError:
pass
self._chrome_path = None
if self.proc is None:
return
if self.proc.poll() is None:
self.proc.terminate()
try:
self.proc.wait(timeout=3)
except subprocess.TimeoutExpired:
self.proc.kill()
self.proc.wait(timeout=2)
self.proc = None
def _gst_env() -> dict:
env = os.environ.copy()
# Never try X11/Wayland; kmssink is the point.
env.pop("DISPLAY", None)
env.pop("WAYLAND_DISPLAY", None)
env.setdefault("GST_GL_PLATFORM", "egl")
env.setdefault("LIBVA_DRIVER_NAME", "iHD")
env.setdefault("LIBVA_DRM_DEVICE", "/dev/dri/renderD129")
return env