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.
122 lines
3.6 KiB
Python
122 lines
3.6 KiB
Python
"""Int16 PCM hop rings in mmap (single writer, one reader).
|
|
|
|
Unlike video latest-frame slots, audio must not skip hops. The ring keeps
|
|
a few hops; a slow reader drops the oldest.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import mmap
|
|
import os
|
|
import struct
|
|
from typing import Optional
|
|
|
|
MAGIC = b"LKHP"
|
|
VER = 1
|
|
HEADER = 64
|
|
_HDR = struct.Struct("<4sIIIIQ") # magic, ver, hop_bytes, nslots, seq, _pad
|
|
DEFAULT_SLOTS = 8
|
|
|
|
|
|
def hop_bytes(rate: int, channels: int, hop_s: float) -> int:
|
|
n = int(round(float(hop_s) * int(rate)))
|
|
return n * int(channels) * 2
|
|
|
|
|
|
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 HopWriter:
|
|
def __init__(self, path: str, hop_nbytes: int, nslots: int = DEFAULT_SLOTS):
|
|
self.path = path
|
|
self.nbytes = int(hop_nbytes)
|
|
self.nslots = int(nslots)
|
|
os.makedirs(os.path.dirname(path) or ".", exist_ok=True)
|
|
size = HEADER + self.nslots * self.nbytes
|
|
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.nbytes, self.nslots, 0, 0)
|
|
|
|
def write(self, payload: bytes) -> int:
|
|
if len(payload) < self.nbytes:
|
|
return self._seq
|
|
nxt = self._seq + 1
|
|
slot = nxt % self.nslots
|
|
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.nbytes, self.nslots, self._seq, 0)
|
|
return self._seq
|
|
|
|
def close(self) -> None:
|
|
self._map.close()
|
|
|
|
|
|
class HopReader:
|
|
def __init__(self, path: str):
|
|
self.path = path
|
|
self._map: Optional[mmap.mmap] = None
|
|
self.nbytes = 0
|
|
self.nslots = 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
|
|
buf = _mmap(fd, st.st_size, False)
|
|
finally:
|
|
os.close(fd)
|
|
magic, ver, nbytes, nslots, seq, _p = _HDR.unpack_from(buf, 0)
|
|
if magic != MAGIC or ver != VER or nbytes <= 0 or nslots < 2:
|
|
buf.close()
|
|
return False
|
|
self._map = buf
|
|
self.nbytes, self.nslots = nbytes, nslots
|
|
self._last = 0
|
|
return True
|
|
|
|
def _close(self) -> None:
|
|
if self._map is None:
|
|
return
|
|
try:
|
|
self._map.close()
|
|
finally:
|
|
self._map = None
|
|
self.nbytes = 0
|
|
self.nslots = 0
|
|
|
|
def read(self) -> Optional[bytes]:
|
|
if not self._open() or self._map is None:
|
|
return None
|
|
seq = struct.unpack_from("<I", self._map, 16)[0]
|
|
if seq < self._last:
|
|
# Writer restarted (cameras-only): seq went back to 0.
|
|
self._close()
|
|
self._last = 0
|
|
if not self._open() or self._map is None:
|
|
return None
|
|
seq = struct.unpack_from("<I", self._map, 16)[0]
|
|
if seq <= self._last:
|
|
return None
|
|
if seq - self._last > self.nslots:
|
|
self._last = seq - 1
|
|
self._last += 1
|
|
slot = self._last % self.nslots
|
|
off = HEADER + slot * self.nbytes
|
|
return bytes(self._map[off:off + self.nbytes])
|