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.
110 lines
3.4 KiB
Python
110 lines
3.4 KiB
Python
"""Latest-access-unit H.264 slots in mmap (single writer, one reader).
|
|
|
|
Double-buffered like frame_shm: fill the inactive slot, then publish seq.
|
|
Unread AUs are overwritten. Oversize payloads are dropped (seq unchanged).
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import mmap
|
|
import os
|
|
import struct
|
|
from typing import Optional
|
|
|
|
MAGIC = b"LKNA"
|
|
VER = 1
|
|
HEADER = 64
|
|
_HDR = struct.Struct("<4sIIIIQ") # magic, ver, max_payload, seq, nbytes, ts_ns
|
|
SLOTS = 2
|
|
AU_MAX = 256 * 1024
|
|
|
|
|
|
def _mmap(fd: int, size: int, writable: bool) -> mmap.mmap:
|
|
prot = mmap.PROT_READ | (mmap.PROT_WRITE if writable else 0)
|
|
return mmap.mmap(fd, size, flags=mmap.MAP_SHARED, prot=prot)
|
|
|
|
|
|
class NalWriter:
|
|
def __init__(self, path: str, max_payload: int = AU_MAX):
|
|
self.path = path
|
|
self.max_payload = int(max_payload)
|
|
if self.max_payload < 1 or self.max_payload > 2 * 1024 * 1024:
|
|
raise ValueError("max_payload out of range")
|
|
os.makedirs(os.path.dirname(path) or ".", exist_ok=True)
|
|
size = HEADER + SLOTS * self.max_payload
|
|
fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o644)
|
|
try:
|
|
os.ftruncate(fd, size)
|
|
self._map = _mmap(fd, size, True)
|
|
finally:
|
|
os.close(fd)
|
|
self._seq = 0
|
|
_HDR.pack_into(self._map, 0, MAGIC, VER, self.max_payload, 0, 0, 0)
|
|
|
|
def write(self, payload: bytes, ts_ns: int = 0) -> int:
|
|
n = len(payload)
|
|
if n < 1 or n > self.max_payload:
|
|
return self._seq
|
|
nxt = self._seq + 1
|
|
slot = nxt % SLOTS
|
|
off = HEADER + slot * self.max_payload
|
|
self._map[off:off + n] = payload
|
|
self._seq = nxt
|
|
_HDR.pack_into(
|
|
self._map, 0, MAGIC, VER, self.max_payload, self._seq, n, int(ts_ns))
|
|
return self._seq
|
|
|
|
def close(self) -> None:
|
|
self._map.close()
|
|
|
|
|
|
class NalReader:
|
|
def __init__(self, path: str):
|
|
self.path = path
|
|
self._map: Optional[mmap.mmap] = None
|
|
self.max_payload = 0
|
|
self._last = 0
|
|
|
|
def _open(self) -> bool:
|
|
if self._map is not None:
|
|
return True
|
|
try:
|
|
fd = os.open(self.path, os.O_RDONLY)
|
|
except FileNotFoundError:
|
|
return False
|
|
try:
|
|
st = os.fstat(fd)
|
|
if st.st_size < HEADER:
|
|
return False
|
|
self._map = _mmap(fd, st.st_size, False)
|
|
finally:
|
|
os.close(fd)
|
|
magic, ver, max_payload, seq, nbytes, _ts = _HDR.unpack_from(self._map, 0)
|
|
if magic != MAGIC or ver != VER or max_payload < 1:
|
|
self._map.close()
|
|
self._map = None
|
|
return False
|
|
self.max_payload = int(max_payload)
|
|
return True
|
|
|
|
def read(self) -> tuple[int, bytes] | None:
|
|
if not self._open() or self._map is None:
|
|
return None
|
|
magic, ver, max_payload, seq, nbytes, _ts = _HDR.unpack_from(self._map, 0)
|
|
if magic != MAGIC or seq <= self._last or nbytes < 1:
|
|
return None
|
|
if nbytes > max_payload:
|
|
return None
|
|
slot = seq % SLOTS
|
|
off = HEADER + slot * max_payload
|
|
payload = bytes(self._map[off:off + nbytes])
|
|
magic2, _v, _m, seq2, nbytes2, _t = _HDR.unpack_from(self._map, 0)
|
|
if seq2 != seq or nbytes2 != nbytes:
|
|
return None
|
|
self._last = seq
|
|
return seq, payload
|
|
|
|
def close(self) -> None:
|
|
if self._map is not None:
|
|
self._map.close()
|
|
self._map = None
|