Files
element-call/src/lotus/lotusDenoiseProcessor.ts
T

421 lines
15 KiB
TypeScript
Raw Normal View History

/*
Copyright 2026 Lotus Guild
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
type AudioProcessorOptions,
type Track,
type TrackProcessor,
} from "livekit-client";
import { logger } from "matrix-js-sdk/lib/logger";
2026-09-12 11:35:57 -04:00
export type LotusDenoiseModel = "rnnoise" | "speex" | "dtln" | "deepfilternet";
export interface LotusDenoiseConfig {
model: LotusDenoiseModel;
/** Base URL the worklet scripts/wasm/ESM are served from (e.g. "./denoise/"). */
assetBase: string;
gate: boolean;
gateThreshold: number;
/**
* Attenuation floor as a dry/wet mix: the fraction (0..~0.3) of the ORIGINAL
* mic blended back under the denoised signal so full suppression never fully
* collapses the noise floor — this is what kills the "underwater"/pumping
* artifact. 0.15 ≈ a -16 dB floor. 0 = full suppression (no floor).
*/
floor: number;
}
// Flat sapphi worklets (RNNoise/Speex): each registers a processor under `name`
// when its `script` module is added; we feed it the fetched wasm binary.
// ⚠️ CONTRACT: these assets are NOT bundled by the fork build — cinny's
// vite.config.js `lotusDenoise()` plugin copies them from
// `@sapphi-red/web-noise-suppressor` / `@workadventure/noise-suppression` /
// `deepfilternet3-noise-filter` into public/element-call/denoise/. The
// worklet/wasm/ESM versions must match what this processor expects. An
// integration smoke-check should assert GET .../denoise/rnnoise.wasm == 200.
const FLAT: Record<
"rnnoise" | "speex",
{ name: string; script: string; wasm: string; simdWasm?: string }
> = {
rnnoise: {
name: "@sapphi-red/web-noise-suppressor/rnnoise",
script: "rnnoiseWorklet.js",
wasm: "rnnoise.wasm",
simdWasm: "rnnoise_simd.wasm",
},
speex: {
name: "@sapphi-red/web-noise-suppressor/speex",
script: "speexWorklet.js",
wasm: "speex.wasm",
},
};
// The sapphi gate worklet registers under "noise-gate" (hyphenated).
const GATE = {
name: "@sapphi-red/web-noise-suppressor/noise-gate",
script: "noiseGateWorklet.js",
};
// DTLN (@workadventure) targets 16kHz and doesn't resample; RNNoise/Speex and
// DeepFilterNet are 48kHz fullband. The worklets don't resample, so the whole
// graph must run at the model's native rate.
const sampleRateFor = (model: LotusDenoiseModel): number =>
model === "dtln" ? 16_000 : 48_000;
// Cache fetched wasm per URL so a reconnect/device-switch doesn't re-download.
const wasmCache = new Map<string, Promise<ArrayBuffer>>();
async function fetchWasmUncached(url: string): Promise<ArrayBuffer> {
const r = await fetch(url);
if (!r.ok) throw new Error(`denoise wasm ${url} -> ${r.status}`);
return r.arrayBuffer();
}
async function fetchWasm(url: string): Promise<ArrayBuffer> {
let p = wasmCache.get(url);
if (!p) {
p = fetchWasmUncached(url);
// Never cache a REJECTED fetch: a transient failure (e.g. a blip during a
// reconnect) must not permanently disable denoise for the whole session.
// Evict on failure so the next restart/device-switch retries.
void p.catch(() => wasmCache.delete(url));
wasmCache.set(url, p);
}
return p;
}
/**
* Resume an AudioContext, but never block indefinitely. A suspended context can
* only resume after a user gesture; the denoise action can arrive via host
* postMessage (no gesture), so `resume()` may stay pending forever. Since this
* runs inside LiveKit's per-track change lock, a hung resume() would deadlock
* every later mute/unmute/device-switch. Race it against a timeout and proceed
* either way — the processor's `statechange` watcher resumes it once a gesture
* lands, and a still-suspended context degrades to (temporary) silence that the
* watcher heals, not a hang.
*/
async function resumeCtx(ctx: AudioContext, timeoutMs = 3_000): Promise<void> {
await Promise.race([
ctx.resume().catch(() => undefined),
new Promise<void>((resolve) => setTimeout(resolve, timeoutMs)),
]);
}
function supportsSimd(): boolean {
try {
// Minimal SIMD module (v128) — validates only where SIMD is supported.
return WebAssembly.validate(
new Uint8Array([
2026-09-12 11:35:57 -04:00
0, 97, 115, 109, 1, 0, 0, 0, 1, 5, 1, 96, 0, 1, 123, 3, 2, 1, 0, 10, 10,
1, 8, 0, 65, 0, 253, 15, 253, 98, 11,
]),
);
} catch {
return false;
}
}
interface MlNode {
node: AudioNode;
dispose?: () => void;
}
interface Graph {
source: MediaStreamAudioSourceNode;
nodes: AudioNode[];
disposes: (() => void)[];
track: MediaStreamTrack;
}
/**
* A LiveKit audio TrackProcessor that runs Lotus ML noise suppression
* (RNNoise / Speex / DTLN / DeepFilterNet) on the local microphone track, as a
* first-class stage in Element Call's publish pipeline — replacing the host's
* `getUserMedia` monkeypatch.
*
* Because it's a real LiveKit processor, EC re-applies it on every
* (re)publish/restart, so denoise survives EC's mid-call reconnect — the root
* cause of A7. It owns a dedicated AudioContext at the model's required sample
* rate (LiveKit does NOT pass an audioContext to restart()), reused across
* restarts and closed on destroy. restart() never throws and never leaves a
* stopped track on the sender: on failure it degrades to the RAW mic track
* rather than silence.
*/
2026-09-12 11:35:57 -04:00
export class LotusDenoiseProcessor implements TrackProcessor<
Track.Kind.Audio,
AudioProcessorOptions
> {
public readonly name = "lotus-denoise";
public processedTrack?: MediaStreamTrack;
private ctx?: AudioContext;
private graph?: Graph;
private ctxStateHandler?: () => void;
public constructor(private readonly config: LotusDenoiseConfig) {}
public async init(_opts: AudioProcessorOptions): Promise<void> {
try {
await this.ensureContext();
this.graph = await this.buildGraph(_opts.track);
this.processedTrack = this.graph.track;
} catch (e) {
// Don't orphan the owned context if graph construction fails (browsers
// cap live AudioContexts, so repeated failed inits could exhaust them).
// The caller degrades to the raw mic; we just release our resources.
await this.closeContext();
throw e;
}
}
public async restart(opts: AudioProcessorOptions): Promise<void> {
try {
await this.ensureContext();
const next = await this.buildGraph(opts.track);
this.disposeGraph(this.graph);
this.graph = next;
this.processedTrack = next.track;
} catch (e) {
// Never go silent on the A7/device-switch path: fall back to raw audio.
// [lotus] IMPORTANT: never assign `opts.track` (LiveKit-owned) here.
// LiveKit's internalStopProcessor() does `processor.processedTrack?.stop()`
// then re-publishes `_mediaStreamTrack` — the SAME object if we set it as
// processedTrack — which kills the live mic on the next stopProcessor()/
// teardown. Leaving processedTrack undefined makes LiveKit fall through to
// its own `_mediaStreamTrack` instead.
logger.warn("[lotus] denoise restart failed; using raw mic", e);
this.disposeGraph(this.graph);
this.graph = undefined;
this.processedTrack = undefined;
}
}
public async destroy(): Promise<void> {
this.disposeGraph(this.graph);
this.graph = undefined;
this.processedTrack = undefined;
await this.closeContext();
}
/** Remove the state watcher and close the owned context, if any. */
private async closeContext(): Promise<void> {
const ctx = this.ctx;
if (!ctx) return;
if (this.ctxStateHandler) {
ctx.removeEventListener("statechange", this.ctxStateHandler);
this.ctxStateHandler = undefined;
}
this.ctx = undefined;
if (ctx.state !== "closed") await ctx.close().catch(() => undefined);
}
/** Create (once) the model-rate context + register the flat worklet modules. */
private async ensureContext(): Promise<void> {
const rate = sampleRateFor(this.config.model);
2026-09-12 11:35:57 -04:00
if (
this.ctx &&
this.ctx.state !== "closed" &&
this.ctx.sampleRate === rate
) {
if (this.ctx.state === "suspended") await resumeCtx(this.ctx);
return;
}
await this.closeContext();
const ctx = new AudioContext({ sampleRate: rate });
try {
if (ctx.sampleRate !== rate)
throw new Error(`denoise: got ${ctx.sampleRate}Hz, need ${rate}Hz`);
// Auto-resume if the OS/browser suspends the context mid-call (mobile
// backgrounding, audio interruption): the dest node otherwise emits
// silence with no recovery. Only resume while a graph is live.
const onStateChange = (): void => {
if (ctx.state === "suspended" && this.graph)
void ctx.resume().catch(() => undefined);
};
ctx.addEventListener("statechange", onStateChange);
// Flat models register via addModule here; DTLN/DeepFilterNet bring their
// own processor via the dynamic-imported helper (see buildMlNode).
if (this.config.model === "rnnoise" || this.config.model === "speex")
await ctx.audioWorklet.addModule(
this.config.assetBase + FLAT[this.config.model].script,
);
if (this.config.gate)
await ctx.audioWorklet.addModule(this.config.assetBase + GATE.script);
// The action can arrive via host postMessage, not a gesture in this
// iframe, so the context can start suspended — resume without hanging.
if (ctx.state === "suspended") await resumeCtx(ctx);
this.ctx = ctx;
this.ctxStateHandler = onStateChange;
} catch (e) {
// Don't leak a half-initialised context on any failure path.
await ctx.close().catch(() => undefined);
throw e;
}
}
private async buildGraph(track: MediaStreamTrack): Promise<Graph> {
const ctx = this.ctx!;
const source = ctx.createMediaStreamSource(new MediaStream([track]));
const dest = ctx.createMediaStreamDestination();
const nodes: AudioNode[] = [];
const disposes: (() => void)[] = [];
try {
// Wet (denoised) path: source → ml → [gate] → wetGain.
const ml = await this.buildMlNode(ctx);
source.connect(ml.node);
nodes.push(ml.node);
if (ml.dispose) disposes.push(ml.dispose);
let wetHead: AudioNode = ml.node;
// Gate AFTER the ML model, not before: gating the raw noisy signal fed
// hard-zeroed frames into the model (discontinuities it must fight) and
// made the threshold operate on pre-denoise levels. Gate the residual.
if (this.config.gate) {
const gate = new AudioWorkletNode(ctx, GATE.name, {
processorOptions: {
openThreshold: this.config.gateThreshold,
closeThreshold: this.config.gateThreshold - 5,
holdMs: 150,
maxChannels: 1,
},
});
wetHead.connect(gate);
wetHead = gate;
nodes.push(gate);
}
// Only mix a dry floor for the LOW-LATENCY flat models (RNNoise/Speex).
// DTLN/DeepFilterNet add tens of ms of algorithmic latency, so summing an
// undelayed dry copy would comb-filter the voice — for those we rely on
// the model's own level (e.g. DFN noiseReductionLevel) instead. RNNoise is
// also where the "robotic/underwater" reports come from, so this targets it.
const lowLatency =
this.config.model === "rnnoise" || this.config.model === "speex";
const floor = lowLatency
? Math.min(0.5, Math.max(0, this.config.floor))
: 0;
if (floor > 0) {
// Dry/wet mix: blend a small amount of the ORIGINAL mic under the
// denoised signal so suppression can't fully collapse the noise floor
// (kills the "underwater"/pumping artifact). During speech (denoised ≈
// original) the two sum back to ~unity; in noise-only gaps the output
// floors at `floor` × original instead of digital silence.
const wetGain = ctx.createGain();
wetGain.gain.value = 1 - floor;
wetHead.connect(wetGain);
wetGain.connect(dest);
nodes.push(wetGain);
const dryGain = ctx.createGain();
dryGain.gain.value = floor;
source.connect(dryGain);
dryGain.connect(dest);
nodes.push(dryGain);
} else {
wetHead.connect(dest);
}
logger.info(
`[lotus] denoise processor active (${this.config.model}, floor=${floor})`,
);
2026-09-12 11:35:57 -04:00
return {
source,
nodes,
disposes,
track: dest.stream.getAudioTracks()[0],
};
} catch (e) {
// A node constructor / model load can throw mid-build; clean up the
// partially-built graph so it doesn't leak (init/restart still fall back
// to the raw mic on the rejection).
this.disposeGraph({
source,
nodes,
disposes,
track: dest.stream.getAudioTracks()[0],
});
throw e;
}
}
private async buildMlNode(ctx: AudioContext): Promise<MlNode> {
const base = this.config.assetBase;
const model = this.config.model;
if (model === "dtln") {
// Self-contained ESM that resolves its own processor + LiteRT wasm +
// TFLite models. bypassUntilReady passes raw audio until the model loads.
2026-09-12 11:35:57 -04:00
const mod = await import(
/* @vite-ignore */ `${base}workadventure/audio-worklet.js`
);
return (await mod.createNoiseSuppressionAudioWorklet(ctx, {
bypassUntilReady: true,
})) as MlNode;
}
if (model === "deepfilternet") {
2026-09-12 11:35:57 -04:00
const dfnBase = new URL(`${base}deepfilternet`, window.location.href)
.href;
const mod = await import(
/* @vite-ignore */ `${base}deepfilternet/index.esm.js`
);
const core = new mod.DeepFilterNet3Core({
sampleRate: 48_000,
// 60, not 80: full-strength suppression is the main source of the
// "over-processed" character; a lower level keeps voice natural while
// the dry/wet floor handles the noise tail.
noiseReductionLevel: 60,
assetConfig: { cdnUrl: dfnBase },
});
await core.initialize();
const node = (await core.createAudioWorkletNode(ctx)) as AudioNode;
2026-09-12 11:35:57 -04:00
return { node, dispose: () => safeCall(() => core.destroy()) };
}
// Flat sapphi worklet (rnnoise/speex).
const flat = FLAT[model];
const useSimd = model === "rnnoise" && !!flat.simdWasm && supportsSimd();
const wasmFile = useSimd ? flat.simdWasm! : flat.wasm;
let wasmBinary: ArrayBuffer;
try {
wasmBinary = await fetchWasm(base + wasmFile);
} catch (e) {
if (useSimd) {
wasmCache.delete(base + wasmFile);
wasmBinary = await fetchWasm(base + flat.wasm); // fall back to non-SIMD
} else throw e;
}
const node = new AudioWorkletNode(ctx, flat.name, {
channelCount: 1,
numberOfInputs: 1,
numberOfOutputs: 1,
processorOptions: { maxChannels: 1, wasmBinary },
});
2026-09-12 11:35:57 -04:00
return {
node,
dispose: () => safeCall(() => node.port.postMessage("destroy")),
};
}
private disposeGraph(graph: Graph | undefined): void {
if (!graph) return;
for (const dispose of graph.disposes) safeCall(dispose);
for (const node of graph.nodes) safeCall(() => node.disconnect());
safeCall(() => graph.source.disconnect());
graph.track.stop();
}
}
function safeCall(fn: () => void): void {
try {
fn();
} catch {
/* ignore */
}
}