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.
137 lines
4.1 KiB
Python
137 lines
4.1 KiB
Python
"""Latest-frame I420 slots in mmap files (single writer, many readers).
|
|
|
|
Double-buffered: the writer fills the inactive slot, then publishes seq.
|
|
Readers copy the active slot; a torn read (seq changed) is discarded.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import mmap
|
|
import os
|
|
import struct
|
|
import time
|
|
from typing import Optional
|
|
|
|
MAGIC = b"LKFR"
|
|
VER = 2
|
|
HEADER = 64
|
|
_HDR = struct.Struct("<4sIIIIIQ") # magic, ver, w, h, seq, nbytes, ts_ns
|
|
SLOTS = 2
|
|
|
|
|
|
def frame_bytes(width: int, height: int) -> int:
|
|
return int(width) * int(height) * 3 // 2
|
|
|
|
|
|
def slot_size(width: int, height: int) -> int:
|
|
return HEADER + SLOTS * frame_bytes(width, height)
|
|
|
|
|
|
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 LatestFrameWriter:
|
|
def __init__(self, path: str, width: int, height: int):
|
|
self.path = path
|
|
self.width = int(width)
|
|
self.height = int(height)
|
|
self.nbytes = frame_bytes(self.width, self.height)
|
|
os.makedirs(os.path.dirname(path) or ".", exist_ok=True)
|
|
size = slot_size(self.width, self.height)
|
|
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
|
|
self._pack_header(0)
|
|
|
|
def _pack_header(self, seq: int) -> None:
|
|
_HDR.pack_into(
|
|
self._map, 0, MAGIC, VER, self.width, self.height, seq, self.nbytes, 0)
|
|
|
|
def write(self, payload: bytes, ts_ns: int = 0) -> int:
|
|
if len(payload) < self.nbytes:
|
|
return self._seq
|
|
nxt = self._seq + 1
|
|
slot = nxt % SLOTS
|
|
off = HEADER + slot * self.nbytes
|
|
self._map[off:off + self.nbytes] = payload[:self.nbytes]
|
|
self._seq = nxt
|
|
_HDR.pack_into(
|
|
self._map, 0, MAGIC, VER, self.width, self.height,
|
|
self._seq, self.nbytes, int(ts_ns))
|
|
return self._seq
|
|
|
|
def close(self) -> None:
|
|
self._map.close()
|
|
|
|
|
|
class LatestFrameReader:
|
|
def __init__(self, path: str):
|
|
self.path = path
|
|
self._map: Optional[mmap.mmap] = None
|
|
self.width = 0
|
|
self.height = 0
|
|
self.nbytes = 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
|
|
buf = _mmap(fd, st.st_size, False)
|
|
finally:
|
|
os.close(fd)
|
|
magic, ver, w, h, _seq, nbytes, _ts = _HDR.unpack_from(buf, 0)
|
|
if magic != MAGIC or ver != VER or nbytes <= 0:
|
|
buf.close()
|
|
return False
|
|
if st.st_size < HEADER + SLOTS * nbytes:
|
|
buf.close()
|
|
return False
|
|
self._map = buf
|
|
self.width, self.height, self.nbytes = w, h, nbytes
|
|
return True
|
|
|
|
def latest(self) -> Optional[tuple[int, bytes]]:
|
|
buf = bytearray(self.nbytes or frame_bytes(16, 16))
|
|
seq = self.copy_into(buf)
|
|
if seq is None:
|
|
return None
|
|
return seq, bytes(buf[:self.nbytes])
|
|
|
|
def copy_into(self, dest: bytearray) -> Optional[int]:
|
|
"""Copy the active slot into dest (must be >= nbytes). None if torn/empty."""
|
|
if not self._open() or self._map is None:
|
|
return None
|
|
if len(dest) < self.nbytes:
|
|
raise ValueError("dest too small")
|
|
s1 = struct.unpack_from("<I", self._map, 16)[0]
|
|
if s1 == 0:
|
|
return None
|
|
slot = s1 % SLOTS
|
|
off = HEADER + slot * self.nbytes
|
|
dest[:self.nbytes] = self._map[off:off + self.nbytes]
|
|
s2 = struct.unpack_from("<I", self._map, 16)[0]
|
|
if s1 != s2:
|
|
return None
|
|
return s1
|
|
|
|
|
|
def wait_for_ready(path: str, timeout: float = 90.0) -> bool:
|
|
deadline = time.time() + timeout
|
|
while time.time() < deadline:
|
|
if os.path.exists(path):
|
|
return True
|
|
time.sleep(0.05)
|
|
return False
|