Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions .github/workflows/CI.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
8 changes: 4 additions & 4 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
37 changes: 37 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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+.
Expand Down
228 changes: 228 additions & 0 deletions __test__/dispose.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,228 @@
import test from 'ava'

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() {
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()))
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()
})

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)
processor.dispose()

await t.throwsAsync(() => processor.process(new Float32Array(blockSize)), {
message: /disposed/,
})
await t.throwsAsync(() => processor.terminateSession(), { message: /disposed/ })
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)
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()
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.
const 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)
processor.dispose()
t.pass()
})
2 changes: 1 addition & 1 deletion __test__/index.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
})
Expand Down
4 changes: 2 additions & 2 deletions __test__/models.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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' },
Expand Down
Loading