Files
livekit-cameras/audio_cleanup.py
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

435 lines
17 KiB
Python

"""Audio cleanup for camera mics, using speechbrain.
Pipeline per audio hop (16 kHz from the C920; stereo ALSA, LiveKit wants mono):
int16 stereo -> float32 [-1, 1]
-> SpeechBrain AudioNormalizer (keep channels, 16 kHz)
-> Delay-and-Sum beamform (tutorial recipe: STFT + Covariance +
GCC-PHAT TDOA + DelaySum + ISTFT) -> mono
-> MetricGAN-plus neural enhancement
-> int16 mono
The Delay-and-Sum path matches
https://speechbrain.readthedocs.io/en/latest/tutorials/preprocessing/multi-microphone-beamforming.html
("Delay-and-Sum Beamforming" / GCC-PHAT). Wrapped as
speechbrain.lobes.beamform_multimic.DelaySum_Beamformer (a torch.nn.Module).
MVDR / GEV from that tutorial need a separate noise covariance, which a
live C920 hop does not have, so they are not used online.
Two enhancement implementations, selected by ENHANCE_MODE:
* "model" - speechbrain SpectralMaskEnhancement (metricGAN-plus)
* "gain" - high-pass + soft-clip when the model is unavailable
Beamforming is selected independently by BEAMFORM_MODE (auto|force|off).
The preprocessor is created once per worker process and reused.
"""
from __future__ import annotations
import logging
import os
import time
import numpy as np
from audio_gate import hop_mono_int16
log = logging.getLogger("audio.cleanup")
SAMPLE_RATE = 16000
def apply_torch_thread_limits(n: int = 1) -> None:
"""Cap BLAS/Torch pools so N camera workers do not oversubscribe."""
n = max(1, int(n))
os.environ.setdefault("OMP_NUM_THREADS", str(n))
os.environ.setdefault("MKL_NUM_THREADS", str(n))
os.environ.setdefault("OPENBLAS_NUM_THREADS", str(n))
try:
import torch
torch.set_num_threads(n)
torch.set_num_interop_threads(1)
except Exception: # noqa: BLE001
pass
class AudioCleaner:
"""Stateful audio cleaner. Call .process(...) -> int16 mono.
Accepts int16 mono ``[N]`` or stereo ``[N, 2]`` (C920 ALSA). Returns
int16 mono of length N for LiveKit.
"""
def __init__(self, mode: str = "auto",
model_source: str = "speechbrain/metricgan-plus-voicebank",
sample_rate: int = SAMPLE_RATE,
chunk_s: float = 1.0,
hop_s: float = 0.2,
beamform: str = "auto",
torch_num_threads: int = 1,
tdoa_every: int = 5,
enhance_infer: str = "auto",
enhance_socket: str = "",
identity: str = "",
gate=None):
self.mode = mode
self.beamform_mode = beamform
self.sample_rate = sample_rate
self.chunk = max(1, int(chunk_s * sample_rate))
self.hop = max(1, int(hop_s * sample_rate))
self._ctx = np.zeros(0, dtype=np.float32)
self._enhancer = None
self._beamformer = None
self._normalizer = None
self._torch = None
self.backend = "off"
self.last_dt_ms = -1.0
self.torch_num_threads = max(1, int(torch_num_threads))
self.tdoa_every = max(1, int(tdoa_every))
self._stft = None
self._cov = None
self._gccphat = None
self._delaysum = None
self._istft = None
self._last_tdoas = None
self._tdoa_age = 0
self.enhance_infer = (enhance_infer or "auto").lower()
self.infer_backend = "torch"
self._ov_compiled = None
self._ov_request = None
self._ov_input = None
self.enhance_socket = (enhance_socket or "").strip()
self.identity = identity or ""
self._client = None
self._gate = gate
if self.enhance_socket and mode != "off":
self._connect_daemon()
if self._client is not None:
return
if mode == "force":
raise RuntimeError(
f"ENHANCE_MODE=force but enhance daemon "
f"{self.enhance_socket!r} is unreachable")
self._load_torch()
if beamform == "force":
self._load_beamformer()
if self._beamformer is None:
raise RuntimeError(
"BEAMFORM_MODE=force but failed to load "
"SpeechBrain DelaySum_Beamformer")
elif beamform == "auto":
self._load_beamformer()
if mode == "force":
self._load_model(model_source)
if self._enhancer is None:
raise RuntimeError(
f"ENHANCE_MODE=force but failed to load model "
f"'{model_source}'")
elif mode == "auto":
self._load_model(model_source)
if self._enhancer is not None and self.enhance_infer != "torch":
self._load_openvino()
if self.enhance_infer == "openvino" and self._ov_compiled is None:
raise RuntimeError(
"ENHANCE_INFER=openvino but OpenVINO CPU mask failed to load")
self._refresh_backend()
def _connect_daemon(self) -> None:
from enhance_ipc import EnhanceClient
try:
self._client = EnhanceClient(self.enhance_socket)
except Exception as e: # noqa: BLE001
log.warning("enhance daemon %s unreachable (%s)", self.enhance_socket, e)
self._client = None
return
self.backend = "daemon"
log.info("audio enhance via daemon %s identity=%s",
self.enhance_socket, self.identity or "-")
def _refresh_backend(self) -> None:
parts = []
if self._enhancer is not None:
parts.append("model")
if self._ov_compiled is not None:
parts.append("ov")
if self._beamformer is not None:
parts.append("beamform")
self.backend = "+".join(parts) if parts else "off"
def _load_torch(self) -> None:
apply_torch_thread_limits(self.torch_num_threads)
try:
import torch
self._torch = torch
except Exception as e: # noqa: BLE001
log.info("torch unavailable (%s)", e)
self._torch = None
def _load_beamformer(self) -> None:
"""Tutorial Delay-and-Sum: STFT, Covariance, GCC-PHAT, DelaySum, ISTFT."""
if self._torch is None:
log.info("no torch; skipping DelaySum beamformer")
return
try:
from speechbrain.dataio.preprocess import AudioNormalizer
from speechbrain.lobes.beamform_multimic import DelaySum_Beamformer
except Exception as e: # noqa: BLE001
log.info("speechbrain beamformer unavailable (%s)", e)
return
try:
log.info("loading speechbrain DelaySum_Beamformer (GCC-PHAT)")
# Native 16 kHz C920s: AudioNormalizer would resample 16k->16k.
if self.sample_rate != SAMPLE_RATE:
self._normalizer = AudioNormalizer(
sample_rate=self.sample_rate, mix="keep")
else:
self._normalizer = None
bf = DelaySum_Beamformer(sampling_rate=self.sample_rate)
bf.eval()
self._beamformer = bf
self._stft = bf.stft
self._cov = bf.cov
self._gccphat = bf.gccphat
self._delaysum = bf.delaysum
self._istft = bf.istft
log.info(
"DelaySum beamformer ready (stereo in, mono out, tdoa_every=%d)",
self.tdoa_every)
except Exception as e: # noqa: BLE001
log.warning("failed to load DelaySum_Beamformer: %s", e)
self._beamformer = None
self._normalizer = None
def _load_model(self, source: str) -> None:
if self._torch is None:
log.info("torch unavailable; using fallback")
self._enhancer = None
return
try:
from speechbrain.inference.enhancement import (
SpectralMaskEnhancement,
)
except Exception as e: # noqa: BLE001
log.info("speechbrain unavailable (%s); using fallback", e)
self._enhancer = None
return
try:
log.info("loading speechbrain enhancement model %s ...", source)
model = SpectralMaskEnhancement.from_hparams(
source=source,
savedir="/root/.cache/speechbrain-enhancement",
)
model.mods.enhance_model.eval()
self._enhancer = model
log.info("speechbrain enhancement ready (16 kHz in/out)")
except Exception as e: # noqa: BLE001
log.warning("failed to load enhancement model %s: %s "
"(using fallback)", source, e)
self._enhancer = None
def _load_openvino(self) -> None:
"""Compile MetricGAN mask on OpenVINO CPU. STFT/ISTFT stay in Torch."""
try:
from ov_metricgan import compile_cpu, ensure_onnx
except Exception as e: # noqa: BLE001
log.info("openvino helper unavailable (%s)", e)
return
try:
onnx_path = ensure_onnx(self._enhancer, self._torch)
compiled, request, ov_input = compile_cpu(
onnx_path, num_threads=self.torch_num_threads)
self._ov_compiled = compiled
self._ov_request = request
self._ov_input = ov_input
self.infer_backend = "openvino"
except Exception as e: # noqa: BLE001
log.warning("OpenVINO CPU mask failed (%s); using torch enhance", e)
self._ov_compiled = None
self._ov_request = None
self._ov_input = None
self.infer_backend = "torch"
# -- DSP fallback ------------------------------------------------------
@staticmethod
def _simple_clean(x: np.ndarray) -> np.ndarray:
"""Lightweight cleanup: high-pass + soft knee + gentle gain.
Good enough to remove DC offset/low rumble and taming clipping
without any ML model.
"""
# first-order high-pass at ~80 Hz
if len(x) > 1:
alpha = 0.998
y = np.empty_like(x)
y[0] = x[0]
for i in range(1, len(x)):
y[i] = alpha * (y[i - 1] + x[i] - x[i - 1])
x = y
# soft-clip to avoid harsh distortion
x = np.tanh(x * 0.9) / np.tanh(0.9)
# gentle makeup gain
x = x * 1.1
return np.clip(x, -1.0, 1.0)
# -- beamform (tutorial Delay-and-Sum) --------------------------------
def _to_mono(self, x: np.ndarray) -> np.ndarray:
"""x: float32 [T, C] -> float32 [T]."""
if x.shape[1] == 1:
return x[:, 0]
if self._beamformer is None or self._torch is None:
return x.mean(axis=1)
try:
return self._beamform(x)
except Exception as e: # noqa: BLE001
log.warning("DelaySum beamform failed (%s); averaging channels", e)
return x.mean(axis=1)
def _beamform(self, stereo: np.ndarray) -> np.ndarray:
"""SpeechBrain Delay-and-Sum on [T, 2] float32 -> [T] float32.
Recipe from the multi-microphone beamforming tutorial:
STFT -> Covariance -> GCC-PHAT TDOA -> DelaySum -> ISTFT.
"""
torch = self._torch
n = stereo.shape[0]
wav = torch.from_numpy(np.ascontiguousarray(stereo, dtype=np.float32))
if self._normalizer is not None:
wav = self._normalizer(wav, self.sample_rate)
# DelaySum expects [batch, time, channels]
xs = wav.unsqueeze(0)
with torch.no_grad():
Xs = self._stft(xs)
if (self._last_tdoas is None
or self._tdoa_age >= self.tdoa_every):
XXs = self._cov(Xs)
self._last_tdoas = self._gccphat(XXs)
self._tdoa_age = 0
self._tdoa_age += 1
Ys = self._delaysum(Xs, self._last_tdoas)
y = self._istft(Ys)
out = y.squeeze().detach().cpu().numpy().astype(np.float32, copy=False)
if out.ndim > 1:
out = np.reshape(out, -1)
if out.shape[0] != n:
if out.shape[0] > n:
out = out[:n]
else:
out = np.pad(out, (0, n - out.shape[0]))
return out
def _enhance_openvino(self, w):
"""w: torch [1, T] on CPU. Returns numpy [T]."""
from ov_metricgan import infer_mask
torch = self._torch
with torch.no_grad():
feats = self._enhancer.compute_features(w)
feats_np = np.ascontiguousarray(feats.detach().cpu().numpy())
mask_np = infer_mask(self._ov_request, self._ov_input, feats_np)
mask = torch.from_numpy(mask_np)
enhanced = torch.mul(mask, feats)
wav = self._enhancer.hparams.resynth(torch.expm1(enhanced), w)
return wav.squeeze(0).detach().cpu().numpy()
def _to_int16(self, y: np.ndarray) -> np.ndarray:
return (np.clip(y, -1.0, 1.0) * 32767.0).astype(np.int16)
# -- public API ---------------------------------------------------------
def process(self, pcm: np.ndarray) -> np.ndarray:
"""pcm: int16 mono [N] or stereo [N, 2]. Returns int16 mono [N].
Accumulates a `chunk_s` window and runs beamform + model on that
window, emitting only the latest hop (the input length). Until a
full window is buffered, the DSP fallback is used so the stream
is not silent.
"""
pcm = np.asarray(pcm)
t0 = time.perf_counter()
try:
if self._gate is not None:
n = pcm.shape[0]
self._gate.apply(hop_mono_int16(pcm))
if not self._gate.open:
return np.zeros(n, dtype=np.int16)
return self._process_inner(pcm)
finally:
self.last_dt_ms = (time.perf_counter() - t0) * 1000.0
def _process_inner(self, pcm: np.ndarray) -> np.ndarray:
if pcm.size == 0:
return np.zeros(0, dtype=np.int16)
if pcm.ndim == 1:
frames = pcm.reshape(-1, 1)
elif pcm.ndim == 2:
frames = pcm
else:
raise ValueError(f"pcm must be 1-D or 2-D, got shape {pcm.shape}")
n = frames.shape[0]
x = frames.astype(np.float32) / 32768.0
if self._client is not None:
try:
out = self._client.process(self.identity, frames.astype(np.int16))
out = np.asarray(out, dtype=np.int16)
if out.shape[0] != n:
if out.shape[0] > n:
out = out[:n]
else:
out = np.pad(out, (0, n - out.shape[0]))
return out
except Exception as e: # noqa: BLE001
log.warning("enhance daemon failed (%s); using fallback", e)
mono_hop = frames.astype(np.float32).reshape(n, -1).mean(axis=1) / 32768.0
return self._to_int16(self._simple_clean(mono_hop))
# Delay-and-Sum the *hop* (200 ms is ~60 ms here). Beamforming the
# full 1 s window every hop is ~330 ms — not 20-camera real-time.
# Same tutorial modules; TDOA is estimated per hop via GCC-PHAT.
mono_hop = self._to_mono(x)
if self._enhancer is None:
return self._to_int16(self._simple_clean(mono_hop))
self._ctx = np.concatenate([self._ctx, mono_hop])
if len(self._ctx) < self.chunk:
return self._to_int16(self._simple_clean(mono_hop))
self._ctx = self._ctx[-self.chunk:]
torch = self._torch
try:
w = torch.from_numpy(self._ctx.copy()).unsqueeze(0)
if hasattr(self._enhancer, "device"):
w = w.to(self._enhancer.device)
if self._ov_compiled is not None:
enhanced = self._enhance_openvino(w)
else:
# speechbrain >= 1.1 (e.g. metricGAN-plus) requires `lengths`
# on enhance_batch. The value is RELATIVE (multiplied by the
# number of frames internally), so 1.0 = "the whole sequence".
# Older versions accept the call without it (TypeError fallback).
lengths = torch.tensor([1.0]).to(w.device)
try:
with torch.no_grad():
out = self._enhancer.enhance_batch(w, lengths=lengths)
except TypeError:
with torch.no_grad():
out = self._enhancer.enhance_batch(w)
enhanced = out.squeeze(0).detach().cpu().numpy()
except Exception as e: # noqa: BLE001
log.warning("enhance failed (%s); using fallback", e)
return self._to_int16(self._simple_clean(mono_hop))
hop_out = enhanced[-n:] if len(enhanced) >= n else np.pad(
enhanced, (0, n - len(enhanced)))
hop_out = hop_out[:n]
return self._to_int16(hop_out)