2026-06-29 23:54:46 -04:00
|
|
|
|
/*
|
|
|
|
|
|
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";
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
|
|
|
|
|
export interface LotusDenoiseConfig {
|
|
|
|
|
|
model: LotusDenoiseModel;
|
2026-06-30 00:28:10 -04:00
|
|
|
|
/** Base URL the worklet scripts/wasm/ESM are served from (e.g. "./denoise/"). */
|
2026-06-29 23:54:46 -04:00
|
|
|
|
assetBase: string;
|
|
|
|
|
|
gate: boolean;
|
|
|
|
|
|
gateThreshold: number;
|
2026-07-01 00:16:23 -04:00
|
|
|
|
/**
|
|
|
|
|
|
* 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;
|
2026-06-29 23:54:46 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-30 00:28:10 -04:00
|
|
|
|
// 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 }
|
2026-06-29 23:54:46 -04:00
|
|
|
|
> = {
|
|
|
|
|
|
rnnoise: {
|
|
|
|
|
|
name: "@sapphi-red/web-noise-suppressor/rnnoise",
|
|
|
|
|
|
script: "rnnoiseWorklet.js",
|
|
|
|
|
|
wasm: "rnnoise.wasm",
|
2026-06-30 00:28:10 -04:00
|
|
|
|
simdWasm: "rnnoise_simd.wasm",
|
2026-06-29 23:54:46 -04:00
|
|
|
|
},
|
|
|
|
|
|
speex: {
|
|
|
|
|
|
name: "@sapphi-red/web-noise-suppressor/speex",
|
|
|
|
|
|
script: "speexWorklet.js",
|
|
|
|
|
|
wasm: "speex.wasm",
|
|
|
|
|
|
},
|
|
|
|
|
|
};
|
2026-06-30 00:28:10 -04:00
|
|
|
|
// The sapphi gate worklet registers under "noise-gate" (hyphenated).
|
|
|
|
|
|
const GATE = {
|
|
|
|
|
|
name: "@sapphi-red/web-noise-suppressor/noise-gate",
|
|
|
|
|
|
script: "noiseGateWorklet.js",
|
|
|
|
|
|
};
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
2026-06-30 00:28:10 -04:00
|
|
|
|
// 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;
|
2026-06-30 00:15:59 -04:00
|
|
|
|
|
2026-06-30 00:28:10 -04:00
|
|
|
|
// Cache fetched wasm per URL so a reconnect/device-switch doesn't re-download.
|
2026-06-30 00:15:59 -04:00
|
|
|
|
const wasmCache = new Map<string, Promise<ArrayBuffer>>();
|
2026-06-30 23:27:44 -04:00
|
|
|
|
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> {
|
2026-06-30 00:15:59 -04:00
|
|
|
|
let p = wasmCache.get(url);
|
|
|
|
|
|
if (!p) {
|
2026-06-30 23:27:44 -04:00
|
|
|
|
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));
|
2026-06-30 00:15:59 -04:00
|
|
|
|
wasmCache.set(url, p);
|
|
|
|
|
|
}
|
|
|
|
|
|
return p;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-30 23:27:44 -04:00
|
|
|
|
/**
|
|
|
|
|
|
* 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.
|
|
|
|
|
|
*/
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// [lotus] 500 ms, not 3 s: this still runs under LiveKit's trackChangeLock (#7),
|
|
|
|
|
|
// and the statechange watcher heals a still-suspended context later anyway.
|
|
|
|
|
|
async function resumeCtx(ctx: AudioContext, timeoutMs = 500): Promise<void> {
|
2026-06-30 23:27:44 -04:00
|
|
|
|
await Promise.race([
|
|
|
|
|
|
ctx.resume().catch(() => undefined),
|
|
|
|
|
|
new Promise<void>((resolve) => setTimeout(resolve, timeoutMs)),
|
|
|
|
|
|
]);
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-30 00:28:10 -04:00
|
|
|
|
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,
|
2026-06-30 00:28:10 -04:00
|
|
|
|
]),
|
|
|
|
|
|
);
|
|
|
|
|
|
} catch {
|
|
|
|
|
|
return false;
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
interface MlNode {
|
|
|
|
|
|
node: AudioNode;
|
|
|
|
|
|
dispose?: () => void;
|
|
|
|
|
|
}
|
2026-06-30 00:15:59 -04:00
|
|
|
|
interface Graph {
|
|
|
|
|
|
source: MediaStreamAudioSourceNode;
|
2026-06-30 00:28:10 -04:00
|
|
|
|
nodes: AudioNode[];
|
|
|
|
|
|
disposes: (() => void)[];
|
2026-06-30 00:15:59 -04:00
|
|
|
|
track: MediaStreamTrack;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// [lotus #26] Minimal local contracts for the two dynamically-imported ESM
|
|
|
|
|
|
// helpers (not bundled here — see the CONTRACT note above). Their exports are
|
|
|
|
|
|
// asserted at runtime so an asset bump that renames/removes one fails loudly
|
|
|
|
|
|
// (and flows into the rnnoise fallback in lotusDenoise.ts) instead of as a
|
|
|
|
|
|
// vague TypeError deep inside `init()`.
|
|
|
|
|
|
interface DtlnModule {
|
|
|
|
|
|
createNoiseSuppressionAudioWorklet: (
|
|
|
|
|
|
ctx: AudioContext,
|
|
|
|
|
|
opts: { bypassUntilReady: boolean },
|
|
|
|
|
|
) => Promise<MlNode>;
|
|
|
|
|
|
}
|
|
|
|
|
|
interface DfnCore {
|
|
|
|
|
|
initialize: () => Promise<void>;
|
|
|
|
|
|
createAudioWorkletNode: (ctx: AudioContext) => Promise<AudioNode>;
|
|
|
|
|
|
destroy: () => void;
|
|
|
|
|
|
}
|
|
|
|
|
|
interface DfnModule {
|
|
|
|
|
|
DeepFilterNet3Core: new (opts: {
|
|
|
|
|
|
sampleRate: number;
|
|
|
|
|
|
noiseReductionLevel: number;
|
|
|
|
|
|
assetConfig: { cdnUrl: string };
|
|
|
|
|
|
}) => DfnCore;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/** Throw a clear error if a dynamic-import module lacks an expected export. */
|
|
|
|
|
|
export function assertModuleExport<T>(
|
|
|
|
|
|
mod: unknown,
|
|
|
|
|
|
name: string,
|
|
|
|
|
|
url: string,
|
|
|
|
|
|
): T {
|
|
|
|
|
|
const exp = (mod as Record<string, unknown> | null | undefined)?.[name];
|
|
|
|
|
|
if (typeof exp !== "function")
|
|
|
|
|
|
throw new Error(
|
|
|
|
|
|
`denoise: ${url} does not export ${name} (got ${typeof exp}) — asset/version mismatch`,
|
|
|
|
|
|
);
|
|
|
|
|
|
return mod as T;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async function loadDfnCore(config: LotusDenoiseConfig): Promise<DfnCore> {
|
|
|
|
|
|
const base = config.assetBase;
|
|
|
|
|
|
const url = `${base}deepfilternet/index.esm.js`;
|
|
|
|
|
|
const dfnBase = new URL(`${base}deepfilternet`, window.location.href).href;
|
|
|
|
|
|
const mod = assertModuleExport<DfnModule>(
|
|
|
|
|
|
await import(/* @vite-ignore */ url),
|
|
|
|
|
|
"DeepFilterNet3Core",
|
|
|
|
|
|
url,
|
|
|
|
|
|
);
|
|
|
|
|
|
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();
|
|
|
|
|
|
return core;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async function loadDtlnModule(config: LotusDenoiseConfig): Promise<DtlnModule> {
|
|
|
|
|
|
const url = `${config.assetBase}workadventure/audio-worklet.js`;
|
|
|
|
|
|
return assertModuleExport<DtlnModule>(
|
|
|
|
|
|
await import(/* @vite-ignore */ url),
|
|
|
|
|
|
"createNoiseSuppressionAudioWorklet",
|
|
|
|
|
|
url,
|
|
|
|
|
|
);
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/** Which wasm file a flat model uses (SIMD build when supported). */
|
|
|
|
|
|
function flatWasmFiles(model: "rnnoise" | "speex"): {
|
|
|
|
|
|
primary: string;
|
|
|
|
|
|
fallback?: string;
|
|
|
|
|
|
} {
|
|
|
|
|
|
const flat = FLAT[model];
|
|
|
|
|
|
const useSimd = model === "rnnoise" && !!flat.simdWasm && supportsSimd();
|
|
|
|
|
|
return useSimd
|
|
|
|
|
|
? { primary: flat.simdWasm!, fallback: flat.wasm }
|
|
|
|
|
|
: { primary: flat.wasm };
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// [lotus #24] Force every node in the graph to a single, explicitly-downmixed
|
|
|
|
|
|
// channel. Without `channelCountMode: "explicit"` the default ("max") IGNORES
|
|
|
|
|
|
// `channelCount`, so a stereo capture device would feed 2 channels into a
|
|
|
|
|
|
// worklet configured with `maxChannels: 1` and sum a stereo dry copy against a
|
|
|
|
|
|
// mono wet one at the destination.
|
|
|
|
|
|
const MONO: AudioNodeOptions = {
|
|
|
|
|
|
channelCount: 1,
|
|
|
|
|
|
channelCountMode: "explicit",
|
|
|
|
|
|
channelInterpretation: "speakers",
|
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
// [lotus #25] Algorithmic latency of each model in samples at its native rate,
|
|
|
|
|
|
// used to delay the DRY copy of the floor mix so it lines up with the wet path
|
|
|
|
|
|
// (otherwise the sum comb-filters — a hollow/phasey colouration on voice).
|
|
|
|
|
|
// - rnnoise: 480-sample (10 ms @ 48 kHz) frames; the sapphi worklet buffers
|
|
|
|
|
|
// 128-sample quanta up to one frame, so the wet path lags by one frame.
|
|
|
|
|
|
// - speex: the sapphi speex worklet uses the same 480-sample framing.
|
|
|
|
|
|
// - dtln: 512-sample block / 128 hop @ 16 kHz (~32 ms) per the DTLN paper —
|
|
|
|
|
|
// best-known, unmeasured (the floor is not mixed for dtln, see buildGraph).
|
|
|
|
|
|
// - deepfilternet: 480-sample hop + 2-frame lookahead @ 48 kHz (~30 ms) per
|
|
|
|
|
|
// DeepFilterNet3 — best-known, unmeasured (floor not mixed for dfn either).
|
|
|
|
|
|
const DRY_DELAY_SAMPLES: Record<LotusDenoiseModel, number> = {
|
|
|
|
|
|
rnnoise: 480,
|
|
|
|
|
|
speex: 480,
|
|
|
|
|
|
dtln: 512,
|
|
|
|
|
|
deepfilternet: 1440,
|
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
|
* Create the model-rate context and register the flat/gate worklet modules.
|
|
|
|
|
|
* Closes the context (and rethrows) on any failure so nothing half-built leaks.
|
|
|
|
|
|
*/
|
|
|
|
|
|
async function createModelContext(
|
|
|
|
|
|
config: LotusDenoiseConfig,
|
|
|
|
|
|
): Promise<AudioContext> {
|
|
|
|
|
|
const rate = sampleRateFor(config.model);
|
|
|
|
|
|
const ctx = new AudioContext({ sampleRate: rate });
|
|
|
|
|
|
try {
|
|
|
|
|
|
if (ctx.sampleRate !== rate)
|
|
|
|
|
|
throw new Error(`denoise: got ${ctx.sampleRate}Hz, need ${rate}Hz`);
|
|
|
|
|
|
// Flat models register via addModule here; DTLN/DeepFilterNet bring their
|
|
|
|
|
|
// own processor via the dynamic-imported helper (see buildMlNode).
|
|
|
|
|
|
if (config.model === "rnnoise" || config.model === "speex")
|
|
|
|
|
|
await ctx.audioWorklet.addModule(
|
|
|
|
|
|
config.assetBase + FLAT[config.model].script,
|
|
|
|
|
|
);
|
|
|
|
|
|
if (config.gate)
|
|
|
|
|
|
await ctx.audioWorklet.addModule(config.assetBase + GATE.script);
|
|
|
|
|
|
return ctx;
|
|
|
|
|
|
} catch (e) {
|
|
|
|
|
|
await ctx.close().catch(() => undefined);
|
|
|
|
|
|
throw e;
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// [lotus #7] Everything heavy that `init()` needs but that does NOT depend on
|
|
|
|
|
|
// the mic track: the AudioContext + worklet modules, the flat wasm binary, and
|
|
|
|
|
|
// (DFN) the fully-initialised model core. `LocalAudioTrack.setProcessor()`
|
|
|
|
|
|
// holds LiveKit's `trackChangeLock` while awaiting `init()`, so every
|
|
|
|
|
|
// mute/unmute/device-switch queues behind it — prepare these as soon as the
|
|
|
|
|
|
// flag is seen (before any track exists) so `init()` only wires them up.
|
|
|
|
|
|
interface PreparedAssets {
|
|
|
|
|
|
ctx: AudioContext;
|
|
|
|
|
|
dfnCore?: DfnCore;
|
|
|
|
|
|
}
|
|
|
|
|
|
const preparedAssets = new Map<string, Promise<PreparedAssets>>();
|
|
|
|
|
|
const preparedKey = (c: LotusDenoiseConfig): string =>
|
|
|
|
|
|
`${c.model}|${c.gate ? 1 : 0}|${c.assetBase}`;
|
|
|
|
|
|
|
|
|
|
|
|
async function prepareUncached(
|
|
|
|
|
|
config: LotusDenoiseConfig,
|
|
|
|
|
|
): Promise<PreparedAssets> {
|
|
|
|
|
|
const ctx = await createModelContext(config);
|
|
|
|
|
|
try {
|
|
|
|
|
|
let dfnCore: DfnCore | undefined;
|
|
|
|
|
|
if (config.model === "rnnoise" || config.model === "speex") {
|
|
|
|
|
|
const { primary, fallback } = flatWasmFiles(config.model);
|
|
|
|
|
|
// Warm the wasm cache; a SIMD miss is fine — buildMlNode falls back.
|
|
|
|
|
|
await fetchWasm(config.assetBase + primary).catch(async () =>
|
|
|
|
|
|
fallback ? fetchWasm(config.assetBase + fallback) : undefined,
|
|
|
|
|
|
);
|
|
|
|
|
|
} else if (config.model === "dtln") {
|
|
|
|
|
|
await loadDtlnModule(config); // warms the browser's module map
|
|
|
|
|
|
} else {
|
|
|
|
|
|
dfnCore = await loadDfnCore(config);
|
|
|
|
|
|
}
|
|
|
|
|
|
return { ctx, dfnCore };
|
|
|
|
|
|
} catch (e) {
|
|
|
|
|
|
await ctx.close().catch(() => undefined);
|
|
|
|
|
|
throw e;
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
|
* Prefetch/prepare the assets for `config` (idempotent per model). Never
|
|
|
|
|
|
* rejects: a failed prepare is evicted so `init()` simply loads inline and
|
|
|
|
|
|
* surfaces the real error there.
|
|
|
|
|
|
*/
|
|
|
|
|
|
export async function prepareDenoiseAssets(
|
|
|
|
|
|
config: LotusDenoiseConfig,
|
|
|
|
|
|
): Promise<void> {
|
|
|
|
|
|
const key = preparedKey(config);
|
|
|
|
|
|
let p = preparedAssets.get(key);
|
|
|
|
|
|
if (!p) {
|
|
|
|
|
|
p = prepareUncached(config);
|
|
|
|
|
|
void p.catch((e) => {
|
|
|
|
|
|
if (preparedAssets.get(key) === p) preparedAssets.delete(key);
|
|
|
|
|
|
logger.warn(`[lotus] denoise prepare failed (${config.model})`, e);
|
|
|
|
|
|
});
|
|
|
|
|
|
preparedAssets.set(key, p);
|
|
|
|
|
|
}
|
|
|
|
|
|
await p.catch(() => undefined);
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/** Take (one-shot) the prepared assets for `config`, if any were prepared. */
|
|
|
|
|
|
function claimPreparedAssets(
|
|
|
|
|
|
config: LotusDenoiseConfig,
|
|
|
|
|
|
): Promise<PreparedAssets> | undefined {
|
|
|
|
|
|
const key = preparedKey(config);
|
|
|
|
|
|
const p = preparedAssets.get(key);
|
|
|
|
|
|
if (p) preparedAssets.delete(key);
|
|
|
|
|
|
return p;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/** Close any prepared-but-unclaimed contexts (call on feature teardown). */
|
|
|
|
|
|
export async function releasePreparedDenoiseAssets(): Promise<void> {
|
|
|
|
|
|
const all = [...preparedAssets.values()];
|
|
|
|
|
|
preparedAssets.clear();
|
|
|
|
|
|
await Promise.all(
|
|
|
|
|
|
all.map(async (p) =>
|
|
|
|
|
|
p
|
|
|
|
|
|
.then(async (a) => {
|
|
|
|
|
|
safeCall(() => a.dfnCore?.destroy());
|
|
|
|
|
|
if (a.ctx.state !== "closed") await a.ctx.close();
|
|
|
|
|
|
})
|
|
|
|
|
|
.catch(() => undefined),
|
|
|
|
|
|
),
|
|
|
|
|
|
);
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-29 23:54:46 -04:00
|
|
|
|
/**
|
|
|
|
|
|
* A LiveKit audio TrackProcessor that runs Lotus ML noise suppression
|
2026-06-30 00:28:10 -04:00
|
|
|
|
* (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.
|
2026-06-29 23:54:46 -04:00
|
|
|
|
*
|
2026-06-30 00:28:10 -04:00
|
|
|
|
* 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-06-29 23:54:46 -04:00
|
|
|
|
*/
|
2026-09-12 11:35:57 -04:00
|
|
|
|
export class LotusDenoiseProcessor implements TrackProcessor<
|
|
|
|
|
|
Track.Kind.Audio,
|
|
|
|
|
|
AudioProcessorOptions
|
|
|
|
|
|
> {
|
2026-06-29 23:54:46 -04:00
|
|
|
|
public readonly name = "lotus-denoise";
|
|
|
|
|
|
public processedTrack?: MediaStreamTrack;
|
|
|
|
|
|
|
2026-06-30 00:15:59 -04:00
|
|
|
|
private ctx?: AudioContext;
|
|
|
|
|
|
private graph?: Graph;
|
2026-06-30 23:27:44 -04:00
|
|
|
|
private ctxStateHandler?: () => void;
|
2026-09-13 01:22:20 -04:00
|
|
|
|
private preparedDfnCore?: DfnCore;
|
|
|
|
|
|
// [lotus #9] True while the mic is muted: we suspend our own context so the
|
|
|
|
|
|
// worklet stops running inference on silence, and the statechange watcher
|
|
|
|
|
|
// must not "heal" that intentional suspension.
|
|
|
|
|
|
private micMuted = false;
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
|
|
|
|
|
public constructor(private readonly config: LotusDenoiseConfig) {}
|
|
|
|
|
|
|
2026-09-13 01:22:20 -04:00
|
|
|
|
/**
|
|
|
|
|
|
* [lotus #9] Mirror the mic's mute state onto the owned context. EC uses
|
|
|
|
|
|
* `stopMicTrackOnMute: false`, so a muted mic keeps producing (silent) frames
|
|
|
|
|
|
* and the ML worklet would otherwise keep running full inference for the
|
|
|
|
|
|
* whole time the user is muted.
|
|
|
|
|
|
*/
|
|
|
|
|
|
public setMicMuted(muted: boolean): void {
|
|
|
|
|
|
this.micMuted = muted;
|
|
|
|
|
|
const ctx = this.ctx;
|
|
|
|
|
|
if (!ctx || ctx.state === "closed") return;
|
|
|
|
|
|
if (muted) {
|
|
|
|
|
|
if (ctx.state === "running") void ctx.suspend().catch(() => undefined);
|
|
|
|
|
|
} else if (ctx.state === "suspended" && this.graph) {
|
|
|
|
|
|
void ctx.resume().catch(() => undefined);
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-30 00:15:59 -04:00
|
|
|
|
public async init(_opts: AudioProcessorOptions): Promise<void> {
|
2026-07-01 00:46:39 -04:00
|
|
|
|
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.
|
2026-09-13 01:22:20 -04:00
|
|
|
|
const core = this.preparedDfnCore;
|
|
|
|
|
|
this.preparedDfnCore = undefined;
|
|
|
|
|
|
if (core) safeCall(() => core.destroy());
|
2026-07-01 00:46:39 -04:00
|
|
|
|
await this.closeContext();
|
|
|
|
|
|
throw e;
|
|
|
|
|
|
}
|
2026-06-29 23:54:46 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
public async restart(opts: AudioProcessorOptions): Promise<void> {
|
2026-06-30 00:15:59 -04:00
|
|
|
|
try {
|
2026-06-30 00:28:10 -04:00
|
|
|
|
await this.ensureContext();
|
2026-06-30 00:15:59 -04:00
|
|
|
|
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.
|
2026-09-12 11:34:08 -04:00
|
|
|
|
// [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.
|
2026-06-30 00:15:59 -04:00
|
|
|
|
logger.warn("[lotus] denoise restart failed; using raw mic", e);
|
|
|
|
|
|
this.disposeGraph(this.graph);
|
|
|
|
|
|
this.graph = undefined;
|
2026-09-12 11:34:08 -04:00
|
|
|
|
this.processedTrack = undefined;
|
2026-06-30 00:15:59 -04:00
|
|
|
|
}
|
2026-06-29 23:54:46 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
public async destroy(): Promise<void> {
|
2026-06-30 00:15:59 -04:00
|
|
|
|
this.disposeGraph(this.graph);
|
|
|
|
|
|
this.graph = undefined;
|
|
|
|
|
|
this.processedTrack = undefined;
|
2026-09-13 01:22:20 -04:00
|
|
|
|
const core = this.preparedDfnCore;
|
|
|
|
|
|
this.preparedDfnCore = undefined;
|
|
|
|
|
|
if (core) safeCall(() => core.destroy());
|
2026-06-30 23:27:44 -04:00
|
|
|
|
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;
|
|
|
|
|
|
}
|
2026-06-30 00:15:59 -04:00
|
|
|
|
this.ctx = undefined;
|
2026-06-30 23:27:44 -04:00
|
|
|
|
if (ctx.state !== "closed") await ctx.close().catch(() => undefined);
|
2026-06-29 23:54:46 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-13 01:22:20 -04:00
|
|
|
|
/** Adopt the prepared context (or create one) + install the state watcher. */
|
2026-06-30 00:28:10 -04:00
|
|
|
|
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
|
|
|
|
|
|
) {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
if (this.ctx.state === "suspended" && !this.micMuted)
|
|
|
|
|
|
await resumeCtx(this.ctx);
|
2026-06-30 00:28:10 -04:00
|
|
|
|
return;
|
|
|
|
|
|
}
|
2026-06-30 23:27:44 -04:00
|
|
|
|
await this.closeContext();
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// [lotus #7] Prefer the context/modules/model prepared before
|
|
|
|
|
|
// setProcessor() was called; only load inline if nothing was prepared
|
|
|
|
|
|
// (e.g. the rnnoise fallback path, or a second processor after a
|
|
|
|
|
|
// republish).
|
|
|
|
|
|
const claimed = await claimPreparedAssets(this.config)?.catch(
|
|
|
|
|
|
() => undefined,
|
|
|
|
|
|
);
|
|
|
|
|
|
let ctx: AudioContext;
|
|
|
|
|
|
if (claimed && claimed.ctx.state !== "closed") {
|
|
|
|
|
|
ctx = claimed.ctx;
|
|
|
|
|
|
this.preparedDfnCore = claimed.dfnCore;
|
|
|
|
|
|
} else {
|
|
|
|
|
|
ctx = await createModelContext(this.config);
|
|
|
|
|
|
}
|
2026-06-30 23:27:44 -04:00
|
|
|
|
try {
|
|
|
|
|
|
// Auto-resume if the OS/browser suspends the context mid-call (mobile
|
|
|
|
|
|
// backgrounding, audio interruption): the dest node otherwise emits
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// silence with no recovery. Only resume while a graph is live and the
|
|
|
|
|
|
// suspension isn't our own mute suspension (#9).
|
2026-06-30 23:27:44 -04:00
|
|
|
|
const onStateChange = (): void => {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
if (ctx.state === "suspended" && this.graph && !this.micMuted)
|
2026-06-30 23:27:44 -04:00
|
|
|
|
void ctx.resume().catch(() => undefined);
|
|
|
|
|
|
};
|
|
|
|
|
|
ctx.addEventListener("statechange", onStateChange);
|
|
|
|
|
|
// The action can arrive via host postMessage, not a gesture in this
|
|
|
|
|
|
// iframe, so the context can start suspended — resume without hanging.
|
2026-09-13 01:22:20 -04:00
|
|
|
|
if (ctx.state === "suspended" && !this.micMuted) await resumeCtx(ctx);
|
|
|
|
|
|
// Attached while already muted (#9): don't let a prepared, running
|
|
|
|
|
|
// context burn inference until the first unmute.
|
|
|
|
|
|
else if (ctx.state === "running" && this.micMuted)
|
|
|
|
|
|
await ctx.suspend().catch(() => undefined);
|
2026-06-30 23:27:44 -04:00
|
|
|
|
|
|
|
|
|
|
this.ctx = ctx;
|
|
|
|
|
|
this.ctxStateHandler = onStateChange;
|
|
|
|
|
|
} catch (e) {
|
|
|
|
|
|
// Don't leak a half-initialised context on any failure path.
|
2026-06-30 00:28:10 -04:00
|
|
|
|
await ctx.close().catch(() => undefined);
|
2026-06-30 23:27:44 -04:00
|
|
|
|
throw e;
|
2026-06-30 00:28:10 -04:00
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
2026-06-30 00:28:10 -04:00
|
|
|
|
private async buildGraph(track: MediaStreamTrack): Promise<Graph> {
|
|
|
|
|
|
const ctx = this.ctx!;
|
2026-06-29 23:54:46 -04:00
|
|
|
|
const source = ctx.createMediaStreamSource(new MediaStream([track]));
|
2026-09-13 01:22:20 -04:00
|
|
|
|
const dest = new MediaStreamAudioDestinationNode(ctx, MONO);
|
2026-06-30 00:28:10 -04:00
|
|
|
|
const nodes: AudioNode[] = [];
|
|
|
|
|
|
const disposes: (() => void)[] = [];
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
2026-07-01 00:46:39 -04:00
|
|
|
|
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;
|
2026-07-01 00:16:23 -04:00
|
|
|
|
|
2026-07-01 00:46:39 -04:00
|
|
|
|
// 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, {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
...MONO,
|
2026-07-01 00:46:39 -04:00
|
|
|
|
processorOptions: {
|
|
|
|
|
|
openThreshold: this.config.gateThreshold,
|
|
|
|
|
|
closeThreshold: this.config.gateThreshold - 5,
|
|
|
|
|
|
holdMs: 150,
|
|
|
|
|
|
maxChannels: 1,
|
|
|
|
|
|
},
|
|
|
|
|
|
});
|
|
|
|
|
|
wetHead.connect(gate);
|
|
|
|
|
|
wetHead = gate;
|
|
|
|
|
|
nodes.push(gate);
|
|
|
|
|
|
}
|
2026-06-29 23:54:46 -04:00
|
|
|
|
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// Only mix a dry floor for the flat models (RNNoise/Speex), whose
|
|
|
|
|
|
// framing latency is known exactly (DRY_DELAY_SAMPLES); the DTLN/DFN
|
|
|
|
|
|
// figures are best-known estimates, so 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.
|
2026-07-01 00:46:39 -04:00
|
|
|
|
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.
|
2026-09-13 01:22:20 -04:00
|
|
|
|
const wetGain = new GainNode(ctx, { ...MONO, gain: 1 - floor });
|
2026-07-01 00:46:39 -04:00
|
|
|
|
wetHead.connect(wetGain);
|
|
|
|
|
|
wetGain.connect(dest);
|
|
|
|
|
|
nodes.push(wetGain);
|
2026-06-30 00:28:10 -04:00
|
|
|
|
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// [lotus #25] Delay the dry copy by the model's algorithmic latency so
|
|
|
|
|
|
// it sums in phase with the (framed, hence delayed) wet path instead
|
|
|
|
|
|
// of comb-filtering against it.
|
|
|
|
|
|
const delaySec = DRY_DELAY_SAMPLES[this.config.model] / ctx.sampleRate;
|
|
|
|
|
|
const dryDelay = new DelayNode(ctx, {
|
|
|
|
|
|
...MONO,
|
|
|
|
|
|
maxDelayTime: Math.max(delaySec, 1 / ctx.sampleRate),
|
|
|
|
|
|
delayTime: delaySec,
|
|
|
|
|
|
});
|
|
|
|
|
|
const dryGain = new GainNode(ctx, { ...MONO, gain: floor });
|
|
|
|
|
|
source.connect(dryDelay);
|
|
|
|
|
|
dryDelay.connect(dryGain);
|
2026-07-01 00:46:39 -04:00
|
|
|
|
dryGain.connect(dest);
|
2026-09-13 01:22:20 -04:00
|
|
|
|
nodes.push(dryDelay, dryGain);
|
2026-07-01 00:46:39 -04:00
|
|
|
|
} else {
|
|
|
|
|
|
wetHead.connect(dest);
|
|
|
|
|
|
}
|
2026-07-01 00:16:23 -04:00
|
|
|
|
|
2026-07-01 00:46:39 -04:00
|
|
|
|
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],
|
|
|
|
|
|
};
|
2026-07-01 00:46:39 -04:00
|
|
|
|
} 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;
|
|
|
|
|
|
}
|
2026-06-30 00:28:10 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
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-13 01:22:20 -04:00
|
|
|
|
const mod = await loadDtlnModule(this.config);
|
|
|
|
|
|
return await mod.createNoiseSuppressionAudioWorklet(ctx, {
|
2026-06-30 00:28:10 -04:00
|
|
|
|
bypassUntilReady: true,
|
2026-09-13 01:22:20 -04:00
|
|
|
|
});
|
2026-06-30 00:28:10 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if (model === "deepfilternet") {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
// [lotus #7] Use the core initialised by prepareDenoiseAssets() if we
|
|
|
|
|
|
// have one (first graph); later rebuilds (restart) load a fresh core.
|
|
|
|
|
|
const prepared = this.preparedDfnCore;
|
|
|
|
|
|
this.preparedDfnCore = undefined;
|
|
|
|
|
|
const core = prepared ?? (await loadDfnCore(this.config));
|
|
|
|
|
|
const node = await core.createAudioWorkletNode(ctx);
|
2026-09-12 11:35:57 -04:00
|
|
|
|
return { node, dispose: () => safeCall(() => core.destroy()) };
|
2026-06-30 00:28:10 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Flat sapphi worklet (rnnoise/speex).
|
|
|
|
|
|
const flat = FLAT[model];
|
2026-09-13 01:22:20 -04:00
|
|
|
|
const { primary, fallback } = flatWasmFiles(model);
|
2026-06-30 00:28:10 -04:00
|
|
|
|
let wasmBinary: ArrayBuffer;
|
|
|
|
|
|
try {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
wasmBinary = await fetchWasm(base + primary);
|
2026-06-30 00:28:10 -04:00
|
|
|
|
} catch (e) {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
if (fallback) {
|
|
|
|
|
|
wasmCache.delete(base + primary);
|
|
|
|
|
|
wasmBinary = await fetchWasm(base + fallback); // fall back to non-SIMD
|
2026-06-30 00:28:10 -04:00
|
|
|
|
} else throw e;
|
|
|
|
|
|
}
|
|
|
|
|
|
const node = new AudioWorkletNode(ctx, flat.name, {
|
2026-09-13 01:22:20 -04:00
|
|
|
|
...MONO,
|
2026-06-29 23:54:46 -04:00
|
|
|
|
numberOfInputs: 1,
|
|
|
|
|
|
numberOfOutputs: 1,
|
2026-06-30 00:28:10 -04:00
|
|
|
|
processorOptions: { maxChannels: 1, wasmBinary },
|
2026-06-29 23:54:46 -04:00
|
|
|
|
});
|
2026-09-12 11:35:57 -04:00
|
|
|
|
return {
|
|
|
|
|
|
node,
|
|
|
|
|
|
dispose: () => safeCall(() => node.port.postMessage("destroy")),
|
|
|
|
|
|
};
|
2026-06-29 23:54:46 -04:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-30 00:15:59 -04:00
|
|
|
|
private disposeGraph(graph: Graph | undefined): void {
|
|
|
|
|
|
if (!graph) return;
|
2026-06-30 00:28:10 -04:00
|
|
|
|
for (const dispose of graph.disposes) safeCall(dispose);
|
|
|
|
|
|
for (const node of graph.nodes) safeCall(() => node.disconnect());
|
|
|
|
|
|
safeCall(() => graph.source.disconnect());
|
2026-06-30 00:15:59 -04:00
|
|
|
|
graph.track.stop();
|
2026-06-29 23:54:46 -04:00
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-06-30 00:28:10 -04:00
|
|
|
|
|
|
|
|
|
|
function safeCall(fn: () => void): void {
|
|
|
|
|
|
try {
|
|
|
|
|
|
fn();
|
|
|
|
|
|
} catch {
|
|
|
|
|
|
/* ignore */
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|