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.
325 lines
12 KiB
Python
325 lines
12 KiB
Python
"""I420 SHM -> iGPU H.264 (or Arc AV1) -> latest-AU mmap. Does not open /dev/camNN."""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import json
|
|
import logging
|
|
import os
|
|
import select
|
|
import shutil
|
|
import signal
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from pathlib import Path
|
|
|
|
from annexb import AnnexBSplitter, nal_unit_type
|
|
from av1_obu import Av1TuSplitter, is_keyframe as av1_is_keyframe
|
|
from frame_shm import LatestFrameReader, frame_bytes, wait_for_ready
|
|
from i420_nv12 import coded_dim, i420_to_nv12_into, nv12_bytes
|
|
from nal_shm import AU_MAX, NalWriter
|
|
from portrait import i420_source_path
|
|
|
|
log = logging.getLogger("cameras.encode-daemon")
|
|
GST = shutil.which("gst-launch-1.0") or "gst-launch-1.0"
|
|
READY_NAME = "ready"
|
|
GROUP = 5
|
|
|
|
|
|
def _pipes(n: int) -> tuple[list[int], list[int]]:
|
|
rs, ws = [], []
|
|
for _ in range(n):
|
|
r, w = os.pipe()
|
|
try:
|
|
import fcntl
|
|
fcntl.fcntl(r, fcntl.F_SETPIPE_SZ, 1048576)
|
|
fcntl.fcntl(w, fcntl.F_SETPIPE_SZ, 1048576)
|
|
except OSError:
|
|
pass
|
|
rs.append(r)
|
|
ws.append(w)
|
|
return rs, ws
|
|
|
|
|
|
def _one_enc(width: int, height: int, fps: int, bitrate_kbps: int,
|
|
fd_in: int, fd_out: int, codec: str = "h264",
|
|
render: str = "/dev/dri/renderD129", hi_kbps: int = 20000) -> list[str]:
|
|
need = frame_bytes(width, height)
|
|
kbps = max(100, int(bitrate_kbps))
|
|
hi = int(height) >= 720
|
|
if hi:
|
|
kbps = max(kbps, max(100, int(hi_kbps)))
|
|
# GOP=1: a 1s GOP left motion trails until the next IDR. Rally is PTZ.
|
|
key_int = 1
|
|
usage = 2 if hi else 7
|
|
cpb = max(80, kbps // 8)
|
|
if codec == "av1":
|
|
# HW AV1 is NV12-only (no I420 entrypoint). Convert is CPU I420→NV12
|
|
# in this process; gst is parse → vaav1enc only (no vapostproc).
|
|
enc = [
|
|
"!", "vaav1enc", f"bitrate={kbps}", f"key-int-max={int(fps)}",
|
|
"target-usage=7", "hierarchical-level=1", "ref-frames=1",
|
|
"gf-group-size=1", "rate-control=cbr",
|
|
f"cpb-size={cpb}",
|
|
]
|
|
fmt = "nv12"
|
|
else:
|
|
# iGPU H.264. NV12-only. GOP=1 so fdsrc preroll does not hang.
|
|
# vah264enc.device-path is read-only; use the per-node factory.
|
|
enc_name = f"va{Path(render).name}h264enc"
|
|
enc = [
|
|
"!", enc_name, f"bitrate={kbps}",
|
|
f"key-int-max={key_int}", "b-frames=0", "ref-frames=1",
|
|
f"target-usage={usage}", "rate-control=cbr", f"cpb-size={cpb}",
|
|
"!", "h264parse", "config-interval=-1",
|
|
"!", "video/x-h264,stream-format=byte-stream,alignment=au",
|
|
]
|
|
fmt = "nv12"
|
|
need = nv12_bytes(width, height)
|
|
head = [
|
|
"fdsrc", f"fd={int(fd_in)}", f"blocksize={need}",
|
|
"do-timestamp=true",
|
|
"!", "rawvideoparse", f"width={int(width)}", f"height={int(height)}",
|
|
f"format={fmt}", f"framerate={int(fps)}/1",
|
|
"!", "queue", "max-size-buffers=1", "max-size-bytes=0", "max-size-time=0",
|
|
"leaky=downstream",
|
|
]
|
|
tail = [
|
|
"!", "queue", "max-size-buffers=1", "max-size-bytes=0", "max-size-time=0",
|
|
"leaky=downstream",
|
|
"!", "fdsink", f"fd={int(fd_out)}", "sync=false",
|
|
]
|
|
return head + enc + tail
|
|
|
|
|
|
def grouped(cams: list[dict], size: int) -> list[list[dict]]:
|
|
by: dict[tuple[int, int], list[dict]] = {}
|
|
for c in cams:
|
|
key = (int(c["width"]), int(c["height"]))
|
|
by.setdefault(key, []).append(c)
|
|
out: list[list[dict]] = []
|
|
for group in by.values():
|
|
for i in range(0, len(group), size):
|
|
out.append(group[i:i + size])
|
|
return out
|
|
|
|
|
|
def main(argv: list[str] | None = None) -> int:
|
|
ap = argparse.ArgumentParser()
|
|
ap.add_argument("--cameras-json", required=True)
|
|
ap.add_argument("--i420-dir", default="/run/livekit-cameras/raw")
|
|
ap.add_argument("--h264-dir", default="/run/livekit-cameras/h264")
|
|
ap.add_argument("--fps", type=int, default=30)
|
|
ap.add_argument("--bitrate", type=int, default=2_000_000)
|
|
ap.add_argument("--vaapi-device", default="/dev/dri/renderD129")
|
|
ap.add_argument("--group-size", type=int, default=GROUP)
|
|
ap.add_argument("--codec", default="h264")
|
|
ap.add_argument("--hi-bitrate", type=int, default=20_000_000,
|
|
help="bps for height>=720 (Rally). Rollback: 12000000")
|
|
ap.add_argument("--portrait-dir", default="",
|
|
help="C920 I420 from portrait daemon; rally stays in --i420-dir")
|
|
args = ap.parse_args(argv)
|
|
codec = (args.codec or "h264").lower()
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [encode-daemon] %(levelname)s %(message)s")
|
|
cams = json.loads(args.cameras_json)
|
|
if not cams:
|
|
log.error("no cameras")
|
|
return 2
|
|
|
|
from vaapi_pin import igpu_render
|
|
|
|
os.environ["LIBVA_DRIVER_NAME"] = "iHD"
|
|
os.environ["GST_VA_ALL_DRIVERS"] = "1"
|
|
os.environ.pop("LIBVA_DRM_DEVICE", None)
|
|
encode_render = igpu_render()
|
|
if codec == "av1":
|
|
log.info("encoder=vaav1enc card0=/dev/dri/renderD128 codec=%s", codec)
|
|
else:
|
|
log.info("encoder=%s codec=%s", f"va{Path(encode_render).name}h264enc", codec)
|
|
|
|
os.makedirs(args.h264_dir, exist_ok=True)
|
|
kbps = max(100, int(args.bitrate) // 1000)
|
|
stop = False
|
|
|
|
def _stop(*_a):
|
|
nonlocal stop
|
|
stop = True
|
|
|
|
signal.signal(signal.SIGTERM, _stop)
|
|
signal.signal(signal.SIGINT, _stop)
|
|
|
|
gst_procs: list[subprocess.Popen] = []
|
|
states = []
|
|
extra_fds: list[int] = []
|
|
|
|
for group in grouped(cams, 1 if codec == "av1" else max(1, args.group_size)):
|
|
n = len(group)
|
|
in_r, in_w = _pipes(n)
|
|
out_r, out_w = _pipes(n)
|
|
for fd in in_r + out_w:
|
|
os.set_inheritable(fd, True)
|
|
for fd in in_w + out_r:
|
|
os.set_inheritable(fd, False)
|
|
for cam, w_in in zip(group, in_w):
|
|
w, h = int(cam["width"]), int(cam["height"])
|
|
if codec in ("av1", "h264"):
|
|
need = nv12_bytes(coded_dim(w), coded_dim(h))
|
|
else:
|
|
need = frame_bytes(w, h)
|
|
try:
|
|
import fcntl
|
|
fcntl.fcntl(w_in, fcntl.F_SETPIPE_SZ, max(1048576, need + 4096))
|
|
except OSError:
|
|
pass
|
|
os.write(w_in, bytes(need))
|
|
for fd in in_w + out_r:
|
|
os.set_blocking(fd, False)
|
|
cmd = [GST, "-q"]
|
|
for cam, r, wfd in zip(group, in_r, out_w):
|
|
w, h = int(cam["width"]), int(cam["height"])
|
|
if codec in ("av1", "h264"):
|
|
w, h = coded_dim(w), coded_dim(h)
|
|
cmd += _one_enc(w, h, args.fps, kbps, r, wfd, codec,
|
|
render=encode_render,
|
|
hi_kbps=max(100, int(args.hi_bitrate) // 1000))
|
|
env = dict(os.environ)
|
|
env["GST_REGISTRY_UPDATE"] = "no"
|
|
proc = subprocess.Popen(
|
|
cmd, pass_fds=tuple(in_r + out_w), env=env,
|
|
stderr=subprocess.STDOUT)
|
|
gst_procs.append(proc)
|
|
for fd in in_r + out_w:
|
|
extra_fds.append(fd)
|
|
for cam, w_in, r_out in zip(group, in_w, out_r):
|
|
ident = cam["identity"]
|
|
w, h = int(cam["width"]), int(cam["height"])
|
|
i420 = i420_source_path(ident, args.i420_dir, args.portrait_dir)
|
|
h264 = os.path.join(args.h264_dir, f"{ident}.h264")
|
|
au_max = 512 * 1024 if codec == "av1" else AU_MAX
|
|
i420_n = frame_bytes(w, h)
|
|
if codec in ("av1", "h264"):
|
|
cw, ch = coded_dim(w), coded_dim(h)
|
|
pipe_n = nv12_bytes(cw, ch)
|
|
nv12 = bytearray(pipe_n)
|
|
else:
|
|
pipe_n = i420_n
|
|
nv12 = None
|
|
states.append({
|
|
"ident": ident, "w": w, "h": h,
|
|
"need": pipe_n,
|
|
"reader": LatestFrameReader(i420),
|
|
"writer": NalWriter(h264, max_payload=au_max),
|
|
"in_w": w_in, "out_r": r_out,
|
|
"scratch": bytearray(i420_n),
|
|
"nv12": nv12,
|
|
"last": -1,
|
|
"split": Av1TuSplitter() if codec == "av1" else AnnexBSplitter(),
|
|
"sps": b"", "pps": b"", "woff": 0, "wseq": -1,
|
|
"codec": codec,
|
|
})
|
|
extra_fds.extend([w_in, r_out])
|
|
log.info("gst pid=%s cams=%s", proc.pid,
|
|
",".join(c["identity"] for c in group))
|
|
|
|
states.sort(key=lambda s: 0 if s["ident"] == "rally" else 1)
|
|
|
|
ready = os.path.join(args.h264_dir, READY_NAME)
|
|
Path(ready).write_text("ok\n")
|
|
|
|
out_map = {s["out_r"]: s for s in states}
|
|
t0 = time.monotonic()
|
|
seqs = {s["ident"]: 0 for s in states}
|
|
try:
|
|
while not stop:
|
|
for s in states:
|
|
pipe = s["nv12"] if s.get("nv12") is not None else s["scratch"]
|
|
if s["woff"]:
|
|
view = memoryview(pipe)
|
|
try:
|
|
n = os.write(s["in_w"], view[s["woff"]:])
|
|
s["woff"] += n
|
|
if s["woff"] >= s["need"]:
|
|
s["woff"] = 0
|
|
s["last"] = s["wseq"]
|
|
except BlockingIOError:
|
|
pass
|
|
except OSError:
|
|
pass
|
|
continue
|
|
seq = s["reader"].copy_into(s["scratch"])
|
|
if seq is None or seq == s["last"]:
|
|
continue
|
|
if s.get("nv12") is not None:
|
|
i420_to_nv12_into(s["scratch"], s["w"], s["h"], s["nv12"])
|
|
pipe = s["nv12"]
|
|
view = memoryview(pipe)
|
|
try:
|
|
n = os.write(s["in_w"], view)
|
|
if n < s["need"]:
|
|
s["woff"] = n
|
|
s["wseq"] = seq
|
|
else:
|
|
s["last"] = seq
|
|
except BlockingIOError:
|
|
s["woff"] = 0
|
|
except OSError:
|
|
pass
|
|
rds, _, _ = select.select(list(out_map), [], [], 0.01)
|
|
for fd in rds:
|
|
s = out_map[fd]
|
|
try:
|
|
chunk = os.read(fd, 65536)
|
|
except OSError:
|
|
continue
|
|
if not chunk:
|
|
continue
|
|
for nal in s["split"].push(chunk):
|
|
if s.get("codec") == "av1":
|
|
if av1_is_keyframe(nal) or nal:
|
|
seqs[s["ident"]] = s["writer"].write(nal)
|
|
continue
|
|
ntype = nal_unit_type(nal)
|
|
if ntype == 7:
|
|
s["sps"] = nal
|
|
continue
|
|
if ntype == 8:
|
|
s["pps"] = nal
|
|
continue
|
|
if ntype == 5:
|
|
payload = s["sps"] + s["pps"] + nal
|
|
elif ntype in (1, 2, 3, 4):
|
|
payload = nal
|
|
else:
|
|
continue
|
|
seqs[s["ident"]] = s["writer"].write(payload)
|
|
if time.monotonic() - t0 >= 10:
|
|
log.info("seqs %s", " ".join(
|
|
f"{k}:{v}" for k, v in sorted(seqs.items())))
|
|
t0 = time.monotonic()
|
|
for p in gst_procs:
|
|
if p.poll() is not None:
|
|
log.error("gst exited rc=%s", p.returncode)
|
|
stop = True
|
|
break
|
|
finally:
|
|
for p in gst_procs:
|
|
if p.poll() is None:
|
|
p.terminate()
|
|
for fd in extra_fds:
|
|
try:
|
|
os.close(fd)
|
|
except OSError:
|
|
pass
|
|
try:
|
|
os.unlink(ready)
|
|
except OSError:
|
|
pass
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|