From 75bbc09e06335faaf9ff7fade76d836acecb86c6 Mon Sep 17 00:00:00 2001 From: steckes Date: Thu, 3 Sep 2026 14:15:13 +0200 Subject: [PATCH 1/7] gc fixes --- README.md | 37 ++++++++ __test__/dispose.spec.ts | 123 +++++++++++++++++++++++++++ index.d.ts | 47 +++++++++++ src/analyzer.rs | 85 ++++++++++++++----- src/error.rs | 8 ++ src/lib.rs | 1 + src/mem.rs | 70 ++++++++++++++++ src/model.rs | 68 ++++++++++++--- src/processor.rs | 84 ++++++++++++++----- src/processor_async.rs | 176 ++++++++++++++++++++++++++++++++++----- src/vad.rs | 72 +++++++++++----- src/vad_async.rs | 94 ++++++++++++++++----- 12 files changed, 747 insertions(+), 118 deletions(-) create mode 100644 __test__/dispose.spec.ts create mode 100644 src/mem.rs diff --git a/README.md b/README.md index 565fee7..20d2e0d 100644 --- a/README.md +++ b/README.md @@ -258,6 +258,43 @@ If your license key is a JWT, refresh it in place instead of rebuilding the obje context.updateBearerToken(renewedJwt) ``` +## Memory management + +`Model`, `Processor`, `ProcessorAsync`, `Vad`, `VadAsync` and `Analyzer` hold large native +allocations behind small JavaScript objects. The binding reports each object's native +footprint to V8 (`napi_adjust_external_memory`), so the garbage collector applies the right +amount of pressure and reclaims dropped instances promptly instead of letting native memory +grow unbounded. + +For deterministic cleanup, every one of these classes also exposes `dispose()`, which +destroys the native object immediately instead of waiting for garbage collection: + +```javascript +const processor = new Processor(model, licenseKey) +try { + processor.initialize(sampleRate, blockSize) + processor.process(block) +} finally { + processor.dispose() +} +``` + +After `dispose()`, every method on the object throws; calling `dispose()` again does +nothing. On the async classes it blocks until in-flight work on the libuv pool finishes. + +Two things to know about cleanup timing: + +- Native cleanup runs on the event loop when the object is finalized, not synchronously at + garbage collection. Finalizers run on event-loop turns, so batches that create many of + these objects back to back hold native memory until the loop turns. An `await` on an + already-resolved promise (a microtask) is not enough; real async boundaries such as I/O, + `setTimeout` or `setImmediate` are. A tight synchronous loop that creates thousands of + objects accumulates their native memory for the duration of the loop; create these + objects per unit of work behind real async boundaries, or reuse a single instance. +- RSS reflects the peak of simultaneously live (or not-yet-finalized) instances: the + allocator reuses freed native memory rather than returning it to the OS, so a burst of N + concurrent instances costs about N x their footprint even after they are dropped. + ## Development Requires a recent Rust toolchain and Node 18+. diff --git a/__test__/dispose.spec.ts b/__test__/dispose.spec.ts new file mode 100644 index 0000000..8994e2f --- /dev/null +++ b/__test__/dispose.spec.ts @@ -0,0 +1,123 @@ +import test from 'ava' + +import { Model, Processor, ProcessorAsync, Vad, VadAsync } from '../index.js' +import { licenseKey, modelPath } from './common.js' + +const enhancementModel = () => Model.fromFile(modelPath('enhancement')) +const vadModel = () => Model.fromFile(modelPath('vad')) + +/** An initialized processor at the model's optimal settings, plus that block size. */ +function initializedProcessor() { + const model = enhancementModel() + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + + const processor = new Processor(model, licenseKey()) + processor.initialize(sampleRate, blockSize) + + return { processor, sampleRate, blockSize, model } +} + +test('sync processor dispose releases memory and rejects later use', (t) => { + const { processor, blockSize, model } = initializedProcessor() + processor.dispose() + + t.throws(() => processor.process(new Float32Array(blockSize)), { message: /disposed/ }) + t.throws(() => processor.initialize(16000, blockSize), { message: /disposed/ }) + t.throws(() => processor.getContext(), { message: /disposed/ }) + t.throws(() => processor.terminateSession(), { message: /disposed/ }) + + // Dispose is idempotent: a second call does nothing and must not throw. + processor.dispose() + + // The model stays usable after one of its processors is disposed. + const second = new Processor(model, licenseKey()) + second.dispose() + t.pass() +}) + +test('sync vad dispose releases memory and rejects later use', (t) => { + const model = vadModel() + const vad = new Vad(model, licenseKey()) + vad.initialize(model.getOptimalSampleRate(), model.getOptimalBlockSize(model.getOptimalSampleRate())) + vad.dispose() + + t.throws(() => vad.process(new Float32Array(160)), { message: /disposed/ }) + t.throws(() => vad.getContext(), { message: /disposed/ }) + vad.dispose() + t.pass() +}) + +test('model dispose rejects later use but keeps created processors working', (t) => { + const model = enhancementModel() + const { processor, blockSize } = initializedProcessorOn(model) + model.dispose() + + t.throws(() => model.getId(), { message: /disposed/ }) + t.throws(() => model.getOptimalSampleRate(), { message: /disposed/ }) + + // The SDK reference-counts the weights, so the processor keeps working. + processor.process(new Float32Array(blockSize)) + t.pass() + + function initializedProcessorOn(model: Model) { + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + const processor = new Processor(model, licenseKey()) + processor.initialize(sampleRate, blockSize) + return { processor, blockSize } + } +}) + +test('async processor dispose rejects later use and is idempotent', async (t) => { + const model = enhancementModel() + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + + const processor = new ProcessorAsync(model, licenseKey()) + await processor.initialize(sampleRate, blockSize) + await processor.dispose() + + await t.throwsAsync(() => processor.process(new Float32Array(blockSize)), { + message: /disposed/, + }) + await t.throwsAsync(() => processor.terminateSession(), { message: /disposed/ }) + await processor.dispose() + t.pass() +}) + +test('async vad dispose rejects later use', async (t) => { + const model = vadModel() + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + + const vad = new VadAsync(model, licenseKey()) + await vad.initialize(sampleRate, blockSize) + await vad.dispose() + + await t.throwsAsync(() => vad.process(new Float32Array(blockSize)), { message: /disposed/ }) + t.pass() +}) + +test('a withConfig handle outlives its original handle being collected', async (t) => { + const model = enhancementModel() + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + + // The chaining form leaves the constructor's handle as garbage. Its finalizer must + // give back only its own claim, not destroy the shared native processor. + let processor = await new ProcessorAsync(model, licenseKey()).withConfig(sampleRate, blockSize) + + // Encourage a GC so the dropped handle's finalizer runs. + const churn: unknown[] = [] + for (let i = 0; i < 10_000; i++) churn.push({ i }) + churn.length = 0 + if (typeof global.gc === 'function') global.gc() + await new Promise((resolve) => setImmediate(resolve)) + + const block = new Float32Array(blockSize) + await processor.process(block) + await processor.dispose() + processor = undefined as unknown as ProcessorAsync + t.pass() +}) diff --git a/index.d.ts b/index.d.ts index 9724270..3c50fe5 100644 --- a/index.d.ts +++ b/index.d.ts @@ -20,6 +20,14 @@ export declare class Analyzer { /** Creates an analyzer from an analysis model. Other model types are rejected. */ constructor(model: Model, licenseKey: string) + /** + * Destroys the native collector and analyzer immediately, releasing their memory + * without waiting for garbage collection. + * + * Every later method throws; calling `dispose()` again does nothing. Blocks until an + * in-flight `analyzeAsync` on a worker thread finishes. + */ + dispose(): void /** * Configures the analyzer for an audio format. Must be called before buffering. * @@ -89,6 +97,15 @@ export declare class Model { * {@link Model.download}. */ static fromFile(path: string): Model + /** + * Unmaps the model file immediately, releasing its footprint without waiting for + * garbage collection. + * + * Objects already created from the model keep working: the SDK keeps the weights + * alive through internal reference counting. Every later method on this handle + * throws; calling `dispose()` again does nothing. + */ + dispose(): void /** * Downloads a model from the ai-coustics artifact CDN and resolves to its path. * @@ -134,6 +151,13 @@ export declare class Processor { * instance. */ constructor(model: Model, licenseKey: string, otelConfig?: OtelConfig | undefined | null) + /** + * Destroys the native processor immediately, releasing its memory and telemetry + * session without waiting for garbage collection. + * + * Every later method throws; calling `dispose()` again does nothing. + */ + dispose(): void /** * Configures the processor for an audio format. Must be called before processing. * @@ -198,6 +222,14 @@ export declare class ProcessorAsync { * instance. */ constructor(model: Model, licenseKey: string, otelConfig?: OtelConfig | undefined | null) + /** + * Destroys the native processor immediately, releasing its memory and telemetry + * session without waiting for garbage collection. + * + * Every later method throws; calling `dispose()` again does nothing. The synchronous + * variant blocks until in-flight work on the libuv pool finishes. + */ + dispose(): void /** * Initializes the processor and resolves to a handle onto it, for chaining off the * constructor: @@ -305,6 +337,13 @@ export declare class Vad { * instance. */ constructor(model: Model, licenseKey: string, otelConfig?: OtelConfig | undefined | null) + /** + * Destroys the native VAD immediately, releasing its memory and telemetry session + * without waiting for garbage collection. + * + * Every later method throws; calling `dispose()` again does nothing. + */ + dispose(): void /** * Configures the VAD for an audio format. Must be called before processing. * @@ -359,6 +398,14 @@ export declare class VadAsync { * instance. */ constructor(model: Model, licenseKey: string, otelConfig?: OtelConfig | undefined | null) + /** + * Destroys the native VAD immediately, releasing its memory and telemetry session + * without waiting for garbage collection. + * + * Every later method throws; calling `dispose()` again does nothing. The synchronous + * variant blocks until in-flight work on the libuv pool finishes. + */ + dispose(): void /** * Initializes the VAD and resolves to a handle onto it, for chaining off the * constructor: diff --git a/src/analyzer.rs b/src/analyzer.rs index 4e01c16..190eb16 100644 --- a/src/analyzer.rs +++ b/src/analyzer.rs @@ -1,6 +1,7 @@ use crate::{ claim_sdk_id, - error::{Result, map_err}, + error::{Result, disposed_error, map_err}, + mem, model::Model, processor::audio_config, processor_async::{Shared, lock}, @@ -8,6 +9,7 @@ use crate::{ use napi::{ Env, Task, + bindgen_prelude::ObjectFinalize, bindgen_prelude::{AsyncTask, Float32Array}, }; use napi_derive::napi; @@ -66,7 +68,7 @@ impl From for AnalysisResult { /// as one object here, but the split still shows through: {@link Analyzer#buffer} drives /// the collector on the calling thread, while {@link Analyzer#analyzeAsync} moves the /// analyzer half onto a worker. The SDK guarantees the two are safe to use concurrently. -#[napi] +#[napi(custom_finalize)] pub struct Analyzer { // Neither half borrows the other, nor the model: `'static` here is the model weights' // lifetime, which `Model.fromFile` satisfies by memory-mapping the file. @@ -74,24 +76,49 @@ pub struct Analyzer { // Only the analyzer half is shared. The collector is owned outright, so `buffer`, the // one call on the audio path, takes no lock and cannot contend with an analysis running // on a worker thread. - collector: aic_sdk::Collector, - analyzer: Shared>, + collector: Option, + analyzer: Shared>>, +} + +impl ObjectFinalize for Analyzer { + fn finalize(self, env: Env) -> Result<()> { + // `dispose()` already gave the footprint back when the collector is gone. + if self.collector.is_some() { + mem::adjust(env, -mem::ANALYZER_BYTES); + } + Ok(()) + } } #[napi] impl Analyzer { /// Creates an analyzer from an analysis model. Other model types are rejected. #[napi(constructor)] - pub fn new(model: &Model, license_key: String) -> Result { + pub fn new(env: Env, model: &Model, license_key: String) -> Result { + let model_inner = model.sdk()?; claim_sdk_id(); - let (collector, analyzer) = map_err(aic_sdk::analyzer_pair(&model.inner, &license_key))?; + let (collector, analyzer) = map_err(aic_sdk::analyzer_pair(model_inner, &license_key))?; + mem::adjust(env, mem::ANALYZER_BYTES); Ok(Self { - collector, - analyzer: Arc::new(Mutex::new(analyzer)), + collector: Some(collector), + analyzer: Arc::new(Mutex::new(Some(analyzer))), }) } + /// Destroys the native collector and analyzer immediately, releasing their memory + /// without waiting for garbage collection. + /// + /// Every later method throws; calling `dispose()` again does nothing. Blocks until an + /// in-flight `analyzeAsync` on a worker thread finishes. + #[napi] + pub fn dispose(&mut self, env: Env) { + if self.collector.take().is_some() { + lock(&self.analyzer).take(); // dropped on scope exit + mem::adjust(env, -mem::ANALYZER_BYTES); + } + } + /// Configures the analyzer for an audio format. Must be called before buffering. /// /// The model's optimal sample rate and block size avoid internal resampling and @@ -103,17 +130,31 @@ impl Analyzer { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err( - self - .collector - .initialize(&audio_config(sample_rate, block_size, variable_block_size)), - ) + let collector = self + .collector + .as_mut() + .ok_or_else(|| disposed_error("Analyzer"))?; + map_err(collector.initialize(&audio_config(sample_rate, block_size, variable_block_size))) + } + + /// Runs `f` with the analyzer half, or fails with the disposed error. + fn with_analyzer( + &self, + f: impl FnOnce(&mut aic_sdk::Analyzer<'static>) -> Result, + ) -> Result { + let mut guard = lock(&self.analyzer); + let analyzer = guard.as_mut().ok_or_else(|| disposed_error("Analyzer"))?; + f(analyzer) } /// Buffers a mono audio block for later analysis, leaving the audio unmodified. #[napi] pub fn buffer(&mut self, audio: Float32Array) -> Result<()> { - map_err(self.collector.buffer(&audio)) + let collector = self + .collector + .as_mut() + .ok_or_else(|| disposed_error("Analyzer"))?; + map_err(collector.buffer(&audio)) } /// Runs the analysis model over the buffered audio, on the calling thread. @@ -127,7 +168,9 @@ impl Analyzer { /// nothing else is waiting on the event loop. #[napi] pub fn analyze(&self) -> Result { - map_err(lock(&self.analyzer).analyze_buffered()).map(AnalysisResult::from) + self + .with_analyzer(|analyzer| map_err(analyzer.analyze_buffered())) + .map(AnalysisResult::from) } /// Runs the analysis model over the buffered audio on a worker thread. @@ -153,7 +196,7 @@ impl Analyzer { /// Clears buffered audio and internal state, keeping the configured audio settings. #[napi] pub fn reset(&self) -> Result<()> { - map_err(lock(&self.analyzer).reset()) + self.with_analyzer(|analyzer| map_err(analyzer.reset())) } /// Swaps in a renewed JWT without tearing down the analyzer. @@ -162,13 +205,13 @@ impl Analyzer { /// call is a no-op and the previous token stays active. #[napi] pub fn update_bearer_token(&self, token: String) -> Result<()> { - map_err(lock(&self.analyzer).update_bearer_token(&token)) + self.with_analyzer(|analyzer| map_err(analyzer.update_bearer_token(&token))) } /// Ends this analyzer's telemetry session, after which it can no longer analyze audio. #[napi] pub fn terminate_session(&self) -> Result<()> { - map_err(lock(&self.analyzer).terminate_session()) + self.with_analyzer(|analyzer| map_err(analyzer.terminate_session())) } } @@ -177,7 +220,7 @@ impl Analyzer { /// Holds only the analyzer half, so the collector stays on the JS thread where `buffer` can /// keep reaching it while this runs. pub struct AnalyzeTask { - analyzer: Shared>, + analyzer: Shared>>, } impl Task for AnalyzeTask { @@ -185,7 +228,9 @@ impl Task for AnalyzeTask { type JsValue = AnalysisResult; fn compute(&mut self) -> Result { - map_err(lock(&self.analyzer).analyze_buffered()) + let mut guard = lock(&self.analyzer); + let analyzer = guard.as_mut().ok_or_else(|| disposed_error("Analyzer"))?; + map_err(analyzer.analyze_buffered()) } fn resolve(&mut self, _env: Env, result: aic_sdk::AnalysisResult) -> Result { diff --git a/src/error.rs b/src/error.rs index 2c6f409..a2d3050 100644 --- a/src/error.rs +++ b/src/error.rs @@ -32,3 +32,11 @@ pub type Result = napi::Result; pub fn map_err(result: std::result::Result) -> Result { result.map_err(|error| JsAicError(error).into()) } + +/// The error thrown by any method called after `dispose()`. +pub fn disposed_error(class: &str) -> napi::Error { + napi::Error::new( + napi::Status::GenericFailure, + format!("{class} has been disposed"), + ) +} diff --git a/src/lib.rs b/src/lib.rs index 8c43479..6736539 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -9,6 +9,7 @@ use napi_derive::napi; mod analyzer; mod error; +mod mem; mod model; mod processor; mod processor_async; diff --git a/src/mem.rs b/src/mem.rs new file mode 100644 index 0000000..530f6dd --- /dev/null +++ b/src/mem.rs @@ -0,0 +1,70 @@ +//! Reports the native footprint of SDK objects to V8's garbage collector. +//! +//! Each binding class holds a small native allocation behind a small JS object: the JS +//! side is a few dozen bytes, while the native side ranges from ~200 KiB for a Processor +//! to the model weights themselves. Without a signal, V8's GC heuristics never see that +//! cost: `heapUsed` and `external` barely move, so the collector feels no pressure to +//! reclaim dropped instances, and a workload that creates processors per unit of work +//! ratchets RSS up until the process is OOM-killed. +//! +//! `Env::adjust_external_memory` is the Node-API mechanism for this (`napi_adjust_external_memory`, +//! the same fix the Ruby binding ships as `rb_gc_adjust_memory_usage`). Each constructor +//! reports its object's footprint; `ObjectFinalize::finalize` reports the negation when the +//! instance is collected, so the ledger balances. +//! +//! Values are estimates keyed to measurement and deliberately err high: over-reporting +//! only makes V8 collect a little more eagerly, while under-reporting is what caused the +//! unbounded growth. Per-class constants are used because the SDK does not (yet) expose a +//! per-instance memory query. + +use std::path::Path; + +use napi::Env; + +const MIB: i64 = 1024 * 1024; +const KIB: i64 = 1024; + +/// A `Processor` or `Vad` instance. Measured by holding N instances and reading the RSS +/// delta, at initialize + one `process`: +/// +/// - `quail-vf-2.2-s` (5 MiB model): ~190 KiB +/// - `quail-vf-2.2-l` (20 MiB model): ~462 KiB +/// - `vad-2.1-xxs` (0.6 MiB model): ~197 KiB +/// +/// The workspace barely scales with the weights (those stay file-backed under `Model`), +/// so 512 KiB covers the largest measured enhancement model with headroom while staying +/// within ~3x of the smallest. +pub(crate) const PROCESSOR_BYTES: i64 = 512 * KIB; + +/// An `Analyzer` (collector + analyzer pair). Measured ~8.2 MiB for `tyto-1.1-l` at +/// construction and ~8.9 MiB with the collector initialized and holding 5 s of audio. +/// 16 MiB gives ~2x headroom for larger analysis models. +pub(crate) const ANALYZER_BYTES: i64 = 16 * MIB; + +/// Fallback footprint for a `Model` when its file cannot be stat'd. Deliberately +/// conservative: the loaded model is memory-mapped, so its resident share approaches the +/// file size as pages are touched. +const MODEL_FALLBACK_BYTES: i64 = 64 * MIB; + +/// Tells V8 that `delta_bytes` of external (native) memory changed. +/// +/// Positive when an object is created, negative (the same value) when it is finalized. +/// The result is a GC hint only: a failed adjustment is never worth failing the API call +/// over, so errors are ignored. +pub(crate) fn adjust(env: Env, delta_bytes: i64) { + if delta_bytes == 0 { + return; + } + + // A non-Ok status only means the hint was not applied; correctness does not depend on + // it, and there is nothing useful to do about a failed hint anyway. + let _ = env.adjust_external_memory(delta_bytes); +} + +/// The resident footprint to report for a model loaded from `path`: the file size, since +/// the weights are memory-mapped, falling back to a conservative estimate. +pub(crate) fn model_bytes(path: &Path) -> i64 { + std::fs::metadata(path) + .map(|meta| meta.len() as i64) + .unwrap_or(MODEL_FALLBACK_BYTES) +} diff --git a/src/model.rs b/src/model.rs index 7cafddd..e846a19 100644 --- a/src/model.rs +++ b/src/model.rs @@ -1,6 +1,9 @@ -use crate::error::{JsAicError, Result, map_err}; +use crate::{ + error::{JsAicError, Result, disposed_error, map_err}, + mem, +}; -use napi::{Env, Task, bindgen_prelude::AsyncTask}; +use napi::{Env, Task, bindgen_prelude::AsyncTask, bindgen_prelude::ObjectFinalize}; use napi_derive::napi; /// A loaded ai-coustics model. @@ -8,11 +11,31 @@ use napi_derive::napi; /// One model can back multiple processors, VADs and analyzers, according to its type. /// The underlying model data is kept alive by every object created from it through /// internal reference counting, so this handle may be released first. -#[napi] +#[napi(custom_finalize)] pub struct Model { // `from_file` memory-maps the file rather than borrowing a caller-owned buffer, so // the SDK model is `'static` and needs no lifetime plumbing here. - pub(crate) inner: aic_sdk::Model<'static>, + pub(crate) inner: Option>, + /// Native footprint reported to V8's GC while this instance is alive (the mmap'd + /// weights). Reported back on dispose or finalize so the accounting balances. + reported_bytes: i64, +} + +impl ObjectFinalize for Model { + fn finalize(self, env: Env) -> Result<()> { + // `dispose()` already gave the footprint back when the inner is gone. + if self.inner.is_some() { + mem::adjust(env, -self.reported_bytes); + } + Ok(()) + } +} + +impl Model { + /// The SDK model, or the disposed error once `dispose()` ran. + pub(crate) fn sdk(&self) -> std::result::Result<&aic_sdk::Model<'static>, napi::Error> { + self.inner.as_ref().ok_or_else(|| disposed_error("Model")) + } } #[napi] @@ -26,12 +49,30 @@ impl Model { /// Browse available models at , or fetch one with /// {@link Model.download}. #[napi(factory)] - pub fn from_file(path: String) -> Result { + pub fn from_file(env: Env, path: String) -> Result { + let inner = map_err(aic_sdk::Model::from_file(&path))?; + let reported_bytes = mem::model_bytes(std::path::Path::new(&path)); + mem::adjust(env, reported_bytes); + Ok(Self { - inner: map_err(aic_sdk::Model::from_file(path))?, + inner: Some(inner), + reported_bytes, }) } + /// Unmaps the model file immediately, releasing its footprint without waiting for + /// garbage collection. + /// + /// Objects already created from the model keep working: the SDK keeps the weights + /// alive through internal reference counting. Every later method on this handle + /// throws; calling `dispose()` again does nothing. + #[napi] + pub fn dispose(&mut self, env: Env) { + if self.inner.take().is_some() { + mem::adjust(env, -self.reported_bytes); + } + } + /// Downloads a model from the ai-coustics artifact CDN and resolves to its path. /// /// The manifest is re-fetched on every call so the newest compatible model version @@ -51,8 +92,9 @@ impl Model { /// The model identifier, e.g. `quail-vf-2.2-s-16khz`. #[napi] - pub fn get_id(&self) -> String { - self.inner.id().to_owned() + pub fn get_id(&self) -> Result { + let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Model"))?; + Ok(inner.id().to_owned()) } /// The sample rate in Hz the model was trained for. @@ -60,8 +102,9 @@ impl Model { /// Audio at any rate can be processed, but a model only enhances frequencies up to /// its own Nyquist limit, so matching this rate gives the best quality. #[napi] - pub fn get_optimal_sample_rate(&self) -> u32 { - self.inner.optimal_sample_rate() + pub fn get_optimal_sample_rate(&self) -> Result { + let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Model"))?; + Ok(inner.optimal_sample_rate()) } /// The block size that avoids internal buffering at `sampleRate`. @@ -70,12 +113,13 @@ impl Model { /// The value changes with the sample rate, because the model works on a fixed time /// window: a 10 ms window is 480 samples at 48 kHz but 160 at 16 kHz. #[napi] - pub fn get_optimal_block_size(&self, sample_rate: u32) -> u32 { + pub fn get_optimal_block_size(&self, sample_rate: u32) -> Result { + let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Model"))?; // The SDK reports sizes as `usize`, which napi would marshal as a JS BigInt. // A BigInt block size would throw on `new Float32Array(n)` and on arithmetic // against plain numbers, so it crosses the boundary as u32. Block sizes are a // few thousand samples at most. - self.inner.optimal_block_size(sample_rate) as u32 + Ok(inner.optimal_block_size(sample_rate) as u32) } } diff --git a/src/processor.rs b/src/processor.rs index 9c7c487..7154e91 100644 --- a/src/processor.rs +++ b/src/processor.rs @@ -1,10 +1,11 @@ use crate::{ claim_sdk_id, - error::{Result, map_err}, + error::{Result, disposed_error, map_err}, + mem, model::Model, }; -use napi::bindgen_prelude::Float32Array; +use napi::{Env, bindgen_prelude::Float32Array, bindgen_prelude::ObjectFinalize}; use napi_derive::napi; /// Enhancement parameters, all changeable while audio is being processed. @@ -78,9 +79,19 @@ pub(crate) fn audio_config( /// and {@link Analyzer} for analysis models; passing the wrong kind throws. /// /// Create several processors to handle multiple streams or to switch models at runtime. -#[napi] +#[napi(custom_finalize)] pub struct Processor { - inner: aic_sdk::Processor<'static>, + inner: Option>, +} + +impl ObjectFinalize for Processor { + fn finalize(self, env: Env) -> Result<()> { + // `dispose()` already gave the footprint back when the inner is gone. + if self.inner.is_some() { + mem::adjust(env, -mem::PROCESSOR_BYTES); + } + Ok(()) + } } #[napi] @@ -90,17 +101,35 @@ impl Processor { /// Telemetry follows the runtime environment; pass `otelConfig` to override it for this /// instance. #[napi(constructor)] - pub fn new(model: &Model, license_key: String, otel_config: Option) -> Result { + pub fn new( + env: Env, + model: &Model, + license_key: String, + otel_config: Option, + ) -> Result { + let model_inner = model.sdk()?; claim_sdk_id(); let inner = match otel_config { Some(config) => { - aic_sdk::Processor::with_otel_config(&model.inner, &license_key, &config.into()) + aic_sdk::Processor::with_otel_config(model_inner, &license_key, &config.into()) } - None => aic_sdk::Processor::new(&model.inner, &license_key), + None => aic_sdk::Processor::new(model_inner, &license_key), }; - Ok(Self { - inner: map_err(inner)?, - }) + let inner = map_err(inner)?; + mem::adjust(env, mem::PROCESSOR_BYTES); + + Ok(Self { inner: Some(inner) }) + } + + /// Destroys the native processor immediately, releasing its memory and telemetry + /// session without waiting for garbage collection. + /// + /// Every later method throws; calling `dispose()` again does nothing. + #[napi] + pub fn dispose(&mut self, env: Env) { + if self.inner.take().is_some() { + mem::adjust(env, -mem::PROCESSOR_BYTES); + } } /// Configures the processor for an audio format. Must be called before processing. @@ -117,11 +146,11 @@ impl Processor { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err( - self - .inner - .initialize(&audio_config(sample_rate, block_size, variable_block_size)), - ) + let inner = self + .inner + .as_mut() + .ok_or_else(|| disposed_error("Processor"))?; + map_err(inner.initialize(&audio_config(sample_rate, block_size, variable_block_size))) } /// Enhances a mono audio block in place. @@ -139,18 +168,27 @@ impl Processor { // its own isolate. A SharedArrayBuffer written by another worker mid-call would // break that assumption, which is inherent to processing JS-owned buffers in place. let samples = unsafe { audio.as_mut() }; + let inner = self + .inner + .as_mut() + .ok_or_else(|| disposed_error("Processor"))?; - map_err(self.inner.process(samples)) + map_err(inner.process(samples)) } /// Creates a handle for reading and writing this processor's parameters and state. /// /// Each call returns an independent handle onto the same processor. #[napi] - pub fn get_context(&self) -> ProcessorContext { - ProcessorContext { - inner: self.inner.context(), - } + pub fn get_context(&self) -> Result { + let inner = self + .inner + .as_ref() + .ok_or_else(|| disposed_error("Processor"))?; + + Ok(ProcessorContext { + inner: inner.context(), + }) } /// Ends this processor's telemetry session, after which it can no longer process audio. @@ -159,7 +197,11 @@ impl Processor { /// is collected, but GC timing is not guaranteed. May block, so keep it off the audio path. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - map_err(self.inner.terminate_session()) + let inner = self + .inner + .as_mut() + .ok_or_else(|| disposed_error("Processor"))?; + map_err(inner.terminate_session()) } } diff --git a/src/processor_async.rs b/src/processor_async.rs index 7c82d10..72d31a8 100644 --- a/src/processor_async.rs +++ b/src/processor_async.rs @@ -1,16 +1,20 @@ use crate::{ claim_sdk_id, - error::{Result, map_err}, + error::{Result, disposed_error, map_err}, + mem, model::Model, processor::{OtelConfig, ProcessorContext, audio_config}, }; use napi::{ Env, Task, - bindgen_prelude::{AsyncTask, Float32Array}, + bindgen_prelude::{AsyncTask, Float32Array, ObjectFinalize}, }; use napi_derive::napi; -use std::sync::{Arc, Mutex, MutexGuard}; +use std::sync::{ + Arc, Mutex, MutexGuard, + atomic::{AtomicI64, Ordering}, +}; /// The SDK object shared between a binding class and the tasks it spawns. /// @@ -28,6 +32,83 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { .unwrap_or_else(|poisoned| poisoned.into_inner()) } +/// An SDK object whose native lifetime several JS handles share, plus the external-memory +/// claims those handles have made against it. +/// +/// `ProcessorAsync` and `VadAsync` can hand out second JS handles onto the same native +/// object (`withConfig`), and each handle reports the object's footprint at construction +/// and gives it back at finalization. The claims therefore live on the shared object, and +/// are drained exactly once — by whichever comes first, `dispose()` or the first +/// finalizer — so the ledger always balances. +pub(crate) struct Held { + inner: Mutex>, + outstanding: AtomicI64, +} + +impl Held { + pub(crate) fn new(inner: T) -> Self { + Self { + inner: Mutex::new(Some(inner)), + outstanding: AtomicI64::new(0), + } + } + + /// Records `bytes` of external memory claimed by a JS handle onto this object. + pub(crate) fn claim(&self, bytes: i64) { + self.outstanding.fetch_add(bytes, Ordering::SeqCst); + } + + /// Gives back this handle's claim, if the claims are still outstanding. A no-op once + /// `release()` drained the ledger. + pub(crate) fn unclaim(&self, env: Env, bytes: i64) { + let previous = self.outstanding.fetch_sub(bytes, Ordering::SeqCst); + if previous >= bytes { + mem::adjust(env, -bytes); + } else { + // Claims were already drained; undo the underflow so the ledger stays at zero. + self.outstanding.fetch_add(bytes, Ordering::SeqCst); + } + } + + /// Destroys the native object, if it is still live, giving back every outstanding + /// claim exactly once. Idempotent. + /// + /// Called by `dispose()`, which destroys the object regardless of other handles, and + /// by the finalizer of the last surviving handle. + pub(crate) fn release(&self, env: Env) { + if lock(&self.inner).take().is_some() { + mem::adjust(env, -self.outstanding.swap(0, Ordering::SeqCst)); + } + } + + /// Whether the native object is still live. + pub(crate) fn is_live(&self) -> bool { + lock(&self.inner).is_some() + } + + /// Runs `f` with the native object, or fails with the disposed error. + pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { + let mut guard = lock(&self.inner); + let inner = guard.as_mut().ok_or_else(|| disposed_error(class))?; + f(inner) + } +} + +impl Drop for Held { + fn drop(&mut self) { + // Frees the native object when the last `Arc` handle goes away without any + // finalizer having released it — e.g. the last JS handle was collected while a task + // still held a clone. Freeing matters more than exact bookkeeping here: the + // outstanding claim, if any, stays reported, which only makes V8 a little more + // eager for the rest of the process. + match self.inner.get_mut() { + Ok(slot) => slot.take(), + // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. + Err(poisoned) => poisoned.into_inner().take(), + }; + } +} + /// Speech enhancement processor that keeps its work off the main thread. /// /// The same processing as {@link Processor}, but each call returns a promise and runs on @@ -47,9 +128,23 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { /// Raise `UV_THREADPOOL_SIZE` before Node starts to run more streams in parallel. The /// SDK's own `AIC_NUM_THREADS` does not apply here: that variable sizes a rayon pool this /// binding deliberately does not use. -#[napi] +#[napi(custom_finalize)] pub struct ProcessorAsync { - inner: Shared>, + inner: Arc>>, +} + +impl ObjectFinalize for ProcessorAsync { + fn finalize(self, env: Env) -> Result<()> { + if Arc::strong_count(&self.inner) == 1 { + // Last handle onto the native object: destroy it and drain every claim. + self.inner.release(env); + } else { + // Other handles (or in-flight tasks) keep the object alive; only this handle's + // claim goes back. + self.inner.unclaim(env, mem::PROCESSOR_BYTES); + } + Ok(()) + } } #[napi] @@ -62,17 +157,35 @@ impl ProcessorAsync { /// Telemetry follows the runtime environment; pass `otelConfig` to override it for this /// instance. #[napi(constructor)] - pub fn new(model: &Model, license_key: String, otel_config: Option) -> Result { + pub fn new( + env: Env, + model: &Model, + license_key: String, + otel_config: Option, + ) -> Result { + let model_inner = model.sdk()?; claim_sdk_id(); let inner = match otel_config { Some(config) => { - aic_sdk::Processor::with_otel_config(&model.inner, &license_key, &config.into()) + aic_sdk::Processor::with_otel_config(model_inner, &license_key, &config.into()) } - None => aic_sdk::Processor::new(&model.inner, &license_key), + None => aic_sdk::Processor::new(model_inner, &license_key), }; - Ok(Self { - inner: Arc::new(Mutex::new(map_err(inner)?)), - }) + let inner = Arc::new(Held::new(map_err(inner)?)); + inner.claim(mem::PROCESSOR_BYTES); + mem::adjust(env, mem::PROCESSOR_BYTES); + + Ok(Self { inner }) + } + + /// Destroys the native processor immediately, releasing its memory and telemetry + /// session without waiting for garbage collection. + /// + /// Every later method throws; calling `dispose()` again does nothing. The synchronous + /// variant blocks until in-flight work on the libuv pool finishes. + #[napi] + pub fn dispose(&self, env: Env) { + self.inner.release(env); } /// Initializes the processor and resolves to a handle onto it, for chaining off the @@ -168,7 +281,7 @@ impl ProcessorAsync { /// Backs {@link ProcessorAsync#withConfig}. pub struct ProcessorWithConfigTask { - inner: Shared>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -177,10 +290,21 @@ impl Task for ProcessorWithConfigTask { type JsValue = ProcessorAsync; fn compute(&mut self) -> Result<()> { - map_err(lock(&self.inner).initialize(&self.config)) + self.inner.with("ProcessorAsync", |inner| { + map_err(inner.initialize(&self.config)) + }) } - fn resolve(&mut self, _env: Env, _: ()) -> Result { + fn resolve(&mut self, env: Env, _: ()) -> Result { + // A second JS instance wrapping the same processor: it reports the same footprint at + // finalize, so it must claim it here to keep the external-memory ledger balanced. + // Born disposed when the processor was disposed mid-flight, in which case there is + // nothing to claim. + if self.inner.is_live() { + self.inner.claim(mem::PROCESSOR_BYTES); + mem::adjust(env, mem::PROCESSOR_BYTES); + } + Ok(ProcessorAsync { inner: self.inner.clone(), }) @@ -189,7 +313,7 @@ impl Task for ProcessorWithConfigTask { /// Backs {@link ProcessorAsync#initialize}. pub struct ProcessorInitializeTask { - inner: Shared>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -198,7 +322,9 @@ impl Task for ProcessorInitializeTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - map_err(lock(&self.inner).initialize(&self.config)) + self.inner.with("ProcessorAsync", |inner| { + map_err(inner.initialize(&self.config)) + }) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { @@ -208,7 +334,7 @@ impl Task for ProcessorInitializeTask { /// Backs {@link ProcessorAsync#process}. pub struct ProcessorProcessTask { - inner: Shared>, + inner: Arc>>, audio: Vec, } @@ -220,7 +346,9 @@ impl Task for ProcessorProcessTask { // Moved out rather than borrowed so the buffer can be handed to V8 in `resolve` // without another copy. The task is used once, so leaving an empty Vec behind is fine. let mut audio = std::mem::take(&mut self.audio); - map_err(lock(&self.inner).process(&mut audio))?; + self + .inner + .with("ProcessorAsync", |inner| map_err(inner.process(&mut audio)))?; Ok(audio) } @@ -234,7 +362,7 @@ impl Task for ProcessorProcessTask { /// Backs {@link ProcessorAsync#getContext}. pub struct ProcessorContextTask { - inner: Shared>, + inner: Arc>>, } impl Task for ProcessorContextTask { @@ -242,7 +370,9 @@ impl Task for ProcessorContextTask { type JsValue = ProcessorContext; fn compute(&mut self) -> Result { - Ok(lock(&self.inner).context()) + self + .inner + .with("ProcessorAsync", |inner| Ok(inner.context())) } fn resolve(&mut self, _env: Env, context: aic_sdk::ProcessorContext) -> Result { @@ -252,7 +382,7 @@ impl Task for ProcessorContextTask { /// Backs {@link ProcessorAsync#terminateSession}. pub struct ProcessorTerminateTask { - inner: Shared>, + inner: Arc>>, } impl Task for ProcessorTerminateTask { @@ -260,7 +390,9 @@ impl Task for ProcessorTerminateTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - map_err(lock(&self.inner).terminate_session()) + self + .inner + .with("ProcessorAsync", |inner| map_err(inner.terminate_session())) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { diff --git a/src/vad.rs b/src/vad.rs index 7b29d5f..8f1be12 100644 --- a/src/vad.rs +++ b/src/vad.rs @@ -1,11 +1,12 @@ use crate::{ claim_sdk_id, - error::{Result, map_err}, + error::{Result, disposed_error, map_err}, + mem, model::Model, processor::{OtelConfig, audio_config}, }; -use napi::bindgen_prelude::Float32Array; +use napi::{Env, bindgen_prelude::Float32Array, bindgen_prelude::ObjectFinalize}; use napi_derive::napi; /// Voice activity detection parameters, all changeable while audio is being processed. @@ -53,9 +54,19 @@ impl From for aic_sdk::VadParameter { /// processor's output: enhancement changes the signal the VAD model expects, and stacks /// the processor's delay onto the prediction. Since `process` leaves its input untouched, /// calling it on the same block before `Processor#process` is enough. -#[napi] +#[napi(custom_finalize)] pub struct Vad { - inner: aic_sdk::Vad<'static>, + inner: Option>, +} + +impl ObjectFinalize for Vad { + fn finalize(self, env: Env) -> Result<()> { + // `dispose()` already gave the footprint back when the inner is gone. + if self.inner.is_some() { + mem::adjust(env, -mem::PROCESSOR_BYTES); + } + Ok(()) + } } #[napi] @@ -65,15 +76,33 @@ impl Vad { /// Telemetry follows the runtime environment; pass `otelConfig` to override it for this /// instance. #[napi(constructor)] - pub fn new(model: &Model, license_key: String, otel_config: Option) -> Result { + pub fn new( + env: Env, + model: &Model, + license_key: String, + otel_config: Option, + ) -> Result { + let model_inner = model.sdk()?; claim_sdk_id(); let inner = match otel_config { - Some(config) => aic_sdk::Vad::with_otel_config(&model.inner, &license_key, &config.into()), - None => aic_sdk::Vad::new(&model.inner, &license_key), + Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), + None => aic_sdk::Vad::new(model_inner, &license_key), }; - Ok(Self { - inner: map_err(inner)?, - }) + let inner = map_err(inner)?; + mem::adjust(env, mem::PROCESSOR_BYTES); + + Ok(Self { inner: Some(inner) }) + } + + /// Destroys the native VAD immediately, releasing its memory and telemetry session + /// without waiting for garbage collection. + /// + /// Every later method throws; calling `dispose()` again does nothing. + #[napi] + pub fn dispose(&mut self, env: Env) { + if self.inner.take().is_some() { + mem::adjust(env, -mem::PROCESSOR_BYTES); + } } /// Configures the VAD for an audio format. Must be called before processing. @@ -87,11 +116,8 @@ impl Vad { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err( - self - .inner - .initialize(&audio_config(sample_rate, block_size, variable_block_size)), - ) + let inner = self.inner.as_mut().ok_or_else(|| disposed_error("Vad"))?; + map_err(inner.initialize(&audio_config(sample_rate, block_size, variable_block_size))) } /// Examines a mono audio block and updates the prediction, leaving the audio unmodified. @@ -99,23 +125,27 @@ impl Vad { pub fn process(&mut self, audio: Float32Array) -> Result<()> { // Read-only, so the safe `Deref` to `&[f32]` is enough here. Taking the view by // value does not copy the caller's samples. - map_err(self.inner.process(&audio)) + let inner = self.inner.as_mut().ok_or_else(|| disposed_error("Vad"))?; + map_err(inner.process(&audio)) } /// Creates a handle for reading predictions and controlling this VAD. /// /// Each call returns an independent handle onto the same VAD. #[napi] - pub fn get_context(&self) -> VadContext { - VadContext { - inner: self.inner.context(), - } + pub fn get_context(&self) -> Result { + let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Vad"))?; + + Ok(VadContext { + inner: inner.context(), + }) } /// Ends this VAD's telemetry session, after which it can no longer process audio. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - map_err(self.inner.terminate_session()) + let inner = self.inner.as_mut().ok_or_else(|| disposed_error("Vad"))?; + map_err(inner.terminate_session()) } } diff --git a/src/vad_async.rs b/src/vad_async.rs index c289e7b..36021a6 100644 --- a/src/vad_async.rs +++ b/src/vad_async.rs @@ -1,18 +1,19 @@ use crate::{ claim_sdk_id, error::{Result, map_err}, + mem, model::Model, processor::{OtelConfig, audio_config}, - processor_async::{Shared, lock}, + processor_async::Held, vad::VadContext, }; use napi::{ Env, Task, - bindgen_prelude::{AsyncTask, Float32Array}, + bindgen_prelude::{AsyncTask, Float32Array, ObjectFinalize}, }; use napi_derive::napi; -use std::sync::{Arc, Mutex}; +use std::sync::Arc; /// Voice activity detector that keeps its work off the main thread. /// @@ -36,9 +37,23 @@ use std::sync::{Arc, Mutex}; /// Raise `UV_THREADPOOL_SIZE` before Node starts to run more streams in parallel. The /// SDK's own `AIC_NUM_THREADS` does not apply here: that variable sizes a rayon pool this /// binding deliberately does not use. -#[napi] +#[napi(custom_finalize)] pub struct VadAsync { - inner: Shared>, + inner: Arc>>, +} + +impl ObjectFinalize for VadAsync { + fn finalize(self, env: Env) -> Result<()> { + if Arc::strong_count(&self.inner) == 1 { + // Last handle onto the native object: destroy it and drain every claim. + self.inner.release(env); + } else { + // Other handles (or in-flight tasks) keep the object alive; only this handle's + // claim goes back. + self.inner.unclaim(env, mem::PROCESSOR_BYTES); + } + Ok(()) + } } #[napi] @@ -51,15 +66,33 @@ impl VadAsync { /// Telemetry follows the runtime environment; pass `otelConfig` to override it for this /// instance. #[napi(constructor)] - pub fn new(model: &Model, license_key: String, otel_config: Option) -> Result { + pub fn new( + env: Env, + model: &Model, + license_key: String, + otel_config: Option, + ) -> Result { + let model_inner = model.sdk()?; claim_sdk_id(); let inner = match otel_config { - Some(config) => aic_sdk::Vad::with_otel_config(&model.inner, &license_key, &config.into()), - None => aic_sdk::Vad::new(&model.inner, &license_key), + Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), + None => aic_sdk::Vad::new(model_inner, &license_key), }; - Ok(Self { - inner: Arc::new(Mutex::new(map_err(inner)?)), - }) + let inner = Arc::new(Held::new(map_err(inner)?)); + inner.claim(mem::PROCESSOR_BYTES); + mem::adjust(env, mem::PROCESSOR_BYTES); + + Ok(Self { inner }) + } + + /// Destroys the native VAD immediately, releasing its memory and telemetry session + /// without waiting for garbage collection. + /// + /// Every later method throws; calling `dispose()` again does nothing. The synchronous + /// variant blocks until in-flight work on the libuv pool finishes. + #[napi] + pub fn dispose(&self, env: Env) { + self.inner.release(env); } /// Initializes the VAD and resolves to a handle onto it, for chaining off the @@ -153,7 +186,7 @@ impl VadAsync { /// Backs {@link VadAsync#withConfig}. pub struct VadWithConfigTask { - inner: Shared>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -162,10 +195,21 @@ impl Task for VadWithConfigTask { type JsValue = VadAsync; fn compute(&mut self) -> Result<()> { - map_err(lock(&self.inner).initialize(&self.config)) + self + .inner + .with("VadAsync", |inner| map_err(inner.initialize(&self.config))) } - fn resolve(&mut self, _env: Env, _: ()) -> Result { + fn resolve(&mut self, env: Env, _: ()) -> Result { + // A second JS instance wrapping the same VAD: it reports the same footprint at + // finalize, so it must claim it here to keep the external-memory ledger balanced. + // Born disposed when the VAD was disposed mid-flight, in which case there is + // nothing to claim. + if self.inner.is_live() { + self.inner.claim(mem::PROCESSOR_BYTES); + mem::adjust(env, mem::PROCESSOR_BYTES); + } + Ok(VadAsync { inner: self.inner.clone(), }) @@ -174,7 +218,7 @@ impl Task for VadWithConfigTask { /// Backs {@link VadAsync#initialize}. pub struct VadInitializeTask { - inner: Shared>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -183,7 +227,9 @@ impl Task for VadInitializeTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - map_err(lock(&self.inner).initialize(&self.config)) + self + .inner + .with("VadAsync", |inner| map_err(inner.initialize(&self.config))) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { @@ -193,7 +239,7 @@ impl Task for VadInitializeTask { /// Backs {@link VadAsync#process}. pub struct VadProcessTask { - inner: Shared>, + inner: Arc>>, audio: Vec, } @@ -205,7 +251,9 @@ impl Task for VadProcessTask { // Moved out rather than borrowed so the buffer can be handed to V8 in `resolve` // without another copy. The task is used once, so leaving an empty Vec behind is fine. let audio = std::mem::take(&mut self.audio); - map_err(lock(&self.inner).process(&audio))?; + self + .inner + .with("VadAsync", |inner| map_err(inner.process(&audio)))?; Ok(audio) } @@ -219,7 +267,7 @@ impl Task for VadProcessTask { /// Backs {@link VadAsync#getContext}. pub struct VadContextTask { - inner: Shared>, + inner: Arc>>, } impl Task for VadContextTask { @@ -227,7 +275,7 @@ impl Task for VadContextTask { type JsValue = VadContext; fn compute(&mut self) -> Result { - Ok(lock(&self.inner).context()) + self.inner.with("VadAsync", |inner| Ok(inner.context())) } fn resolve(&mut self, _env: Env, context: aic_sdk::VadContext) -> Result { @@ -237,7 +285,7 @@ impl Task for VadContextTask { /// Backs {@link VadAsync#terminateSession}. pub struct VadTerminateTask { - inner: Shared>, + inner: Arc>>, } impl Task for VadTerminateTask { @@ -245,7 +293,9 @@ impl Task for VadTerminateTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - map_err(lock(&self.inner).terminate_session()) + self + .inner + .with("VadAsync", |inner| map_err(inner.terminate_session())) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { From 9dd5cd98f3e63c993cdc97535ebf17b8586fba9a Mon Sep 17 00:00:00 2001 From: steckes Date: Thu, 3 Sep 2026 15:01:06 +0200 Subject: [PATCH 2/7] fixes --- .github/workflows/CI.yml | 8 +-- .gitignore | 8 +-- __test__/dispose.spec.ts | 119 ++++++++++++++++++++++++++++++++++++--- index.d.ts | 15 +++-- src/analyzer.rs | 13 ++++- src/model.rs | 15 ++--- src/processor.rs | 56 ++++++++++-------- src/processor_async.rs | 12 ++-- src/vad.rs | 31 ++++++---- src/vad_async.rs | 6 +- 10 files changed, 209 insertions(+), 74 deletions(-) diff --git a/.github/workflows/CI.yml b/.github/workflows/CI.yml index 9609e26..0c35a9c 100644 --- a/.github/workflows/CI.yml +++ b/.github/workflows/CI.yml @@ -85,7 +85,7 @@ jobs: # assembles hand-written ARM asm. These images carry a complete sysroot, gcc and # clang, and their older glibc keeps the published binaries widely usable. # - # Both images are amd64 — the aarch64 one bundles a cross toolchain — so the + # Both images are amd64 (the aarch64 one bundles a cross toolchain), so the # host-installed node_modules stay compatible inside the container. - host: ubuntu-latest target: x86_64-unknown-linux-gnu @@ -138,9 +138,9 @@ jobs: run: pnpm install # The build downloads the matching prebuilt native SDK for the target, so these steps # need network access. - # The image ships Rust 1.82, which predates edition 2024 — required by this crate and - # by `aic-sdk` and `aic-sdk-sys` themselves, so lowering our own edition would not - # help. The toolchain is updated in the container before building. + # The image ships Rust 1.82, which predates edition 2024. That edition is required + # by this crate and by `aic-sdk` and `aic-sdk-sys` themselves, so lowering our own + # would not help. The toolchain is updated in the container before building. # # The workspace is mounted at its own path rather than somewhere like /build, so the # absolute paths in pnpm's node_modules symlinks still resolve. diff --git a/.gitignore b/.gitignore index 0546ba4..778b085 100644 --- a/.gitignore +++ b/.gitignore @@ -6,10 +6,10 @@ debug/ target/ # Cargo.lock is deliberately committed, contrary to the template default below. This crate -# is a cdylib published as prebuilt binaries — an artifact, not a library other crates -# compile against — so its builds should be reproducible. Leaving it out meant every CI run -# re-resolved the whole Rust graph from crates.io, so an upstream release could break the -# build with no change to this repository. +# is a cdylib published as prebuilt binaries, an artifact rather than a library other +# crates compile against, so its builds should be reproducible. Leaving it out meant every +# CI run re-resolved the whole Rust graph from crates.io, so an upstream release could +# break the build with no change to this repository. # # Remove Cargo.lock from gitignore if creating an executable, leave it for libraries # More information here https://doc.rust-lang.org/cargo/guide/cargo-toml-vs-cargo-lock.html diff --git a/__test__/dispose.spec.ts b/__test__/dispose.spec.ts index 8994e2f..9b4cc06 100644 --- a/__test__/dispose.spec.ts +++ b/__test__/dispose.spec.ts @@ -1,10 +1,42 @@ import test from 'ava' -import { Model, Processor, ProcessorAsync, Vad, VadAsync } from '../index.js' +import { Analyzer, Model, Processor, ProcessorAsync, Vad, VadAsync } from '../index.js' import { licenseKey, modelPath } from './common.js' const enhancementModel = () => Model.fromFile(modelPath('enhancement')) const vadModel = () => Model.fromFile(modelPath('vad')) +const analysisModel = () => Model.fromFile(modelPath('analysis')) + +/** An initialized analyzer with a span of audio already buffered. */ +function bufferedAnalyzer() { + const model = analysisModel() + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + + const analyzer = new Analyzer(model, licenseKey()) + analyzer.initialize(sampleRate, blockSize) + + const audio = Float32Array.from({ length: blockSize }, (_, i) => Math.sin(i / 10) * 0.5) + for (let block = 0; block < 100; block += 1) { + analyzer.buffer(audio) + } + + return { analyzer, audio, blockSize } +} + +/** Milliseconds since a `process.hrtime.bigint()` reading. */ +function elapsedMs(since: bigint): number { + return Number(process.hrtime.bigint() - since) / 1e6 +} + +/** Holds the calling thread for `ms`, leaving worker threads free to run. */ +function spin(ms: number) { + const until = process.hrtime.bigint() + BigInt(Math.round(ms * 1e6)) + while (process.hrtime.bigint() < until) { + // Deliberately busy: yielding here would hand the event loop back and a timer would + // overshoot, either of which lets the analysis finish before dispose() is called. + } +} /** An initialized processor at the model's optimal settings, plus that block size. */ function initializedProcessor() { @@ -40,11 +72,16 @@ test('sync vad dispose releases memory and rejects later use', (t) => { const model = vadModel() const vad = new Vad(model, licenseKey()) vad.initialize(model.getOptimalSampleRate(), model.getOptimalBlockSize(model.getOptimalSampleRate())) + const context = vad.getContext() vad.dispose() t.throws(() => vad.process(new Float32Array(160)), { message: /disposed/ }) t.throws(() => vad.getContext(), { message: /disposed/ }) vad.dispose() + + // The SDK guarantees a VAD context stays valid after its VAD is destroyed + // (the prediction just stops updating), so using it must not crash. + context.isSpeechDetected() t.pass() }) @@ -76,13 +113,13 @@ test('async processor dispose rejects later use and is idempotent', async (t) => const processor = new ProcessorAsync(model, licenseKey()) await processor.initialize(sampleRate, blockSize) - await processor.dispose() + processor.dispose() await t.throwsAsync(() => processor.process(new Float32Array(blockSize)), { message: /disposed/, }) await t.throwsAsync(() => processor.terminateSession(), { message: /disposed/ }) - await processor.dispose() + processor.dispose() t.pass() }) @@ -93,12 +130,81 @@ test('async vad dispose rejects later use', async (t) => { const vad = new VadAsync(model, licenseKey()) await vad.initialize(sampleRate, blockSize) - await vad.dispose() + vad.dispose() await t.throwsAsync(() => vad.process(new Float32Array(blockSize)), { message: /disposed/ }) t.pass() }) +test('analyzer dispose rejects later use and is idempotent', async (t) => { + const model = analysisModel() + const sampleRate = model.getOptimalSampleRate() + const blockSize = model.getOptimalBlockSize(sampleRate) + + const analyzer = new Analyzer(model, licenseKey()) + analyzer.initialize(sampleRate, blockSize) + analyzer.buffer(new Float32Array(blockSize)) + analyzer.dispose() + + t.throws(() => analyzer.buffer(new Float32Array(blockSize)), { message: /disposed/ }) + t.throws(() => analyzer.analyze(), { message: /disposed/ }) + await t.throwsAsync(() => analyzer.analyzeAsync(), { message: /disposed/ }) + + // Dispose is idempotent: a second call does nothing and must not throw. + analyzer.dispose() + t.pass() +}) + +test('analyzer dispose racing an analyzeAsync settles it without crashing', async (t) => { + const { analyzer } = bufferedAnalyzer() + + // Disposing in the same tick leaves who reaches the lock first genuinely undecided, so + // the promise may resolve or reject. Neither outcome is asserted here; that the race is + // resolved safely at all is the point, and the blocking path gets its own test below. + const inFlight = analyzer.analyzeAsync().catch(() => {}) + analyzer.dispose() + await inFlight + + t.throws(() => analyzer.analyze(), { message: /disposed/ }) + t.pass() +}) + +test('analyzer dispose waits for an in-flight analyzeAsync to finish', async (t) => { + const { analyzer } = bufferedAnalyzer() + + // Time a warm analysis so the thresholds below scale to this machine rather than to a + // guessed constant. The first run pays any lazy setup, so the second is the estimate. + await analyzer.analyzeAsync() + const calibrationStarted = process.hrtime.bigint() + await analyzer.analyzeAsync() + const analysisMs = elapsedMs(calibrationStarted) + + // Hand the task to libuv and let the worker take the analyzer lock. The delay is spun + // rather than timed: a timer's granularity could overshoot the whole analysis, which + // would leave dispose() with nothing to wait for and quietly void the test. + const inFlight = analyzer.analyzeAsync() + await new Promise((resolve) => setImmediate(resolve)) + spin(Math.min(analysisMs / 8, 5)) + + const disposeStarted = process.hrtime.bigint() + analyzer.dispose() + const blockedMs = elapsedMs(disposeStarted) + + // dispose() waited for the lock instead of pulling the analyzer out from under the + // worker, so the analysis ran to completion and produced a real result. + const result = await inFlight + t.is(typeof result.riskScore, 'number', 'the in-flight analysis must complete, not be cancelled') + + // And it was still running when dispose() was called, so that wait was the lock being + // held rather than a no-op on an already-finished analysis. + t.true( + blockedMs >= analysisMs / 4, + `dispose() blocked ${blockedMs.toFixed(1)}ms of a ~${analysisMs.toFixed(1)}ms analysis`, + ) + + t.throws(() => analyzer.analyze(), { message: /disposed/ }) +}) + test('a withConfig handle outlives its original handle being collected', async (t) => { const model = enhancementModel() const sampleRate = model.getOptimalSampleRate() @@ -106,7 +212,7 @@ test('a withConfig handle outlives its original handle being collected', async ( // The chaining form leaves the constructor's handle as garbage. Its finalizer must // give back only its own claim, not destroy the shared native processor. - let processor = await new ProcessorAsync(model, licenseKey()).withConfig(sampleRate, blockSize) + const processor = await new ProcessorAsync(model, licenseKey()).withConfig(sampleRate, blockSize) // Encourage a GC so the dropped handle's finalizer runs. const churn: unknown[] = [] @@ -117,7 +223,6 @@ test('a withConfig handle outlives its original handle being collected', async ( const block = new Float32Array(blockSize) await processor.process(block) - await processor.dispose() - processor = undefined as unknown as ProcessorAsync + processor.dispose() t.pass() }) diff --git a/index.d.ts b/index.d.ts index 3c50fe5..442e299 100644 --- a/index.d.ts +++ b/index.d.ts @@ -226,8 +226,8 @@ export declare class ProcessorAsync { * Destroys the native processor immediately, releasing its memory and telemetry * session without waiting for garbage collection. * - * Every later method throws; calling `dispose()` again does nothing. The synchronous - * variant blocks until in-flight work on the libuv pool finishes. + * Every later method throws; calling `dispose()` again does nothing. Blocks until + * in-flight work on the libuv pool finishes. */ dispose(): void /** @@ -285,8 +285,11 @@ export declare class ProcessorAsync { /** * Control handle for a {@link Processor}. * - * Every method may be called while audio is being processed. Releasing the handle does - * not destroy the processor it came from. + * Every method may be called while audio is being processed. The handle and the processor + * have independent lifetimes in both directions: releasing the handle does not destroy + * the processor it came from, and the handle stays valid after its processor is disposed + * or garbage-collected. Calls on it keep succeeding; they just no longer reach a live + * processor. */ export declare class ProcessorContext { /** Sets an enhancement parameter. Throws if the value is out of range. */ @@ -402,8 +405,8 @@ export declare class VadAsync { * Destroys the native VAD immediately, releasing its memory and telemetry session * without waiting for garbage collection. * - * Every later method throws; calling `dispose()` again does nothing. The synchronous - * variant blocks until in-flight work on the libuv pool finishes. + * Every later method throws; calling `dispose()` again does nothing. Blocks until + * in-flight work on the libuv pool finishes. */ dispose(): void /** diff --git a/src/analyzer.rs b/src/analyzer.rs index 190eb16..21411a7 100644 --- a/src/analyzer.rs +++ b/src/analyzer.rs @@ -83,6 +83,11 @@ pub struct Analyzer { impl ObjectFinalize for Analyzer { fn finalize(self, env: Env) -> Result<()> { // `dispose()` already gave the footprint back when the collector is gone. + // + // An in-flight `AnalyzeTask` holds the analyzer `Arc`, so that half is dropped + // only after the worker finishes. The collector drops here, possibly while the + // worker analyzes. That is safe per the C API, which destroys the paired halves + // independently, in any order (`aic_collector_destroy`). if self.collector.is_some() { mem::adjust(env, -mem::ANALYZER_BYTES); } @@ -95,7 +100,7 @@ impl Analyzer { /// Creates an analyzer from an analysis model. Other model types are rejected. #[napi(constructor)] pub fn new(env: Env, model: &Model, license_key: String) -> Result { - let model_inner = model.sdk()?; + let model_inner = model.live()?; claim_sdk_id(); let (collector, analyzer) = map_err(aic_sdk::analyzer_pair(model_inner, &license_key))?; mem::adjust(env, mem::ANALYZER_BYTES); @@ -114,6 +119,12 @@ impl Analyzer { #[napi] pub fn dispose(&mut self, env: Env) { if self.collector.take().is_some() { + // The collector drops before the analyzer lock is taken, so it can be destroyed + // while an `analyzeAsync` is in flight on a worker. That is safe per the C API + // (`aic_collector_destroy`): the paired halves are destroyed independently, in + // any order, and the collector handle itself is only ever used on this thread. + // The analyzer half is destroyed under the lock, which is what blocks until the + // in-flight analysis finishes. lock(&self.analyzer).take(); // dropped on scope exit mem::adjust(env, -mem::ANALYZER_BYTES); } diff --git a/src/model.rs b/src/model.rs index e846a19..f0acb78 100644 --- a/src/model.rs +++ b/src/model.rs @@ -15,7 +15,7 @@ use napi_derive::napi; pub struct Model { // `from_file` memory-maps the file rather than borrowing a caller-owned buffer, so // the SDK model is `'static` and needs no lifetime plumbing here. - pub(crate) inner: Option>, + inner: Option>, /// Native footprint reported to V8's GC while this instance is alive (the mmap'd /// weights). Reported back on dispose or finalize so the accounting balances. reported_bytes: i64, @@ -32,8 +32,8 @@ impl ObjectFinalize for Model { } impl Model { - /// The SDK model, or the disposed error once `dispose()` ran. - pub(crate) fn sdk(&self) -> std::result::Result<&aic_sdk::Model<'static>, napi::Error> { + /// The live SDK model, or the disposed error once `dispose()` ran. + pub(crate) fn live(&self) -> Result<&aic_sdk::Model<'static>> { self.inner.as_ref().ok_or_else(|| disposed_error("Model")) } } @@ -93,8 +93,7 @@ impl Model { /// The model identifier, e.g. `quail-vf-2.2-s-16khz`. #[napi] pub fn get_id(&self) -> Result { - let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Model"))?; - Ok(inner.id().to_owned()) + Ok(self.live()?.id().to_owned()) } /// The sample rate in Hz the model was trained for. @@ -103,8 +102,7 @@ impl Model { /// its own Nyquist limit, so matching this rate gives the best quality. #[napi] pub fn get_optimal_sample_rate(&self) -> Result { - let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Model"))?; - Ok(inner.optimal_sample_rate()) + Ok(self.live()?.optimal_sample_rate()) } /// The block size that avoids internal buffering at `sampleRate`. @@ -114,12 +112,11 @@ impl Model { /// window: a 10 ms window is 480 samples at 48 kHz but 160 at 16 kHz. #[napi] pub fn get_optimal_block_size(&self, sample_rate: u32) -> Result { - let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Model"))?; // The SDK reports sizes as `usize`, which napi would marshal as a JS BigInt. // A BigInt block size would throw on `new Float32Array(n)` and on arithmetic // against plain numbers, so it crosses the boundary as u32. Block sizes are a // few thousand samples at most. - Ok(inner.optimal_block_size(sample_rate) as u32) + Ok(self.live()?.optimal_block_size(sample_rate) as u32) } } diff --git a/src/processor.rs b/src/processor.rs index 7154e91..e550687 100644 --- a/src/processor.rs +++ b/src/processor.rs @@ -94,6 +94,24 @@ impl ObjectFinalize for Processor { } } +impl Processor { + /// The live SDK processor, or the disposed error once `dispose()` ran. + fn live(&self) -> Result<&aic_sdk::Processor<'static>> { + self + .inner + .as_ref() + .ok_or_else(|| disposed_error("Processor")) + } + + /// The same, for a `&mut` call. + fn live_mut(&mut self) -> Result<&mut aic_sdk::Processor<'static>> { + self + .inner + .as_mut() + .ok_or_else(|| disposed_error("Processor")) + } +} + #[napi] impl Processor { /// Creates a processor from an enhancement or bypass model. @@ -107,7 +125,7 @@ impl Processor { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.sdk()?; + let model_inner = model.live()?; claim_sdk_id(); let inner = match otel_config { Some(config) => { @@ -146,11 +164,11 @@ impl Processor { block_size: u32, variable_block_size: Option, ) -> Result<()> { - let inner = self - .inner - .as_mut() - .ok_or_else(|| disposed_error("Processor"))?; - map_err(inner.initialize(&audio_config(sample_rate, block_size, variable_block_size))) + map_err(self.live_mut()?.initialize(&audio_config( + sample_rate, + block_size, + variable_block_size, + ))) } /// Enhances a mono audio block in place. @@ -168,12 +186,8 @@ impl Processor { // its own isolate. A SharedArrayBuffer written by another worker mid-call would // break that assumption, which is inherent to processing JS-owned buffers in place. let samples = unsafe { audio.as_mut() }; - let inner = self - .inner - .as_mut() - .ok_or_else(|| disposed_error("Processor"))?; - map_err(inner.process(samples)) + map_err(self.live_mut()?.process(samples)) } /// Creates a handle for reading and writing this processor's parameters and state. @@ -181,13 +195,8 @@ impl Processor { /// Each call returns an independent handle onto the same processor. #[napi] pub fn get_context(&self) -> Result { - let inner = self - .inner - .as_ref() - .ok_or_else(|| disposed_error("Processor"))?; - Ok(ProcessorContext { - inner: inner.context(), + inner: self.live()?.context(), }) } @@ -197,18 +206,17 @@ impl Processor { /// is collected, but GC timing is not guaranteed. May block, so keep it off the audio path. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - let inner = self - .inner - .as_mut() - .ok_or_else(|| disposed_error("Processor"))?; - map_err(inner.terminate_session()) + map_err(self.live_mut()?.terminate_session()) } } /// Control handle for a {@link Processor}. /// -/// Every method may be called while audio is being processed. Releasing the handle does -/// not destroy the processor it came from. +/// Every method may be called while audio is being processed. The handle and the processor +/// have independent lifetimes in both directions: releasing the handle does not destroy +/// the processor it came from, and the handle stays valid after its processor is disposed +/// or garbage-collected. Calls on it keep succeeding; they just no longer reach a live +/// processor. #[napi] pub struct ProcessorContext { pub(crate) inner: aic_sdk::ProcessorContext, diff --git a/src/processor_async.rs b/src/processor_async.rs index 72d31a8..5cd610e 100644 --- a/src/processor_async.rs +++ b/src/processor_async.rs @@ -38,8 +38,8 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { /// `ProcessorAsync` and `VadAsync` can hand out second JS handles onto the same native /// object (`withConfig`), and each handle reports the object's footprint at construction /// and gives it back at finalization. The claims therefore live on the shared object, and -/// are drained exactly once — by whichever comes first, `dispose()` or the first -/// finalizer — so the ledger always balances. +/// are drained exactly once: by `dispose()`, or by the last handle's finalizer, whichever +/// comes first. The ledger always balances. pub(crate) struct Held { inner: Mutex>, outstanding: AtomicI64, @@ -97,7 +97,7 @@ impl Held { impl Drop for Held { fn drop(&mut self) { // Frees the native object when the last `Arc` handle goes away without any - // finalizer having released it — e.g. the last JS handle was collected while a task + // finalizer having released it, e.g. the last JS handle was collected while a task // still held a clone. Freeing matters more than exact bookkeeping here: the // outstanding claim, if any, stays reported, which only makes V8 a little more // eager for the rest of the process. @@ -163,7 +163,7 @@ impl ProcessorAsync { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.sdk()?; + let model_inner = model.live()?; claim_sdk_id(); let inner = match otel_config { Some(config) => { @@ -181,8 +181,8 @@ impl ProcessorAsync { /// Destroys the native processor immediately, releasing its memory and telemetry /// session without waiting for garbage collection. /// - /// Every later method throws; calling `dispose()` again does nothing. The synchronous - /// variant blocks until in-flight work on the libuv pool finishes. + /// Every later method throws; calling `dispose()` again does nothing. Blocks until + /// in-flight work on the libuv pool finishes. #[napi] pub fn dispose(&self, env: Env) { self.inner.release(env); diff --git a/src/vad.rs b/src/vad.rs index 8f1be12..8da67ce 100644 --- a/src/vad.rs +++ b/src/vad.rs @@ -69,6 +69,18 @@ impl ObjectFinalize for Vad { } } +impl Vad { + /// The live SDK VAD, or the disposed error once `dispose()` ran. + fn live(&self) -> Result<&aic_sdk::Vad<'static>> { + self.inner.as_ref().ok_or_else(|| disposed_error("Vad")) + } + + /// The same, for a `&mut` call. + fn live_mut(&mut self) -> Result<&mut aic_sdk::Vad<'static>> { + self.inner.as_mut().ok_or_else(|| disposed_error("Vad")) + } +} + #[napi] impl Vad { /// Creates a voice activity detector from a dedicated VAD model. @@ -82,7 +94,7 @@ impl Vad { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.sdk()?; + let model_inner = model.live()?; claim_sdk_id(); let inner = match otel_config { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), @@ -116,8 +128,11 @@ impl Vad { block_size: u32, variable_block_size: Option, ) -> Result<()> { - let inner = self.inner.as_mut().ok_or_else(|| disposed_error("Vad"))?; - map_err(inner.initialize(&audio_config(sample_rate, block_size, variable_block_size))) + map_err(self.live_mut()?.initialize(&audio_config( + sample_rate, + block_size, + variable_block_size, + ))) } /// Examines a mono audio block and updates the prediction, leaving the audio unmodified. @@ -125,8 +140,7 @@ impl Vad { pub fn process(&mut self, audio: Float32Array) -> Result<()> { // Read-only, so the safe `Deref` to `&[f32]` is enough here. Taking the view by // value does not copy the caller's samples. - let inner = self.inner.as_mut().ok_or_else(|| disposed_error("Vad"))?; - map_err(inner.process(&audio)) + map_err(self.live_mut()?.process(&audio)) } /// Creates a handle for reading predictions and controlling this VAD. @@ -134,18 +148,15 @@ impl Vad { /// Each call returns an independent handle onto the same VAD. #[napi] pub fn get_context(&self) -> Result { - let inner = self.inner.as_ref().ok_or_else(|| disposed_error("Vad"))?; - Ok(VadContext { - inner: inner.context(), + inner: self.live()?.context(), }) } /// Ends this VAD's telemetry session, after which it can no longer process audio. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - let inner = self.inner.as_mut().ok_or_else(|| disposed_error("Vad"))?; - map_err(inner.terminate_session()) + map_err(self.live_mut()?.terminate_session()) } } diff --git a/src/vad_async.rs b/src/vad_async.rs index 36021a6..1e5b0e8 100644 --- a/src/vad_async.rs +++ b/src/vad_async.rs @@ -72,7 +72,7 @@ impl VadAsync { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.sdk()?; + let model_inner = model.live()?; claim_sdk_id(); let inner = match otel_config { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), @@ -88,8 +88,8 @@ impl VadAsync { /// Destroys the native VAD immediately, releasing its memory and telemetry session /// without waiting for garbage collection. /// - /// Every later method throws; calling `dispose()` again does nothing. The synchronous - /// variant blocks until in-flight work on the libuv pool finishes. + /// Every later method throws; calling `dispose()` again does nothing. Blocks until + /// in-flight work on the libuv pool finishes. #[napi] pub fn dispose(&self, env: Env) { self.inner.release(env); From 4a04faaaff6c7f39dbae2e993d91e62935fad0ad Mon Sep 17 00:00:00 2001 From: steckes Date: Thu, 3 Sep 2026 15:48:02 +0200 Subject: [PATCH 3/7] improvements --- CHANGELOG.md | 12 ++++++ src/analyzer.rs | 2 +- src/model.rs | 10 ++--- src/processor.rs | 16 ++++---- src/processor_async.rs | 90 +++++++++++++----------------------------- src/vad.rs | 16 ++++---- src/vad_async.rs | 29 +++++--------- 7 files changed, 71 insertions(+), 104 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8374ad6..ab86635 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,18 @@ the binding and its type declarations are now generated from annotated Rust. at once. - `Model.download` is asynchronous and resolves to the model path, so a cold download no longer blocks the event loop. +- **`dispose()` on every class holding native resources** (`Model`, `Processor`, + `ProcessorAsync`, `Vad`, `VadAsync`, `Analyzer`), destroying the native object + immediately instead of waiting for garbage collection. Every later method throws and a + repeat `dispose()` does nothing; where work can still be in flight (the async classes, + or an `Analyzer` mid-`analyzeAsync`) it blocks until that work finishes. + `Model.dispose()` unmaps the model file while objects created from it keep working. + `ProcessorContext` and `VadContext` handles stay valid after their object is disposed; + calls on them just no longer reach a live one. +- **Native footprints are reported to V8's garbage collector** (`napi_adjust_external_memory`). + These classes hold large native allocations behind small JavaScript objects; without the + report the collector felt no pressure to reclaim dropped instances, so native memory + grew unbounded. ### Changed diff --git a/src/analyzer.rs b/src/analyzer.rs index 21411a7..aaf6178 100644 --- a/src/analyzer.rs +++ b/src/analyzer.rs @@ -100,7 +100,7 @@ impl Analyzer { /// Creates an analyzer from an analysis model. Other model types are rejected. #[napi(constructor)] pub fn new(env: Env, model: &Model, license_key: String) -> Result { - let model_inner = model.live()?; + let model_inner = model.inner()?; claim_sdk_id(); let (collector, analyzer) = map_err(aic_sdk::analyzer_pair(model_inner, &license_key))?; mem::adjust(env, mem::ANALYZER_BYTES); diff --git a/src/model.rs b/src/model.rs index f0acb78..53ed91b 100644 --- a/src/model.rs +++ b/src/model.rs @@ -32,8 +32,8 @@ impl ObjectFinalize for Model { } impl Model { - /// The live SDK model, or the disposed error once `dispose()` ran. - pub(crate) fn live(&self) -> Result<&aic_sdk::Model<'static>> { + /// The inner SDK model, or the disposed error once `dispose()` ran. + pub(crate) fn inner(&self) -> Result<&aic_sdk::Model<'static>> { self.inner.as_ref().ok_or_else(|| disposed_error("Model")) } } @@ -93,7 +93,7 @@ impl Model { /// The model identifier, e.g. `quail-vf-2.2-s-16khz`. #[napi] pub fn get_id(&self) -> Result { - Ok(self.live()?.id().to_owned()) + Ok(self.inner()?.id().to_owned()) } /// The sample rate in Hz the model was trained for. @@ -102,7 +102,7 @@ impl Model { /// its own Nyquist limit, so matching this rate gives the best quality. #[napi] pub fn get_optimal_sample_rate(&self) -> Result { - Ok(self.live()?.optimal_sample_rate()) + Ok(self.inner()?.optimal_sample_rate()) } /// The block size that avoids internal buffering at `sampleRate`. @@ -116,7 +116,7 @@ impl Model { // A BigInt block size would throw on `new Float32Array(n)` and on arithmetic // against plain numbers, so it crosses the boundary as u32. Block sizes are a // few thousand samples at most. - Ok(self.live()?.optimal_block_size(sample_rate) as u32) + Ok(self.inner()?.optimal_block_size(sample_rate) as u32) } } diff --git a/src/processor.rs b/src/processor.rs index e550687..ed74f24 100644 --- a/src/processor.rs +++ b/src/processor.rs @@ -95,8 +95,8 @@ impl ObjectFinalize for Processor { } impl Processor { - /// The live SDK processor, or the disposed error once `dispose()` ran. - fn live(&self) -> Result<&aic_sdk::Processor<'static>> { + /// The inner SDK processor, or the disposed error once `dispose()` ran. + fn inner(&self) -> Result<&aic_sdk::Processor<'static>> { self .inner .as_ref() @@ -104,7 +104,7 @@ impl Processor { } /// The same, for a `&mut` call. - fn live_mut(&mut self) -> Result<&mut aic_sdk::Processor<'static>> { + fn inner_mut(&mut self) -> Result<&mut aic_sdk::Processor<'static>> { self .inner .as_mut() @@ -125,7 +125,7 @@ impl Processor { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.live()?; + let model_inner = model.inner()?; claim_sdk_id(); let inner = match otel_config { Some(config) => { @@ -164,7 +164,7 @@ impl Processor { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err(self.live_mut()?.initialize(&audio_config( + map_err(self.inner_mut()?.initialize(&audio_config( sample_rate, block_size, variable_block_size, @@ -187,7 +187,7 @@ impl Processor { // break that assumption, which is inherent to processing JS-owned buffers in place. let samples = unsafe { audio.as_mut() }; - map_err(self.live_mut()?.process(samples)) + map_err(self.inner_mut()?.process(samples)) } /// Creates a handle for reading and writing this processor's parameters and state. @@ -196,7 +196,7 @@ impl Processor { #[napi] pub fn get_context(&self) -> Result { Ok(ProcessorContext { - inner: self.live()?.context(), + inner: self.inner()?.context(), }) } @@ -206,7 +206,7 @@ impl Processor { /// is collected, but GC timing is not guaranteed. May block, so keep it off the audio path. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - map_err(self.live_mut()?.terminate_session()) + map_err(self.inner_mut()?.terminate_session()) } } diff --git a/src/processor_async.rs b/src/processor_async.rs index 5cd610e..9b6c0eb 100644 --- a/src/processor_async.rs +++ b/src/processor_async.rs @@ -11,10 +11,7 @@ use napi::{ bindgen_prelude::{AsyncTask, Float32Array, ObjectFinalize}, }; use napi_derive::napi; -use std::sync::{ - Arc, Mutex, MutexGuard, - atomic::{AtomicI64, Ordering}, -}; +use std::sync::{Arc, Mutex, MutexGuard}; /// The SDK object shared between a binding class and the tasks it spawns. /// @@ -32,60 +29,36 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { .unwrap_or_else(|poisoned| poisoned.into_inner()) } -/// An SDK object whose native lifetime several JS handles share, plus the external-memory -/// claims those handles have made against it. +/// An SDK object whose native lifetime several JS handles and in-flight tasks share. /// -/// `ProcessorAsync` and `VadAsync` can hand out second JS handles onto the same native -/// object (`withConfig`), and each handle reports the object's footprint at construction -/// and gives it back at finalization. The claims therefore live on the shared object, and -/// are drained exactly once: by `dispose()`, or by the last handle's finalizer, whichever -/// comes first. The ledger always balances. +/// `ProcessorAsync` and `VadAsync` can hand out a second JS handle onto the same native +/// object (`withConfig`), and tasks on the libuv pool hold `Arc` clones. The footprint is +/// therefore reported to V8 once per object, at construction, and given back exactly +/// once: by `dispose()`, or by the finalizer of the last surviving handle, whichever +/// comes first. `Option::take` makes the give-back idempotent, so the two never +/// double-count. pub(crate) struct Held { inner: Mutex>, - outstanding: AtomicI64, } impl Held { pub(crate) fn new(inner: T) -> Self { Self { inner: Mutex::new(Some(inner)), - outstanding: AtomicI64::new(0), - } - } - - /// Records `bytes` of external memory claimed by a JS handle onto this object. - pub(crate) fn claim(&self, bytes: i64) { - self.outstanding.fetch_add(bytes, Ordering::SeqCst); - } - - /// Gives back this handle's claim, if the claims are still outstanding. A no-op once - /// `release()` drained the ledger. - pub(crate) fn unclaim(&self, env: Env, bytes: i64) { - let previous = self.outstanding.fetch_sub(bytes, Ordering::SeqCst); - if previous >= bytes { - mem::adjust(env, -bytes); - } else { - // Claims were already drained; undo the underflow so the ledger stays at zero. - self.outstanding.fetch_add(bytes, Ordering::SeqCst); } } - /// Destroys the native object, if it is still live, giving back every outstanding - /// claim exactly once. Idempotent. + /// Destroys the native object, if it is still live, and gives its footprint back to + /// V8. Idempotent. /// /// Called by `dispose()`, which destroys the object regardless of other handles, and /// by the finalizer of the last surviving handle. - pub(crate) fn release(&self, env: Env) { + pub(crate) fn release(&self, env: Env, bytes: i64) { if lock(&self.inner).take().is_some() { - mem::adjust(env, -self.outstanding.swap(0, Ordering::SeqCst)); + mem::adjust(env, -bytes); } } - /// Whether the native object is still live. - pub(crate) fn is_live(&self) -> bool { - lock(&self.inner).is_some() - } - /// Runs `f` with the native object, or fails with the disposed error. pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { let mut guard = lock(&self.inner); @@ -96,11 +69,11 @@ impl Held { impl Drop for Held { fn drop(&mut self) { - // Frees the native object when the last `Arc` handle goes away without any - // finalizer having released it, e.g. the last JS handle was collected while a task - // still held a clone. Freeing matters more than exact bookkeeping here: the - // outstanding claim, if any, stays reported, which only makes V8 a little more - // eager for the rest of the process. + // Frees the native object when the last `Arc` goes away without `release` having + // run, e.g. the last JS handle was finalized while a task still held a clone (that + // finalizer saw the extra reference and left the object for the task). The footprint + // report cannot be returned here — that takes the finalizer's `Env` — so those bytes + // stay reported, which only makes V8 a little more eager for the rest of the process. match self.inner.get_mut() { Ok(slot) => slot.take(), // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. @@ -135,13 +108,11 @@ pub struct ProcessorAsync { impl ObjectFinalize for ProcessorAsync { fn finalize(self, env: Env) -> Result<()> { + // Only the last handle onto the native object destroys it; with other handles or + // in-flight tasks holding an `Arc`, this leaves the object (and its footprint + // report) for them. Idempotent against `dispose()`. if Arc::strong_count(&self.inner) == 1 { - // Last handle onto the native object: destroy it and drain every claim. - self.inner.release(env); - } else { - // Other handles (or in-flight tasks) keep the object alive; only this handle's - // claim goes back. - self.inner.unclaim(env, mem::PROCESSOR_BYTES); + self.inner.release(env, mem::PROCESSOR_BYTES); } Ok(()) } @@ -163,7 +134,7 @@ impl ProcessorAsync { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.live()?; + let model_inner = model.inner()?; claim_sdk_id(); let inner = match otel_config { Some(config) => { @@ -172,7 +143,6 @@ impl ProcessorAsync { None => aic_sdk::Processor::new(model_inner, &license_key), }; let inner = Arc::new(Held::new(map_err(inner)?)); - inner.claim(mem::PROCESSOR_BYTES); mem::adjust(env, mem::PROCESSOR_BYTES); Ok(Self { inner }) @@ -185,7 +155,7 @@ impl ProcessorAsync { /// in-flight work on the libuv pool finishes. #[napi] pub fn dispose(&self, env: Env) { - self.inner.release(env); + self.inner.release(env, mem::PROCESSOR_BYTES); } /// Initializes the processor and resolves to a handle onto it, for chaining off the @@ -295,16 +265,10 @@ impl Task for ProcessorWithConfigTask { }) } - fn resolve(&mut self, env: Env, _: ()) -> Result { - // A second JS instance wrapping the same processor: it reports the same footprint at - // finalize, so it must claim it here to keep the external-memory ledger balanced. - // Born disposed when the processor was disposed mid-flight, in which case there is - // nothing to claim. - if self.inner.is_live() { - self.inner.claim(mem::PROCESSOR_BYTES); - mem::adjust(env, mem::PROCESSOR_BYTES); - } - + fn resolve(&mut self, _env: Env, _: ()) -> Result { + // A second JS handle onto the same native processor. The footprint is reported once + // per object at construction, so there is nothing to report here; the last handle's + // finalizer gives it back. Born disposed when the processor was disposed mid-flight. Ok(ProcessorAsync { inner: self.inner.clone(), }) diff --git a/src/vad.rs b/src/vad.rs index 8da67ce..39c68eb 100644 --- a/src/vad.rs +++ b/src/vad.rs @@ -70,13 +70,13 @@ impl ObjectFinalize for Vad { } impl Vad { - /// The live SDK VAD, or the disposed error once `dispose()` ran. - fn live(&self) -> Result<&aic_sdk::Vad<'static>> { + /// The inner SDK VAD, or the disposed error once `dispose()` ran. + fn inner(&self) -> Result<&aic_sdk::Vad<'static>> { self.inner.as_ref().ok_or_else(|| disposed_error("Vad")) } /// The same, for a `&mut` call. - fn live_mut(&mut self) -> Result<&mut aic_sdk::Vad<'static>> { + fn inner_mut(&mut self) -> Result<&mut aic_sdk::Vad<'static>> { self.inner.as_mut().ok_or_else(|| disposed_error("Vad")) } } @@ -94,7 +94,7 @@ impl Vad { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.live()?; + let model_inner = model.inner()?; claim_sdk_id(); let inner = match otel_config { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), @@ -128,7 +128,7 @@ impl Vad { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err(self.live_mut()?.initialize(&audio_config( + map_err(self.inner_mut()?.initialize(&audio_config( sample_rate, block_size, variable_block_size, @@ -140,7 +140,7 @@ impl Vad { pub fn process(&mut self, audio: Float32Array) -> Result<()> { // Read-only, so the safe `Deref` to `&[f32]` is enough here. Taking the view by // value does not copy the caller's samples. - map_err(self.live_mut()?.process(&audio)) + map_err(self.inner_mut()?.process(&audio)) } /// Creates a handle for reading predictions and controlling this VAD. @@ -149,14 +149,14 @@ impl Vad { #[napi] pub fn get_context(&self) -> Result { Ok(VadContext { - inner: self.live()?.context(), + inner: self.inner()?.context(), }) } /// Ends this VAD's telemetry session, after which it can no longer process audio. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - map_err(self.live_mut()?.terminate_session()) + map_err(self.inner_mut()?.terminate_session()) } } diff --git a/src/vad_async.rs b/src/vad_async.rs index 1e5b0e8..d28c5c0 100644 --- a/src/vad_async.rs +++ b/src/vad_async.rs @@ -44,13 +44,11 @@ pub struct VadAsync { impl ObjectFinalize for VadAsync { fn finalize(self, env: Env) -> Result<()> { + // Only the last handle onto the native object destroys it; with other handles or + // in-flight tasks holding an `Arc`, this leaves the object (and its footprint + // report) for them. Idempotent against `dispose()`. if Arc::strong_count(&self.inner) == 1 { - // Last handle onto the native object: destroy it and drain every claim. - self.inner.release(env); - } else { - // Other handles (or in-flight tasks) keep the object alive; only this handle's - // claim goes back. - self.inner.unclaim(env, mem::PROCESSOR_BYTES); + self.inner.release(env, mem::PROCESSOR_BYTES); } Ok(()) } @@ -72,14 +70,13 @@ impl VadAsync { license_key: String, otel_config: Option, ) -> Result { - let model_inner = model.live()?; + let model_inner = model.inner()?; claim_sdk_id(); let inner = match otel_config { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), None => aic_sdk::Vad::new(model_inner, &license_key), }; let inner = Arc::new(Held::new(map_err(inner)?)); - inner.claim(mem::PROCESSOR_BYTES); mem::adjust(env, mem::PROCESSOR_BYTES); Ok(Self { inner }) @@ -92,7 +89,7 @@ impl VadAsync { /// in-flight work on the libuv pool finishes. #[napi] pub fn dispose(&self, env: Env) { - self.inner.release(env); + self.inner.release(env, mem::PROCESSOR_BYTES); } /// Initializes the VAD and resolves to a handle onto it, for chaining off the @@ -200,16 +197,10 @@ impl Task for VadWithConfigTask { .with("VadAsync", |inner| map_err(inner.initialize(&self.config))) } - fn resolve(&mut self, env: Env, _: ()) -> Result { - // A second JS instance wrapping the same VAD: it reports the same footprint at - // finalize, so it must claim it here to keep the external-memory ledger balanced. - // Born disposed when the VAD was disposed mid-flight, in which case there is - // nothing to claim. - if self.inner.is_live() { - self.inner.claim(mem::PROCESSOR_BYTES); - mem::adjust(env, mem::PROCESSOR_BYTES); - } - + fn resolve(&mut self, _env: Env, _: ()) -> Result { + // A second JS handle onto the same native VAD. The footprint is reported once per + // object at construction, so there is nothing to report here; the last handle's + // finalizer gives it back. Born disposed when the VAD was disposed mid-flight. Ok(VadAsync { inner: self.inner.clone(), }) From 5be7a186d8d49ae5e35451a0f2d472788d9bbb8a Mon Sep 17 00:00:00 2001 From: steckes Date: Fri, 4 Sep 2026 10:36:10 +0200 Subject: [PATCH 4/7] pr review --- __test__/index.spec.ts | 2 +- __test__/models.ts | 4 +- src/mem.rs | 92 ++++++++++++++++++++++++++++++++++-------- src/processor_async.rs | 71 +++++--------------------------- src/vad_async.rs | 17 ++++---- 5 files changed, 95 insertions(+), 91 deletions(-) diff --git a/__test__/index.spec.ts b/__test__/index.spec.ts index dc0089e..37a7eb5 100644 --- a/__test__/index.spec.ts +++ b/__test__/index.spec.ts @@ -36,7 +36,7 @@ test('sdk reports a version', (t) => { }) test('sdk expects the model version the fixtures are published under', (t) => { - // Guards against the fixture URLs in models.ts drifting from the addon's expectation, + // Guards against the fixture ids in models.ts drifting from the addon's expectation, // which would otherwise surface as a confusing "model version unsupported" much later. t.is(getCompatibleModelVersion(), TEST_MODEL_VERSION) }) diff --git a/__test__/models.ts b/__test__/models.ts index e56038f..b5208dd 100644 --- a/__test__/models.ts +++ b/__test__/models.ts @@ -17,8 +17,8 @@ export interface TestModel { * `scripts/fetch-test-models.mjs` resolves these through `Model.download()`, which * re-fetches the manifest and pulls the newest compatible model version. The model file * format version is tied to the SDK version, so `getCompatibleModelVersion()` reports - * the version the built addon expects, and `sdk exposes compatible model version` - * asserts the fixtures still match it. + * the version the built addon expects, and the `sdk expects the model version the + * fixtures are published under` test asserts the fixtures still match it. */ export const TEST_MODELS = { enhancement: { id: 'quail-vf-2.2-s-16khz' }, diff --git a/src/mem.rs b/src/mem.rs index 530f6dd..e798d57 100644 --- a/src/mem.rs +++ b/src/mem.rs @@ -1,31 +1,35 @@ //! Reports the native footprint of SDK objects to V8's garbage collector. //! -//! Each binding class holds a small native allocation behind a small JS object: the JS -//! side is a few dozen bytes, while the native side ranges from ~200 KiB for a Processor -//! to the model weights themselves. Without a signal, V8's GC heuristics never see that -//! cost: `heapUsed` and `external` barely move, so the collector feels no pressure to -//! reclaim dropped instances, and a workload that creates processors per unit of work -//! ratchets RSS up until the process is OOM-killed. +//! Each binding class is a small JS object of a few dozen bytes in front of a much larger +//! native allocation, from ~200 KiB for a processor up to the size of the model weights. +//! V8's heuristics only see the JS side: `heapUsed` and `external` barely move, so the +//! collector feels no pressure to reclaim dropped instances, and a workload that creates +//! processors per unit of work ratchets RSS up until the process is OOM-killed. //! -//! `Env::adjust_external_memory` is the Node-API mechanism for this (`napi_adjust_external_memory`, -//! the same fix the Ruby binding ships as `rb_gc_adjust_memory_usage`). Each constructor -//! reports its object's footprint; `ObjectFinalize::finalize` reports the negation when the -//! instance is collected, so the ledger balances. +//! `Env::adjust_external_memory` (`napi_adjust_external_memory`) is the Node-API mechanism +//! for reporting that hidden cost. Each constructor reports its object's footprint, and the +//! negation is reported when the instance goes away: by `dispose()`, by +//! [`MemoryTracked::release`], or by the class finalizer. Either way the ledger balances. //! -//! Values are estimates keyed to measurement and deliberately err high: over-reporting -//! only makes V8 collect a little more eagerly, while under-reporting is what caused the -//! unbounded growth. Per-class constants are used because the SDK does not (yet) expose a -//! per-instance memory query. +//! Footprints are per-class constants because the SDK exposes no per-instance memory +//! query. They are estimates keyed to measurement and deliberately err high: over-reporting +//! only makes V8 collect a little more eagerly, while under-reporting leaves the growth +//! above unchecked. -use std::path::Path; +use std::{path::Path, sync::Mutex}; use napi::Env; +use crate::{ + error::{Result, disposed_error}, + processor_async::lock, +}; + const MIB: i64 = 1024 * 1024; const KIB: i64 = 1024; -/// A `Processor` or `Vad` instance. Measured by holding N instances and reading the RSS -/// delta, at initialize + one `process`: +/// A `Processor` or `Vad` instance. Measured as the RSS delta per instance, at initialize +/// plus one `process` call: /// /// - `quail-vf-2.2-s` (5 MiB model): ~190 KiB /// - `quail-vf-2.2-l` (20 MiB model): ~462 KiB @@ -68,3 +72,57 @@ pub(crate) fn model_bytes(path: &Path) -> i64 { .map(|meta| meta.len() as i64) .unwrap_or(MODEL_FALLBACK_BYTES) } + +/// An SDK object whose native lifetime several JS handles and in-flight tasks share, and +/// whose footprint is reported to V8 for exactly as long as the object lives. +/// +/// `ProcessorAsync` and `VadAsync` can hand out a second JS handle onto the same native +/// object (`withConfig`), and tasks on the libuv pool hold `Arc` clones. The footprint is +/// therefore reported once per object, at construction, and given back exactly once: by +/// `dispose()`, or by the finalizer of the last surviving handle, whichever comes first. +/// `Option::take` makes the give-back idempotent, so the two never double-count. +pub(crate) struct MemoryTracked { + inner: Mutex>, +} + +impl MemoryTracked { + pub(crate) fn new(inner: T) -> Self { + Self { + inner: Mutex::new(Some(inner)), + } + } + + /// Destroys the native object, if it is still live, and gives its footprint back to + /// V8. Idempotent. + /// + /// Called by `dispose()`, which destroys the object regardless of other handles, and + /// by the finalizer of the last surviving handle. + pub(crate) fn release(&self, env: Env, bytes: i64) { + if lock(&self.inner).take().is_some() { + adjust(env, -bytes); + } + } + + /// Runs `f` with the native object, or fails with the disposed error. + pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { + let mut guard = lock(&self.inner); + let inner = guard.as_mut().ok_or_else(|| disposed_error(class))?; + f(inner) + } +} + +impl Drop for MemoryTracked { + fn drop(&mut self) { + // Frees the native object when the last `Arc` goes away without `release` having + // run, e.g. the last JS handle was finalized while a task still held a clone (that + // finalizer saw the extra reference and left the object for the task). The footprint + // report cannot be returned here, since that takes the finalizer's `Env`, so those + // bytes stay reported, which only makes V8 a little more eager for the rest of the + // process. + match self.inner.get_mut() { + Ok(slot) => slot.take(), + // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. + Err(poisoned) => poisoned.into_inner().take(), + }; + } +} diff --git a/src/processor_async.rs b/src/processor_async.rs index 9b6c0eb..dea0a07 100644 --- a/src/processor_async.rs +++ b/src/processor_async.rs @@ -1,7 +1,7 @@ use crate::{ claim_sdk_id, - error::{Result, disposed_error, map_err}, - mem, + error::{Result, map_err}, + mem::{self, MemoryTracked}, model::Model, processor::{OtelConfig, ProcessorContext, audio_config}, }; @@ -29,59 +29,6 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { .unwrap_or_else(|poisoned| poisoned.into_inner()) } -/// An SDK object whose native lifetime several JS handles and in-flight tasks share. -/// -/// `ProcessorAsync` and `VadAsync` can hand out a second JS handle onto the same native -/// object (`withConfig`), and tasks on the libuv pool hold `Arc` clones. The footprint is -/// therefore reported to V8 once per object, at construction, and given back exactly -/// once: by `dispose()`, or by the finalizer of the last surviving handle, whichever -/// comes first. `Option::take` makes the give-back idempotent, so the two never -/// double-count. -pub(crate) struct Held { - inner: Mutex>, -} - -impl Held { - pub(crate) fn new(inner: T) -> Self { - Self { - inner: Mutex::new(Some(inner)), - } - } - - /// Destroys the native object, if it is still live, and gives its footprint back to - /// V8. Idempotent. - /// - /// Called by `dispose()`, which destroys the object regardless of other handles, and - /// by the finalizer of the last surviving handle. - pub(crate) fn release(&self, env: Env, bytes: i64) { - if lock(&self.inner).take().is_some() { - mem::adjust(env, -bytes); - } - } - - /// Runs `f` with the native object, or fails with the disposed error. - pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { - let mut guard = lock(&self.inner); - let inner = guard.as_mut().ok_or_else(|| disposed_error(class))?; - f(inner) - } -} - -impl Drop for Held { - fn drop(&mut self) { - // Frees the native object when the last `Arc` goes away without `release` having - // run, e.g. the last JS handle was finalized while a task still held a clone (that - // finalizer saw the extra reference and left the object for the task). The footprint - // report cannot be returned here — that takes the finalizer's `Env` — so those bytes - // stay reported, which only makes V8 a little more eager for the rest of the process. - match self.inner.get_mut() { - Ok(slot) => slot.take(), - // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. - Err(poisoned) => poisoned.into_inner().take(), - }; - } -} - /// Speech enhancement processor that keeps its work off the main thread. /// /// The same processing as {@link Processor}, but each call returns a promise and runs on @@ -103,7 +50,7 @@ impl Drop for Held { /// binding deliberately does not use. #[napi(custom_finalize)] pub struct ProcessorAsync { - inner: Arc>>, + inner: Arc>>, } impl ObjectFinalize for ProcessorAsync { @@ -142,7 +89,7 @@ impl ProcessorAsync { } None => aic_sdk::Processor::new(model_inner, &license_key), }; - let inner = Arc::new(Held::new(map_err(inner)?)); + let inner = Arc::new(MemoryTracked::new(map_err(inner)?)); mem::adjust(env, mem::PROCESSOR_BYTES); Ok(Self { inner }) @@ -251,7 +198,7 @@ impl ProcessorAsync { /// Backs {@link ProcessorAsync#withConfig}. pub struct ProcessorWithConfigTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -277,7 +224,7 @@ impl Task for ProcessorWithConfigTask { /// Backs {@link ProcessorAsync#initialize}. pub struct ProcessorInitializeTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -298,7 +245,7 @@ impl Task for ProcessorInitializeTask { /// Backs {@link ProcessorAsync#process}. pub struct ProcessorProcessTask { - inner: Arc>>, + inner: Arc>>, audio: Vec, } @@ -326,7 +273,7 @@ impl Task for ProcessorProcessTask { /// Backs {@link ProcessorAsync#getContext}. pub struct ProcessorContextTask { - inner: Arc>>, + inner: Arc>>, } impl Task for ProcessorContextTask { @@ -346,7 +293,7 @@ impl Task for ProcessorContextTask { /// Backs {@link ProcessorAsync#terminateSession}. pub struct ProcessorTerminateTask { - inner: Arc>>, + inner: Arc>>, } impl Task for ProcessorTerminateTask { diff --git a/src/vad_async.rs b/src/vad_async.rs index d28c5c0..4f5e4de 100644 --- a/src/vad_async.rs +++ b/src/vad_async.rs @@ -1,10 +1,9 @@ use crate::{ claim_sdk_id, error::{Result, map_err}, - mem, + mem::{self, MemoryTracked}, model::Model, processor::{OtelConfig, audio_config}, - processor_async::Held, vad::VadContext, }; @@ -39,7 +38,7 @@ use std::sync::Arc; /// binding deliberately does not use. #[napi(custom_finalize)] pub struct VadAsync { - inner: Arc>>, + inner: Arc>>, } impl ObjectFinalize for VadAsync { @@ -76,7 +75,7 @@ impl VadAsync { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), None => aic_sdk::Vad::new(model_inner, &license_key), }; - let inner = Arc::new(Held::new(map_err(inner)?)); + let inner = Arc::new(MemoryTracked::new(map_err(inner)?)); mem::adjust(env, mem::PROCESSOR_BYTES); Ok(Self { inner }) @@ -183,7 +182,7 @@ impl VadAsync { /// Backs {@link VadAsync#withConfig}. pub struct VadWithConfigTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -209,7 +208,7 @@ impl Task for VadWithConfigTask { /// Backs {@link VadAsync#initialize}. pub struct VadInitializeTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -230,7 +229,7 @@ impl Task for VadInitializeTask { /// Backs {@link VadAsync#process}. pub struct VadProcessTask { - inner: Arc>>, + inner: Arc>>, audio: Vec, } @@ -258,7 +257,7 @@ impl Task for VadProcessTask { /// Backs {@link VadAsync#getContext}. pub struct VadContextTask { - inner: Arc>>, + inner: Arc>>, } impl Task for VadContextTask { @@ -276,7 +275,7 @@ impl Task for VadContextTask { /// Backs {@link VadAsync#terminateSession}. pub struct VadTerminateTask { - inner: Arc>>, + inner: Arc>>, } impl Task for VadTerminateTask { From 4ec5e86f3fb80c12d1d50fb9e78936f760c2101d Mon Sep 17 00:00:00 2001 From: steckes Date: Fri, 4 Sep 2026 10:49:18 +0200 Subject: [PATCH 5/7] rename to DisposableSlot --- src/disposable_slot.rs | 69 ++++++++++++++++++++++++++++++++++++++++++ src/lib.rs | 1 + src/mem.rs | 64 ++------------------------------------- src/processor_async.rs | 17 ++++++----- src/vad_async.rs | 17 ++++++----- 5 files changed, 91 insertions(+), 77 deletions(-) create mode 100644 src/disposable_slot.rs diff --git a/src/disposable_slot.rs b/src/disposable_slot.rs new file mode 100644 index 0000000..e5fbef8 --- /dev/null +++ b/src/disposable_slot.rs @@ -0,0 +1,69 @@ +//! One native SDK object shared by several JS handles, destroyed exactly once. +//! +//! `ProcessorAsync` and `VadAsync` are the only classes that need this: `withConfig` hands +//! out a second JS handle onto the same native object, and each in-flight task on the +//! libuv pool holds an `Arc` clone. Whichever owner gets there first destroys the object +//! and gives its footprint back to V8; the rest find it already gone. The sync classes own +//! their object outright and use a plain `Option` field instead. + +use std::sync::Mutex; + +use napi::Env; + +use crate::{ + error::{Result, disposed_error}, + mem::adjust, + processor_async::lock, +}; + +/// A slot holding a native SDK object with several owners, emptied by whichever one gets +/// there first. Every later access finds it empty and fails with the disposed error. +/// +/// Its footprint is reported to V8 once per object, at construction, and given back +/// exactly once: by `dispose()`, or by the finalizer of the last surviving handle. +/// `Option::take` makes the give-back idempotent, so the two never double-count. +pub(crate) struct DisposableSlot { + inner: Mutex>, +} + +impl DisposableSlot { + pub(crate) fn new(inner: T) -> Self { + Self { + inner: Mutex::new(Some(inner)), + } + } + + /// Destroys the native object, if it is still live, and gives its footprint back to + /// V8. Idempotent. + /// + /// Called by `dispose()`, which destroys the object regardless of other handles, and + /// by the finalizer of the last surviving handle. + pub(crate) fn release(&self, env: Env, bytes: i64) { + if lock(&self.inner).take().is_some() { + adjust(env, -bytes); + } + } + + /// Runs `f` with the native object, or fails with the disposed error. + pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { + let mut guard = lock(&self.inner); + let inner = guard.as_mut().ok_or_else(|| disposed_error(class))?; + f(inner) + } +} + +impl Drop for DisposableSlot { + fn drop(&mut self) { + // Frees the native object when the last `Arc` goes away without `release` having + // run, e.g. the last JS handle was finalized while a task still held a clone (that + // finalizer saw the extra reference and left the object for the task). The footprint + // report cannot be returned here, since that takes the finalizer's `Env`, so those + // bytes stay reported, which only makes V8 a little more eager for the rest of the + // process. + match self.inner.get_mut() { + Ok(slot) => slot.take(), + // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. + Err(poisoned) => poisoned.into_inner().take(), + }; + } +} diff --git a/src/lib.rs b/src/lib.rs index 6736539..1e32358 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -8,6 +8,7 @@ use napi_derive::napi; mod analyzer; +mod disposable_slot; mod error; mod mem; mod model; diff --git a/src/mem.rs b/src/mem.rs index e798d57..5b82a2f 100644 --- a/src/mem.rs +++ b/src/mem.rs @@ -9,22 +9,18 @@ //! `Env::adjust_external_memory` (`napi_adjust_external_memory`) is the Node-API mechanism //! for reporting that hidden cost. Each constructor reports its object's footprint, and the //! negation is reported when the instance goes away: by `dispose()`, by -//! [`MemoryTracked::release`], or by the class finalizer. Either way the ledger balances. +//! [`DisposableSlot::release`](crate::disposable_slot::DisposableSlot::release), or by the +//! class finalizer. Either way the ledger balances. //! //! Footprints are per-class constants because the SDK exposes no per-instance memory //! query. They are estimates keyed to measurement and deliberately err high: over-reporting //! only makes V8 collect a little more eagerly, while under-reporting leaves the growth //! above unchecked. -use std::{path::Path, sync::Mutex}; +use std::path::Path; use napi::Env; -use crate::{ - error::{Result, disposed_error}, - processor_async::lock, -}; - const MIB: i64 = 1024 * 1024; const KIB: i64 = 1024; @@ -72,57 +68,3 @@ pub(crate) fn model_bytes(path: &Path) -> i64 { .map(|meta| meta.len() as i64) .unwrap_or(MODEL_FALLBACK_BYTES) } - -/// An SDK object whose native lifetime several JS handles and in-flight tasks share, and -/// whose footprint is reported to V8 for exactly as long as the object lives. -/// -/// `ProcessorAsync` and `VadAsync` can hand out a second JS handle onto the same native -/// object (`withConfig`), and tasks on the libuv pool hold `Arc` clones. The footprint is -/// therefore reported once per object, at construction, and given back exactly once: by -/// `dispose()`, or by the finalizer of the last surviving handle, whichever comes first. -/// `Option::take` makes the give-back idempotent, so the two never double-count. -pub(crate) struct MemoryTracked { - inner: Mutex>, -} - -impl MemoryTracked { - pub(crate) fn new(inner: T) -> Self { - Self { - inner: Mutex::new(Some(inner)), - } - } - - /// Destroys the native object, if it is still live, and gives its footprint back to - /// V8. Idempotent. - /// - /// Called by `dispose()`, which destroys the object regardless of other handles, and - /// by the finalizer of the last surviving handle. - pub(crate) fn release(&self, env: Env, bytes: i64) { - if lock(&self.inner).take().is_some() { - adjust(env, -bytes); - } - } - - /// Runs `f` with the native object, or fails with the disposed error. - pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { - let mut guard = lock(&self.inner); - let inner = guard.as_mut().ok_or_else(|| disposed_error(class))?; - f(inner) - } -} - -impl Drop for MemoryTracked { - fn drop(&mut self) { - // Frees the native object when the last `Arc` goes away without `release` having - // run, e.g. the last JS handle was finalized while a task still held a clone (that - // finalizer saw the extra reference and left the object for the task). The footprint - // report cannot be returned here, since that takes the finalizer's `Env`, so those - // bytes stay reported, which only makes V8 a little more eager for the rest of the - // process. - match self.inner.get_mut() { - Ok(slot) => slot.take(), - // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. - Err(poisoned) => poisoned.into_inner().take(), - }; - } -} diff --git a/src/processor_async.rs b/src/processor_async.rs index dea0a07..b14e900 100644 --- a/src/processor_async.rs +++ b/src/processor_async.rs @@ -1,7 +1,8 @@ use crate::{ claim_sdk_id, + disposable_slot::DisposableSlot, error::{Result, map_err}, - mem::{self, MemoryTracked}, + mem, model::Model, processor::{OtelConfig, ProcessorContext, audio_config}, }; @@ -50,7 +51,7 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { /// binding deliberately does not use. #[napi(custom_finalize)] pub struct ProcessorAsync { - inner: Arc>>, + inner: Arc>>, } impl ObjectFinalize for ProcessorAsync { @@ -89,7 +90,7 @@ impl ProcessorAsync { } None => aic_sdk::Processor::new(model_inner, &license_key), }; - let inner = Arc::new(MemoryTracked::new(map_err(inner)?)); + let inner = Arc::new(DisposableSlot::new(map_err(inner)?)); mem::adjust(env, mem::PROCESSOR_BYTES); Ok(Self { inner }) @@ -198,7 +199,7 @@ impl ProcessorAsync { /// Backs {@link ProcessorAsync#withConfig}. pub struct ProcessorWithConfigTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -224,7 +225,7 @@ impl Task for ProcessorWithConfigTask { /// Backs {@link ProcessorAsync#initialize}. pub struct ProcessorInitializeTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -245,7 +246,7 @@ impl Task for ProcessorInitializeTask { /// Backs {@link ProcessorAsync#process}. pub struct ProcessorProcessTask { - inner: Arc>>, + inner: Arc>>, audio: Vec, } @@ -273,7 +274,7 @@ impl Task for ProcessorProcessTask { /// Backs {@link ProcessorAsync#getContext}. pub struct ProcessorContextTask { - inner: Arc>>, + inner: Arc>>, } impl Task for ProcessorContextTask { @@ -293,7 +294,7 @@ impl Task for ProcessorContextTask { /// Backs {@link ProcessorAsync#terminateSession}. pub struct ProcessorTerminateTask { - inner: Arc>>, + inner: Arc>>, } impl Task for ProcessorTerminateTask { diff --git a/src/vad_async.rs b/src/vad_async.rs index 4f5e4de..d7f9b05 100644 --- a/src/vad_async.rs +++ b/src/vad_async.rs @@ -1,7 +1,8 @@ use crate::{ claim_sdk_id, + disposable_slot::DisposableSlot, error::{Result, map_err}, - mem::{self, MemoryTracked}, + mem, model::Model, processor::{OtelConfig, audio_config}, vad::VadContext, @@ -38,7 +39,7 @@ use std::sync::Arc; /// binding deliberately does not use. #[napi(custom_finalize)] pub struct VadAsync { - inner: Arc>>, + inner: Arc>>, } impl ObjectFinalize for VadAsync { @@ -75,7 +76,7 @@ impl VadAsync { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), None => aic_sdk::Vad::new(model_inner, &license_key), }; - let inner = Arc::new(MemoryTracked::new(map_err(inner)?)); + let inner = Arc::new(DisposableSlot::new(map_err(inner)?)); mem::adjust(env, mem::PROCESSOR_BYTES); Ok(Self { inner }) @@ -182,7 +183,7 @@ impl VadAsync { /// Backs {@link VadAsync#withConfig}. pub struct VadWithConfigTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -208,7 +209,7 @@ impl Task for VadWithConfigTask { /// Backs {@link VadAsync#initialize}. pub struct VadInitializeTask { - inner: Arc>>, + inner: Arc>>, config: aic_sdk::ProcessorConfig, } @@ -229,7 +230,7 @@ impl Task for VadInitializeTask { /// Backs {@link VadAsync#process}. pub struct VadProcessTask { - inner: Arc>>, + inner: Arc>>, audio: Vec, } @@ -257,7 +258,7 @@ impl Task for VadProcessTask { /// Backs {@link VadAsync#getContext}. pub struct VadContextTask { - inner: Arc>>, + inner: Arc>>, } impl Task for VadContextTask { @@ -275,7 +276,7 @@ impl Task for VadContextTask { /// Backs {@link VadAsync#terminateSession}. pub struct VadTerminateTask { - inner: Arc>>, + inner: Arc>>, } impl Task for VadTerminateTask { From 69745e03753c4e816fec6416b076aebef3b03bdd Mon Sep 17 00:00:00 2001 From: steckes Date: Mon, 7 Sep 2026 10:54:02 +0200 Subject: [PATCH 6/7] change `ids` to `IDs` --- __test__/index.spec.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/__test__/index.spec.ts b/__test__/index.spec.ts index 37a7eb5..b7a7e9d 100644 --- a/__test__/index.spec.ts +++ b/__test__/index.spec.ts @@ -36,7 +36,7 @@ test('sdk reports a version', (t) => { }) test('sdk expects the model version the fixtures are published under', (t) => { - // Guards against the fixture ids in models.ts drifting from the addon's expectation, + // Guards against the fixture IDs in models.ts drifting from the addon's expectation, // which would otherwise surface as a confusing "model version unsupported" much later. t.is(getCompatibleModelVersion(), TEST_MODEL_VERSION) }) From 91574e09a291bd52e37fcf733f1172eedfdbd28d Mon Sep 17 00:00:00 2001 From: steckes Date: Mon, 7 Sep 2026 12:07:42 +0200 Subject: [PATCH 7/7] improve DisposableSlot --- src/analyzer.rs | 97 +++++++++++++++++--------------------- src/disposable_slot.rs | 103 ++++++++++++++++++++++++----------------- src/mem.rs | 25 ++++++---- src/model.rs | 31 ++++++------- src/processor.rs | 50 +++++++------------- src/processor_async.rs | 84 +++++++++++++-------------------- src/vad.rs | 44 +++++++----------- src/vad_async.rs | 66 +++++++++++++------------- 8 files changed, 230 insertions(+), 270 deletions(-) diff --git a/src/analyzer.rs b/src/analyzer.rs index aaf6178..132ae38 100644 --- a/src/analyzer.rs +++ b/src/analyzer.rs @@ -1,10 +1,10 @@ use crate::{ claim_sdk_id, - error::{Result, disposed_error, map_err}, + disposable_slot::{DisposableSlot, lock}, + error::{Result, map_err}, mem, model::Model, processor::audio_config, - processor_async::{Shared, lock}, }; use napi::{ @@ -76,20 +76,23 @@ pub struct Analyzer { // Only the analyzer half is shared. The collector is owned outright, so `buffer`, the // one call on the audio path, takes no lock and cannot contend with an analysis running // on a worker thread. - collector: Option, - analyzer: Shared>>, + collector: DisposableSlot, + analyzer: Arc>>>, } impl ObjectFinalize for Analyzer { - fn finalize(self, env: Env) -> Result<()> { - // `dispose()` already gave the footprint back when the collector is gone. + fn finalize(mut self, env: Env) -> Result<()> { + // Each half is released on its own, and each release is a no-op when `dispose()` + // got there first. // // An in-flight `AnalyzeTask` holds the analyzer `Arc`, so that half is dropped - // only after the worker finishes. The collector drops here, possibly while the - // worker analyzes. That is safe per the C API, which destroys the paired halves - // independently, in any order (`aic_collector_destroy`). - if self.collector.is_some() { - mem::adjust(env, -mem::ANALYZER_BYTES); + // only after the worker finishes, and its footprint is left for whoever holds the + // last handle. The collector drops here, possibly while the worker analyzes. That + // is safe per the C API, which destroys the paired halves independently, in any + // order (`aic_collector_destroy`). + self.collector.release(env); + if Arc::strong_count(&self.analyzer) == 1 { + lock(&self.analyzer).release(env); } Ok(()) } @@ -103,11 +106,15 @@ impl Analyzer { let model_inner = model.inner()?; claim_sdk_id(); let (collector, analyzer) = map_err(aic_sdk::analyzer_pair(model_inner, &license_key))?; - mem::adjust(env, mem::ANALYZER_BYTES); Ok(Self { - collector: Some(collector), - analyzer: Arc::new(Mutex::new(Some(analyzer))), + collector: DisposableSlot::new(env, collector, "Analyzer", mem::COLLECTOR_BYTES), + analyzer: Arc::new(Mutex::new(DisposableSlot::new( + env, + analyzer, + "Analyzer", + mem::ANALYZER_BYTES, + ))), }) } @@ -118,16 +125,16 @@ impl Analyzer { /// in-flight `analyzeAsync` on a worker thread finishes. #[napi] pub fn dispose(&mut self, env: Env) { - if self.collector.take().is_some() { - // The collector drops before the analyzer lock is taken, so it can be destroyed - // while an `analyzeAsync` is in flight on a worker. That is safe per the C API - // (`aic_collector_destroy`): the paired halves are destroyed independently, in - // any order, and the collector handle itself is only ever used on this thread. - // The analyzer half is destroyed under the lock, which is what blocks until the - // in-flight analysis finishes. - lock(&self.analyzer).take(); // dropped on scope exit - mem::adjust(env, -mem::ANALYZER_BYTES); - } + // The collector is released before the analyzer lock is taken, so it can be + // destroyed while an `analyzeAsync` is in flight on a worker. That is safe per the + // C API (`aic_collector_destroy`): the paired halves are destroyed independently, in + // any order, and the collector handle itself is only ever used on this thread. The + // analyzer half is destroyed under the lock, which is what blocks until the + // in-flight analysis finishes. + // + // Both releases are idempotent, so a second `dispose()` does nothing. + self.collector.release(env); + lock(&self.analyzer).release(env); } /// Configures the analyzer for an audio format. Must be called before buffering. @@ -141,31 +148,17 @@ impl Analyzer { block_size: u32, variable_block_size: Option, ) -> Result<()> { - let collector = self - .collector - .as_mut() - .ok_or_else(|| disposed_error("Analyzer"))?; - map_err(collector.initialize(&audio_config(sample_rate, block_size, variable_block_size))) - } - - /// Runs `f` with the analyzer half, or fails with the disposed error. - fn with_analyzer( - &self, - f: impl FnOnce(&mut aic_sdk::Analyzer<'static>) -> Result, - ) -> Result { - let mut guard = lock(&self.analyzer); - let analyzer = guard.as_mut().ok_or_else(|| disposed_error("Analyzer"))?; - f(analyzer) + map_err(self.collector.get_mut()?.initialize(&audio_config( + sample_rate, + block_size, + variable_block_size, + ))) } /// Buffers a mono audio block for later analysis, leaving the audio unmodified. #[napi] pub fn buffer(&mut self, audio: Float32Array) -> Result<()> { - let collector = self - .collector - .as_mut() - .ok_or_else(|| disposed_error("Analyzer"))?; - map_err(collector.buffer(&audio)) + map_err(self.collector.get_mut()?.buffer(&audio)) } /// Runs the analysis model over the buffered audio, on the calling thread. @@ -179,9 +172,7 @@ impl Analyzer { /// nothing else is waiting on the event loop. #[napi] pub fn analyze(&self) -> Result { - self - .with_analyzer(|analyzer| map_err(analyzer.analyze_buffered())) - .map(AnalysisResult::from) + map_err(lock(&self.analyzer).get_mut()?.analyze_buffered()).map(AnalysisResult::from) } /// Runs the analysis model over the buffered audio on a worker thread. @@ -207,7 +198,7 @@ impl Analyzer { /// Clears buffered audio and internal state, keeping the configured audio settings. #[napi] pub fn reset(&self) -> Result<()> { - self.with_analyzer(|analyzer| map_err(analyzer.reset())) + map_err(lock(&self.analyzer).get_mut()?.reset()) } /// Swaps in a renewed JWT without tearing down the analyzer. @@ -216,13 +207,13 @@ impl Analyzer { /// call is a no-op and the previous token stays active. #[napi] pub fn update_bearer_token(&self, token: String) -> Result<()> { - self.with_analyzer(|analyzer| map_err(analyzer.update_bearer_token(&token))) + map_err(lock(&self.analyzer).get_mut()?.update_bearer_token(&token)) } /// Ends this analyzer's telemetry session, after which it can no longer analyze audio. #[napi] pub fn terminate_session(&self) -> Result<()> { - self.with_analyzer(|analyzer| map_err(analyzer.terminate_session())) + map_err(lock(&self.analyzer).get_mut()?.terminate_session()) } } @@ -231,7 +222,7 @@ impl Analyzer { /// Holds only the analyzer half, so the collector stays on the JS thread where `buffer` can /// keep reaching it while this runs. pub struct AnalyzeTask { - analyzer: Shared>>, + analyzer: Arc>>>, } impl Task for AnalyzeTask { @@ -239,9 +230,7 @@ impl Task for AnalyzeTask { type JsValue = AnalysisResult; fn compute(&mut self) -> Result { - let mut guard = lock(&self.analyzer); - let analyzer = guard.as_mut().ok_or_else(|| disposed_error("Analyzer"))?; - map_err(analyzer.analyze_buffered()) + map_err(lock(&self.analyzer).get_mut()?.analyze_buffered()) } fn resolve(&mut self, _env: Env, result: aic_sdk::AnalysisResult) -> Result { diff --git a/src/disposable_slot.rs b/src/disposable_slot.rs index e5fbef8..b98853c 100644 --- a/src/disposable_slot.rs +++ b/src/disposable_slot.rs @@ -1,69 +1,88 @@ -//! One native SDK object shared by several JS handles, destroyed exactly once. +//! One native SDK object, destroyed exactly once, with its footprint reported to V8 for +//! as long as it lives. //! -//! `ProcessorAsync` and `VadAsync` are the only classes that need this: `withConfig` hands -//! out a second JS handle onto the same native object, and each in-flight task on the -//! libuv pool holds an `Arc` clone. Whichever owner gets there first destroys the object -//! and gives its footprint back to V8; the rest find it already gone. The sync classes own -//! their object outright and use a plain `Option` field instead. +//! Every binding class holds its SDK object in a slot, so the disposed error and the +//! footprint accounting are written once here rather than per class. The slot itself is +//! a plain owner with no interior mutability; the classes that share their object with +//! tasks on the libuv pool wrap it in an `Arc>` and reach it through [`lock`]. -use std::sync::Mutex; +use std::sync::{Mutex, MutexGuard}; use napi::Env; use crate::{ error::{Result, disposed_error}, mem::adjust, - processor_async::lock, }; -/// A slot holding a native SDK object with several owners, emptied by whichever one gets -/// there first. Every later access finds it empty and fails with the disposed error. +/// A slot holding a native SDK object until it is disposed, after which every access +/// fails with the disposed error. /// -/// Its footprint is reported to V8 once per object, at construction, and given back -/// exactly once: by `dispose()`, or by the finalizer of the last surviving handle. -/// `Option::take` makes the give-back idempotent, so the two never double-count. +/// The slot owns both halves of the object's footprint report: [`new`](Self::new) reports +/// it to V8 and [`release`](Self::release) gives it back. `Option::take` makes the +/// give-back idempotent, so `dispose()` and a class finalizer cannot double-count it. +/// +/// A slot dropped without `release` having run still destroys the object, but those +/// bytes stay reported: giving them back takes an `Env`, which a `Drop` impl does not +/// have. That happens when the last JS handle onto a shared object is finalized while a +/// task still holds a clone; over-reporting only makes V8 a little more eager for the +/// rest of the process. pub(crate) struct DisposableSlot { - inner: Mutex>, + inner: Option, + /// The JS class name, for the disposed error message. + class: &'static str, + /// The footprint reported to V8 while `inner` is live. + bytes: i64, } impl DisposableSlot { - pub(crate) fn new(inner: T) -> Self { + /// Takes ownership of `inner` and reports its `bytes` of native footprint to V8. + pub(crate) fn new(env: Env, inner: T, class: &'static str, bytes: i64) -> Self { + adjust(env, bytes); + Self { - inner: Mutex::new(Some(inner)), + inner: Some(inner), + class, + bytes, } } + /// The native object, or the disposed error once it is gone. + pub(crate) fn get(&self) -> Result<&T> { + self + .inner + .as_ref() + .ok_or_else(|| disposed_error(self.class)) + } + + /// The same, for a `&mut` call. + pub(crate) fn get_mut(&mut self) -> Result<&mut T> { + self + .inner + .as_mut() + .ok_or_else(|| disposed_error(self.class)) + } + /// Destroys the native object, if it is still live, and gives its footprint back to /// V8. Idempotent. /// - /// Called by `dispose()`, which destroys the object regardless of other handles, and - /// by the finalizer of the last surviving handle. - pub(crate) fn release(&self, env: Env, bytes: i64) { - if lock(&self.inner).take().is_some() { - adjust(env, -bytes); + /// Called by `dispose()` and by the class finalizer, in whichever order they happen. + pub(crate) fn release(&mut self, env: Env) { + if self.inner.take().is_some() { + adjust(env, -self.bytes); } } - - /// Runs `f` with the native object, or fails with the disposed error. - pub(crate) fn with(&self, class: &str, f: impl FnOnce(&mut T) -> Result) -> Result { - let mut guard = lock(&self.inner); - let inner = guard.as_mut().ok_or_else(|| disposed_error(class))?; - f(inner) - } } -impl Drop for DisposableSlot { - fn drop(&mut self) { - // Frees the native object when the last `Arc` goes away without `release` having - // run, e.g. the last JS handle was finalized while a task still held a clone (that - // finalizer saw the extra reference and left the object for the task). The footprint - // report cannot be returned here, since that takes the finalizer's `Env`, so those - // bytes stay reported, which only makes V8 a little more eager for the rest of the - // process. - match self.inner.get_mut() { - Ok(slot) => slot.take(), - // Dropping cannot fail on a poisoned lock: the guard's contents are still ours. - Err(poisoned) => poisoned.into_inner().take(), - }; - } +/// Locks a shared slot, recovering the guard if the lock is poisoned. +/// +/// A `Mutex` rather than an async lock: `compute` runs on a libuv worker, where blocking +/// is exactly what that thread is for. +/// +/// Poisoning would mean an earlier call panicked while holding the guard, which the SDK +/// does not do. Recovering keeps one hypothetical failure from turning every later call +/// into a panic, disposal included: a panic would happen inside an SDK call, leaving the +/// slot's own `Option` intact. +pub(crate) fn lock(slot: &Mutex) -> MutexGuard<'_, T> { + slot.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) } diff --git a/src/mem.rs b/src/mem.rs index 5b82a2f..dc13c27 100644 --- a/src/mem.rs +++ b/src/mem.rs @@ -7,10 +7,10 @@ //! processors per unit of work ratchets RSS up until the process is OOM-killed. //! //! `Env::adjust_external_memory` (`napi_adjust_external_memory`) is the Node-API mechanism -//! for reporting that hidden cost. Each constructor reports its object's footprint, and the -//! negation is reported when the instance goes away: by `dispose()`, by -//! [`DisposableSlot::release`](crate::disposable_slot::DisposableSlot::release), or by the -//! class finalizer. Either way the ledger balances. +//! for reporting that hidden cost. Both halves of the ledger live in +//! [`DisposableSlot`](crate::disposable_slot::DisposableSlot): it reports its object's +//! footprint when constructed and reports the negation when released, by `dispose()` or by +//! the class finalizer, whichever gets there first. Either way the ledger balances. //! //! Footprints are per-class constants because the SDK exposes no per-instance memory //! query. They are estimates keyed to measurement and deliberately err high: over-reporting @@ -36,10 +36,19 @@ const KIB: i64 = 1024; /// within ~3x of the smallest. pub(crate) const PROCESSOR_BYTES: i64 = 512 * KIB; -/// An `Analyzer` (collector + analyzer pair). Measured ~8.2 MiB for `tyto-1.1-l` at -/// construction and ~8.9 MiB with the collector initialized and holding 5 s of audio. -/// 16 MiB gives ~2x headroom for larger analysis models. -pub(crate) const ANALYZER_BYTES: i64 = 16 * MIB; +/// The analyzer half of an `Analyzer`, which holds the model workspace. Measured ~8.2 MiB +/// for `tyto-1.1-l` at construction, so 14 MiB leaves ~1.7x headroom for larger analysis +/// models. +/// +/// Reported separately from [`COLLECTOR_BYTES`] because the two halves are destroyed +/// independently: the collector can go while a worker thread still analyzes, so a single +/// report for the pair would be given back too early. +pub(crate) const ANALYZER_BYTES: i64 = 14 * MIB; + +/// The collector half of an `Analyzer`, which holds the buffered audio. Measured as the +/// ~0.7 MiB that `tyto-1.1-l` grows by once the collector is initialized and holding its +/// 5 s span, so 2 MiB leaves ~3x headroom. +pub(crate) const COLLECTOR_BYTES: i64 = 2 * MIB; /// Fallback footprint for a `Model` when its file cannot be stat'd. Deliberately /// conservative: the loaded model is memory-mapped, so its resident share approaches the diff --git a/src/model.rs b/src/model.rs index 53ed91b..10a262b 100644 --- a/src/model.rs +++ b/src/model.rs @@ -1,5 +1,6 @@ use crate::{ - error::{JsAicError, Result, disposed_error, map_err}, + disposable_slot::DisposableSlot, + error::{JsAicError, Result, map_err}, mem, }; @@ -15,18 +16,16 @@ use napi_derive::napi; pub struct Model { // `from_file` memory-maps the file rather than borrowing a caller-owned buffer, so // the SDK model is `'static` and needs no lifetime plumbing here. - inner: Option>, - /// Native footprint reported to V8's GC while this instance is alive (the mmap'd - /// weights). Reported back on dispose or finalize so the accounting balances. - reported_bytes: i64, + // + // The slot's footprint is this instance's alone: the mmap'd weights. It is per-instance + // rather than a per-class constant, since it is the model file's size. + slot: DisposableSlot>, } impl ObjectFinalize for Model { - fn finalize(self, env: Env) -> Result<()> { - // `dispose()` already gave the footprint back when the inner is gone. - if self.inner.is_some() { - mem::adjust(env, -self.reported_bytes); - } + fn finalize(mut self, env: Env) -> Result<()> { + // A no-op when `dispose()` already gave the footprint back. + self.slot.release(env); Ok(()) } } @@ -34,7 +33,7 @@ impl ObjectFinalize for Model { impl Model { /// The inner SDK model, or the disposed error once `dispose()` ran. pub(crate) fn inner(&self) -> Result<&aic_sdk::Model<'static>> { - self.inner.as_ref().ok_or_else(|| disposed_error("Model")) + self.slot.get() } } @@ -51,12 +50,10 @@ impl Model { #[napi(factory)] pub fn from_file(env: Env, path: String) -> Result { let inner = map_err(aic_sdk::Model::from_file(&path))?; - let reported_bytes = mem::model_bytes(std::path::Path::new(&path)); - mem::adjust(env, reported_bytes); + let bytes = mem::model_bytes(std::path::Path::new(&path)); Ok(Self { - inner: Some(inner), - reported_bytes, + slot: DisposableSlot::new(env, inner, "Model", bytes), }) } @@ -68,9 +65,7 @@ impl Model { /// throws; calling `dispose()` again does nothing. #[napi] pub fn dispose(&mut self, env: Env) { - if self.inner.take().is_some() { - mem::adjust(env, -self.reported_bytes); - } + self.slot.release(env); } /// Downloads a model from the ai-coustics artifact CDN and resolves to its path. diff --git a/src/processor.rs b/src/processor.rs index ed74f24..cbf92b2 100644 --- a/src/processor.rs +++ b/src/processor.rs @@ -1,6 +1,7 @@ use crate::{ claim_sdk_id, - error::{Result, disposed_error, map_err}, + disposable_slot::DisposableSlot, + error::{Result, map_err}, mem, model::Model, }; @@ -81,37 +82,19 @@ pub(crate) fn audio_config( /// Create several processors to handle multiple streams or to switch models at runtime. #[napi(custom_finalize)] pub struct Processor { - inner: Option>, + // Owned outright, with no lock: every method here runs on the JS thread. The async + // class shares the same slot with its tasks instead. + slot: DisposableSlot>, } impl ObjectFinalize for Processor { - fn finalize(self, env: Env) -> Result<()> { - // `dispose()` already gave the footprint back when the inner is gone. - if self.inner.is_some() { - mem::adjust(env, -mem::PROCESSOR_BYTES); - } + fn finalize(mut self, env: Env) -> Result<()> { + // A no-op when `dispose()` already gave the footprint back. + self.slot.release(env); Ok(()) } } -impl Processor { - /// The inner SDK processor, or the disposed error once `dispose()` ran. - fn inner(&self) -> Result<&aic_sdk::Processor<'static>> { - self - .inner - .as_ref() - .ok_or_else(|| disposed_error("Processor")) - } - - /// The same, for a `&mut` call. - fn inner_mut(&mut self) -> Result<&mut aic_sdk::Processor<'static>> { - self - .inner - .as_mut() - .ok_or_else(|| disposed_error("Processor")) - } -} - #[napi] impl Processor { /// Creates a processor from an enhancement or bypass model. @@ -134,9 +117,10 @@ impl Processor { None => aic_sdk::Processor::new(model_inner, &license_key), }; let inner = map_err(inner)?; - mem::adjust(env, mem::PROCESSOR_BYTES); - Ok(Self { inner: Some(inner) }) + Ok(Self { + slot: DisposableSlot::new(env, inner, "Processor", mem::PROCESSOR_BYTES), + }) } /// Destroys the native processor immediately, releasing its memory and telemetry @@ -145,9 +129,7 @@ impl Processor { /// Every later method throws; calling `dispose()` again does nothing. #[napi] pub fn dispose(&mut self, env: Env) { - if self.inner.take().is_some() { - mem::adjust(env, -mem::PROCESSOR_BYTES); - } + self.slot.release(env); } /// Configures the processor for an audio format. Must be called before processing. @@ -164,7 +146,7 @@ impl Processor { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err(self.inner_mut()?.initialize(&audio_config( + map_err(self.slot.get_mut()?.initialize(&audio_config( sample_rate, block_size, variable_block_size, @@ -187,7 +169,7 @@ impl Processor { // break that assumption, which is inherent to processing JS-owned buffers in place. let samples = unsafe { audio.as_mut() }; - map_err(self.inner_mut()?.process(samples)) + map_err(self.slot.get_mut()?.process(samples)) } /// Creates a handle for reading and writing this processor's parameters and state. @@ -196,7 +178,7 @@ impl Processor { #[napi] pub fn get_context(&self) -> Result { Ok(ProcessorContext { - inner: self.inner()?.context(), + inner: self.slot.get()?.context(), }) } @@ -206,7 +188,7 @@ impl Processor { /// is collected, but GC timing is not guaranteed. May block, so keep it off the audio path. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - map_err(self.inner_mut()?.terminate_session()) + map_err(self.slot.get_mut()?.terminate_session()) } } diff --git a/src/processor_async.rs b/src/processor_async.rs index b14e900..064e2fa 100644 --- a/src/processor_async.rs +++ b/src/processor_async.rs @@ -1,6 +1,6 @@ use crate::{ claim_sdk_id, - disposable_slot::DisposableSlot, + disposable_slot::{DisposableSlot, lock}, error::{Result, map_err}, mem, model::Model, @@ -12,23 +12,7 @@ use napi::{ bindgen_prelude::{AsyncTask, Float32Array, ObjectFinalize}, }; use napi_derive::napi; -use std::sync::{Arc, Mutex, MutexGuard}; - -/// The SDK object shared between a binding class and the tasks it spawns. -/// -/// A `Mutex` rather than an async lock: `compute` runs on a libuv worker, where blocking -/// is exactly what that thread is for. -pub(crate) type Shared = Arc>; - -/// Locks a shared SDK object, recovering the guard if the lock is poisoned. -/// -/// Poisoning would mean an earlier call panicked mid-process, which the SDK does not do. -/// Recovering keeps one hypothetical failure from turning every later call into a panic. -pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { - shared - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()) -} +use std::sync::{Arc, Mutex}; /// Speech enhancement processor that keeps its work off the main thread. /// @@ -51,7 +35,7 @@ pub(crate) fn lock(shared: &Mutex) -> MutexGuard<'_, T> { /// binding deliberately does not use. #[napi(custom_finalize)] pub struct ProcessorAsync { - inner: Arc>>, + slot: Arc>>>, } impl ObjectFinalize for ProcessorAsync { @@ -59,8 +43,8 @@ impl ObjectFinalize for ProcessorAsync { // Only the last handle onto the native object destroys it; with other handles or // in-flight tasks holding an `Arc`, this leaves the object (and its footprint // report) for them. Idempotent against `dispose()`. - if Arc::strong_count(&self.inner) == 1 { - self.inner.release(env, mem::PROCESSOR_BYTES); + if Arc::strong_count(&self.slot) == 1 { + lock(&self.slot).release(env); } Ok(()) } @@ -90,10 +74,16 @@ impl ProcessorAsync { } None => aic_sdk::Processor::new(model_inner, &license_key), }; - let inner = Arc::new(DisposableSlot::new(map_err(inner)?)); - mem::adjust(env, mem::PROCESSOR_BYTES); - - Ok(Self { inner }) + let inner = map_err(inner)?; + + Ok(Self { + slot: Arc::new(Mutex::new(DisposableSlot::new( + env, + inner, + "ProcessorAsync", + mem::PROCESSOR_BYTES, + ))), + }) } /// Destroys the native processor immediately, releasing its memory and telemetry @@ -103,7 +93,7 @@ impl ProcessorAsync { /// in-flight work on the libuv pool finishes. #[napi] pub fn dispose(&self, env: Env) { - self.inner.release(env, mem::PROCESSOR_BYTES); + lock(&self.slot).release(env); } /// Initializes the processor and resolves to a handle onto it, for chaining off the @@ -123,7 +113,7 @@ impl ProcessorAsync { variable_block_size: Option, ) -> AsyncTask { AsyncTask::new(ProcessorWithConfigTask { - inner: self.inner.clone(), + slot: self.slot.clone(), config: audio_config(sample_rate, block_size, variable_block_size), }) } @@ -139,7 +129,7 @@ impl ProcessorAsync { variable_block_size: Option, ) -> AsyncTask { AsyncTask::new(ProcessorInitializeTask { - inner: self.inner.clone(), + slot: self.slot.clone(), config: audio_config(sample_rate, block_size, variable_block_size), }) } @@ -165,7 +155,7 @@ impl ProcessorAsync { #[napi(ts_return_type = "Promise>")] pub fn process(&self, audio: Float32Array) -> AsyncTask { AsyncTask::new(ProcessorProcessTask { - inner: self.inner.clone(), + slot: self.slot.clone(), // Copied on the JS thread so the worker owns its samples outright. A block is a // couple of kilobytes, far below the cost of running the model over it, and it // removes any chance of JS mutating the buffer mid-process. @@ -182,7 +172,7 @@ impl ProcessorAsync { #[napi(ts_return_type = "Promise")] pub fn get_context(&self) -> AsyncTask { AsyncTask::new(ProcessorContextTask { - inner: self.inner.clone(), + slot: self.slot.clone(), }) } @@ -192,14 +182,14 @@ impl ProcessorAsync { #[napi(ts_return_type = "Promise")] pub fn terminate_session(&self) -> AsyncTask { AsyncTask::new(ProcessorTerminateTask { - inner: self.inner.clone(), + slot: self.slot.clone(), }) } } /// Backs {@link ProcessorAsync#withConfig}. pub struct ProcessorWithConfigTask { - inner: Arc>>, + slot: Arc>>>, config: aic_sdk::ProcessorConfig, } @@ -208,9 +198,7 @@ impl Task for ProcessorWithConfigTask { type JsValue = ProcessorAsync; fn compute(&mut self) -> Result<()> { - self.inner.with("ProcessorAsync", |inner| { - map_err(inner.initialize(&self.config)) - }) + map_err(lock(&self.slot).get_mut()?.initialize(&self.config)) } fn resolve(&mut self, _env: Env, _: ()) -> Result { @@ -218,14 +206,14 @@ impl Task for ProcessorWithConfigTask { // per object at construction, so there is nothing to report here; the last handle's // finalizer gives it back. Born disposed when the processor was disposed mid-flight. Ok(ProcessorAsync { - inner: self.inner.clone(), + slot: self.slot.clone(), }) } } /// Backs {@link ProcessorAsync#initialize}. pub struct ProcessorInitializeTask { - inner: Arc>>, + slot: Arc>>>, config: aic_sdk::ProcessorConfig, } @@ -234,9 +222,7 @@ impl Task for ProcessorInitializeTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - self.inner.with("ProcessorAsync", |inner| { - map_err(inner.initialize(&self.config)) - }) + map_err(lock(&self.slot).get_mut()?.initialize(&self.config)) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { @@ -246,7 +232,7 @@ impl Task for ProcessorInitializeTask { /// Backs {@link ProcessorAsync#process}. pub struct ProcessorProcessTask { - inner: Arc>>, + slot: Arc>>>, audio: Vec, } @@ -258,9 +244,7 @@ impl Task for ProcessorProcessTask { // Moved out rather than borrowed so the buffer can be handed to V8 in `resolve` // without another copy. The task is used once, so leaving an empty Vec behind is fine. let mut audio = std::mem::take(&mut self.audio); - self - .inner - .with("ProcessorAsync", |inner| map_err(inner.process(&mut audio)))?; + map_err(lock(&self.slot).get_mut()?.process(&mut audio))?; Ok(audio) } @@ -274,7 +258,7 @@ impl Task for ProcessorProcessTask { /// Backs {@link ProcessorAsync#getContext}. pub struct ProcessorContextTask { - inner: Arc>>, + slot: Arc>>>, } impl Task for ProcessorContextTask { @@ -282,9 +266,7 @@ impl Task for ProcessorContextTask { type JsValue = ProcessorContext; fn compute(&mut self) -> Result { - self - .inner - .with("ProcessorAsync", |inner| Ok(inner.context())) + Ok(lock(&self.slot).get()?.context()) } fn resolve(&mut self, _env: Env, context: aic_sdk::ProcessorContext) -> Result { @@ -294,7 +276,7 @@ impl Task for ProcessorContextTask { /// Backs {@link ProcessorAsync#terminateSession}. pub struct ProcessorTerminateTask { - inner: Arc>>, + slot: Arc>>>, } impl Task for ProcessorTerminateTask { @@ -302,9 +284,7 @@ impl Task for ProcessorTerminateTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - self - .inner - .with("ProcessorAsync", |inner| map_err(inner.terminate_session())) + map_err(lock(&self.slot).get_mut()?.terminate_session()) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { diff --git a/src/vad.rs b/src/vad.rs index 39c68eb..cb3e01f 100644 --- a/src/vad.rs +++ b/src/vad.rs @@ -1,6 +1,7 @@ use crate::{ claim_sdk_id, - error::{Result, disposed_error, map_err}, + disposable_slot::DisposableSlot, + error::{Result, map_err}, mem, model::Model, processor::{OtelConfig, audio_config}, @@ -56,31 +57,19 @@ impl From for aic_sdk::VadParameter { /// calling it on the same block before `Processor#process` is enough. #[napi(custom_finalize)] pub struct Vad { - inner: Option>, + // Owned outright, with no lock: every method here runs on the JS thread. The async + // class shares the same slot with its tasks instead. + slot: DisposableSlot>, } impl ObjectFinalize for Vad { - fn finalize(self, env: Env) -> Result<()> { - // `dispose()` already gave the footprint back when the inner is gone. - if self.inner.is_some() { - mem::adjust(env, -mem::PROCESSOR_BYTES); - } + fn finalize(mut self, env: Env) -> Result<()> { + // A no-op when `dispose()` already gave the footprint back. + self.slot.release(env); Ok(()) } } -impl Vad { - /// The inner SDK VAD, or the disposed error once `dispose()` ran. - fn inner(&self) -> Result<&aic_sdk::Vad<'static>> { - self.inner.as_ref().ok_or_else(|| disposed_error("Vad")) - } - - /// The same, for a `&mut` call. - fn inner_mut(&mut self) -> Result<&mut aic_sdk::Vad<'static>> { - self.inner.as_mut().ok_or_else(|| disposed_error("Vad")) - } -} - #[napi] impl Vad { /// Creates a voice activity detector from a dedicated VAD model. @@ -101,9 +90,10 @@ impl Vad { None => aic_sdk::Vad::new(model_inner, &license_key), }; let inner = map_err(inner)?; - mem::adjust(env, mem::PROCESSOR_BYTES); - Ok(Self { inner: Some(inner) }) + Ok(Self { + slot: DisposableSlot::new(env, inner, "Vad", mem::PROCESSOR_BYTES), + }) } /// Destroys the native VAD immediately, releasing its memory and telemetry session @@ -112,9 +102,7 @@ impl Vad { /// Every later method throws; calling `dispose()` again does nothing. #[napi] pub fn dispose(&mut self, env: Env) { - if self.inner.take().is_some() { - mem::adjust(env, -mem::PROCESSOR_BYTES); - } + self.slot.release(env); } /// Configures the VAD for an audio format. Must be called before processing. @@ -128,7 +116,7 @@ impl Vad { block_size: u32, variable_block_size: Option, ) -> Result<()> { - map_err(self.inner_mut()?.initialize(&audio_config( + map_err(self.slot.get_mut()?.initialize(&audio_config( sample_rate, block_size, variable_block_size, @@ -140,7 +128,7 @@ impl Vad { pub fn process(&mut self, audio: Float32Array) -> Result<()> { // Read-only, so the safe `Deref` to `&[f32]` is enough here. Taking the view by // value does not copy the caller's samples. - map_err(self.inner_mut()?.process(&audio)) + map_err(self.slot.get_mut()?.process(&audio)) } /// Creates a handle for reading predictions and controlling this VAD. @@ -149,14 +137,14 @@ impl Vad { #[napi] pub fn get_context(&self) -> Result { Ok(VadContext { - inner: self.inner()?.context(), + inner: self.slot.get()?.context(), }) } /// Ends this VAD's telemetry session, after which it can no longer process audio. #[napi] pub fn terminate_session(&mut self) -> Result<()> { - map_err(self.inner_mut()?.terminate_session()) + map_err(self.slot.get_mut()?.terminate_session()) } } diff --git a/src/vad_async.rs b/src/vad_async.rs index d7f9b05..36b347f 100644 --- a/src/vad_async.rs +++ b/src/vad_async.rs @@ -1,6 +1,6 @@ use crate::{ claim_sdk_id, - disposable_slot::DisposableSlot, + disposable_slot::{DisposableSlot, lock}, error::{Result, map_err}, mem, model::Model, @@ -13,7 +13,7 @@ use napi::{ bindgen_prelude::{AsyncTask, Float32Array, ObjectFinalize}, }; use napi_derive::napi; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; /// Voice activity detector that keeps its work off the main thread. /// @@ -39,7 +39,7 @@ use std::sync::Arc; /// binding deliberately does not use. #[napi(custom_finalize)] pub struct VadAsync { - inner: Arc>>, + slot: Arc>>>, } impl ObjectFinalize for VadAsync { @@ -47,8 +47,8 @@ impl ObjectFinalize for VadAsync { // Only the last handle onto the native object destroys it; with other handles or // in-flight tasks holding an `Arc`, this leaves the object (and its footprint // report) for them. Idempotent against `dispose()`. - if Arc::strong_count(&self.inner) == 1 { - self.inner.release(env, mem::PROCESSOR_BYTES); + if Arc::strong_count(&self.slot) == 1 { + lock(&self.slot).release(env); } Ok(()) } @@ -76,10 +76,16 @@ impl VadAsync { Some(config) => aic_sdk::Vad::with_otel_config(model_inner, &license_key, &config.into()), None => aic_sdk::Vad::new(model_inner, &license_key), }; - let inner = Arc::new(DisposableSlot::new(map_err(inner)?)); - mem::adjust(env, mem::PROCESSOR_BYTES); - - Ok(Self { inner }) + let inner = map_err(inner)?; + + Ok(Self { + slot: Arc::new(Mutex::new(DisposableSlot::new( + env, + inner, + "VadAsync", + mem::PROCESSOR_BYTES, + ))), + }) } /// Destroys the native VAD immediately, releasing its memory and telemetry session @@ -89,7 +95,7 @@ impl VadAsync { /// in-flight work on the libuv pool finishes. #[napi] pub fn dispose(&self, env: Env) { - self.inner.release(env, mem::PROCESSOR_BYTES); + lock(&self.slot).release(env); } /// Initializes the VAD and resolves to a handle onto it, for chaining off the @@ -109,7 +115,7 @@ impl VadAsync { variable_block_size: Option, ) -> AsyncTask { AsyncTask::new(VadWithConfigTask { - inner: self.inner.clone(), + slot: self.slot.clone(), config: audio_config(sample_rate, block_size, variable_block_size), }) } @@ -125,7 +131,7 @@ impl VadAsync { variable_block_size: Option, ) -> AsyncTask { AsyncTask::new(VadInitializeTask { - inner: self.inner.clone(), + slot: self.slot.clone(), config: audio_config(sample_rate, block_size, variable_block_size), }) } @@ -149,7 +155,7 @@ impl VadAsync { #[napi(ts_return_type = "Promise>")] pub fn process(&self, audio: Float32Array) -> AsyncTask { AsyncTask::new(VadProcessTask { - inner: self.inner.clone(), + slot: self.slot.clone(), // Copied on the JS thread so the worker owns its samples outright. A block is a // couple of kilobytes, far below the cost of running the model over it, and it // removes any chance of JS mutating the buffer mid-process. @@ -166,7 +172,7 @@ impl VadAsync { #[napi(ts_return_type = "Promise")] pub fn get_context(&self) -> AsyncTask { AsyncTask::new(VadContextTask { - inner: self.inner.clone(), + slot: self.slot.clone(), }) } @@ -176,14 +182,14 @@ impl VadAsync { #[napi(ts_return_type = "Promise")] pub fn terminate_session(&self) -> AsyncTask { AsyncTask::new(VadTerminateTask { - inner: self.inner.clone(), + slot: self.slot.clone(), }) } } /// Backs {@link VadAsync#withConfig}. pub struct VadWithConfigTask { - inner: Arc>>, + slot: Arc>>>, config: aic_sdk::ProcessorConfig, } @@ -192,9 +198,7 @@ impl Task for VadWithConfigTask { type JsValue = VadAsync; fn compute(&mut self) -> Result<()> { - self - .inner - .with("VadAsync", |inner| map_err(inner.initialize(&self.config))) + map_err(lock(&self.slot).get_mut()?.initialize(&self.config)) } fn resolve(&mut self, _env: Env, _: ()) -> Result { @@ -202,14 +206,14 @@ impl Task for VadWithConfigTask { // object at construction, so there is nothing to report here; the last handle's // finalizer gives it back. Born disposed when the VAD was disposed mid-flight. Ok(VadAsync { - inner: self.inner.clone(), + slot: self.slot.clone(), }) } } /// Backs {@link VadAsync#initialize}. pub struct VadInitializeTask { - inner: Arc>>, + slot: Arc>>>, config: aic_sdk::ProcessorConfig, } @@ -218,9 +222,7 @@ impl Task for VadInitializeTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - self - .inner - .with("VadAsync", |inner| map_err(inner.initialize(&self.config))) + map_err(lock(&self.slot).get_mut()?.initialize(&self.config)) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> { @@ -230,7 +232,7 @@ impl Task for VadInitializeTask { /// Backs {@link VadAsync#process}. pub struct VadProcessTask { - inner: Arc>>, + slot: Arc>>>, audio: Vec, } @@ -242,9 +244,7 @@ impl Task for VadProcessTask { // Moved out rather than borrowed so the buffer can be handed to V8 in `resolve` // without another copy. The task is used once, so leaving an empty Vec behind is fine. let audio = std::mem::take(&mut self.audio); - self - .inner - .with("VadAsync", |inner| map_err(inner.process(&audio)))?; + map_err(lock(&self.slot).get_mut()?.process(&audio))?; Ok(audio) } @@ -258,7 +258,7 @@ impl Task for VadProcessTask { /// Backs {@link VadAsync#getContext}. pub struct VadContextTask { - inner: Arc>>, + slot: Arc>>>, } impl Task for VadContextTask { @@ -266,7 +266,7 @@ impl Task for VadContextTask { type JsValue = VadContext; fn compute(&mut self) -> Result { - self.inner.with("VadAsync", |inner| Ok(inner.context())) + Ok(lock(&self.slot).get()?.context()) } fn resolve(&mut self, _env: Env, context: aic_sdk::VadContext) -> Result { @@ -276,7 +276,7 @@ impl Task for VadContextTask { /// Backs {@link VadAsync#terminateSession}. pub struct VadTerminateTask { - inner: Arc>>, + slot: Arc>>>, } impl Task for VadTerminateTask { @@ -284,9 +284,7 @@ impl Task for VadTerminateTask { type JsValue = (); fn compute(&mut self) -> Result<()> { - self - .inner - .with("VadAsync", |inner| map_err(inner.terminate_session())) + map_err(lock(&self.slot).get_mut()?.terminate_session()) } fn resolve(&mut self, _env: Env, _: ()) -> Result<()> {