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.
435 lines
17 KiB
Python
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)
|