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.
88 lines
2.8 KiB
Python
88 lines
2.8 KiB
Python
"""Length-prefixed int16 audio frames for the shared enhance daemon."""
|
|
from __future__ import annotations
|
|
|
|
import socket
|
|
import struct
|
|
|
|
import numpy as np
|
|
|
|
MAGIC = b"ENH1"
|
|
# magic(4) + ident_len(H) + channels(H) + n_samples(I)
|
|
_HEADER = struct.Struct("!4sHHI")
|
|
|
|
|
|
def dump_frame(identity: str, pcm: np.ndarray) -> bytes:
|
|
"""Serialize identity + int16 pcm ([N] or [N, C]) to one frame."""
|
|
pcm = np.asarray(pcm, dtype=np.int16)
|
|
if pcm.ndim == 1:
|
|
pcm = pcm.reshape(-1, 1)
|
|
elif pcm.ndim != 2:
|
|
raise ValueError(f"pcm must be 1-D or 2-D, got {pcm.shape}")
|
|
n_samples, channels = pcm.shape
|
|
ident = identity.encode("utf-8")
|
|
header = _HEADER.pack(MAGIC, len(ident), channels, n_samples)
|
|
return header + ident + np.ascontiguousarray(pcm).tobytes()
|
|
|
|
|
|
def load_frame(blob: bytes) -> tuple[str, np.ndarray]:
|
|
"""Parse a frame into (identity, int16 array)."""
|
|
if len(blob) < _HEADER.size:
|
|
raise ValueError("truncated enhance frame")
|
|
magic, ident_len, channels, n_samples = _HEADER.unpack_from(blob, 0)
|
|
if magic != MAGIC:
|
|
raise ValueError(f"bad enhance magic {magic!r}")
|
|
ident_end = _HEADER.size + ident_len
|
|
payload = blob[ident_end:]
|
|
need = n_samples * channels * 2
|
|
if len(payload) != need:
|
|
raise ValueError(f"pcm size {len(payload)} != {need}")
|
|
ident = blob[_HEADER.size:ident_end].decode("utf-8")
|
|
pcm = np.frombuffer(payload, dtype=np.int16).reshape(n_samples, channels)
|
|
if channels == 1:
|
|
pcm = pcm[:, 0]
|
|
return ident, pcm
|
|
|
|
|
|
def recvall(sock: socket.socket, n: int) -> bytes:
|
|
buf = bytearray()
|
|
while len(buf) < n:
|
|
chunk = sock.recv(n - len(buf))
|
|
if not chunk:
|
|
raise EOFError("enhance socket closed")
|
|
buf.extend(chunk)
|
|
return bytes(buf)
|
|
|
|
|
|
def read_frame(sock: socket.socket) -> tuple[str, np.ndarray]:
|
|
header = recvall(sock, _HEADER.size)
|
|
magic, ident_len, channels, n_samples = _HEADER.unpack(header)
|
|
if magic != MAGIC:
|
|
raise ValueError(f"bad enhance magic {magic!r}")
|
|
rest = recvall(sock, ident_len + n_samples * channels * 2)
|
|
return load_frame(header + rest)
|
|
|
|
|
|
def write_frame(sock: socket.socket, identity: str, pcm: np.ndarray) -> None:
|
|
sock.sendall(dump_frame(identity, pcm))
|
|
|
|
|
|
class EnhanceClient:
|
|
"""One persistent unix-socket connection to the enhance daemon."""
|
|
|
|
def __init__(self, path: str):
|
|
self.path = path
|
|
self._sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
|
self._sock.settimeout(15.0)
|
|
self._sock.connect(path)
|
|
|
|
def process(self, identity: str, pcm: np.ndarray) -> np.ndarray:
|
|
write_frame(self._sock, identity, pcm)
|
|
_, out = read_frame(self._sock)
|
|
return np.asarray(out, dtype=np.int16)
|
|
|
|
def close(self) -> None:
|
|
try:
|
|
self._sock.close()
|
|
except OSError:
|
|
pass
|