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.
451 lines
15 KiB
Python
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
|