Files
livekit-cameras/encode_daemon.py
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

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())