Files
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

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