diff --git a/packages/ac3/src/decoder.ts b/packages/ac3/src/decoder.ts index 73cc9fe..a676073 100644 --- a/packages/ac3/src/decoder.ts +++ b/packages/ac3/src/decoder.ts @@ -13,7 +13,7 @@ import { EncodedPacket, registerDecoder, } from 'mediabunny'; -import { sendCommand } from './worker-client'; +import { sendCommand, refWorker, unrefWorker } from './worker-client'; class Ac3Decoder extends CustomAudioDecoder { private ctx = 0; @@ -23,6 +23,8 @@ class Ac3Decoder extends CustomAudioDecoder { } async init() { + await refWorker(); + const result = await sendCommand({ type: 'init-decoder', data: { codec: this.codec }, @@ -53,8 +55,9 @@ class Ac3Decoder extends CustomAudioDecoder { await sendCommand({ type: 'flush-decoder', data: { ctx: this.ctx } }); } - close() { + async close() { void sendCommand({ type: 'close-decoder', data: { ctx: this.ctx } }); + await unrefWorker(); } } diff --git a/packages/ac3/src/encoder.ts b/packages/ac3/src/encoder.ts index 1f567a0..2db9dfa 100644 --- a/packages/ac3/src/encoder.ts +++ b/packages/ac3/src/encoder.ts @@ -13,7 +13,7 @@ import { EncodedPacket, registerEncoder, } from 'mediabunny'; -import { sendCommand } from './worker-client'; +import { sendCommand, refWorker, unrefWorker } from './worker-client'; import { assert } from './shared'; import { AC3_SAMPLE_RATES, EAC3_REDUCED_SAMPLE_RATES } from '../../../shared/ac3-misc'; @@ -38,11 +38,14 @@ class Ac3Encoder extends CustomAudioEncoder { return (codec === 'ac3' || codec === 'eac3') && config.numberOfChannels >= 1 && config.numberOfChannels <= 8 - && sampleRates.includes(config.sampleRate); + && sampleRates.includes(config.sampleRate) + && config.bitrate !== undefined; } async init() { - assert(this.config.bitrate); + await refWorker(); + + assert(this.config.bitrate !== undefined); this.sampleRate = this.config.sampleRate; this.numberOfChannels = this.config.numberOfChannels; @@ -130,6 +133,7 @@ class Ac3Encoder extends CustomAudioEncoder { close() { void sendCommand({ type: 'close-encoder', data: { ctx: this.ctx } }); + void unrefWorker(); } private async encodeOneFrame() { diff --git a/packages/ac3/src/worker-client.ts b/packages/ac3/src/worker-client.ts index d79e1ae..74a8bab 100644 --- a/packages/ac3/src/worker-client.ts +++ b/packages/ac3/src/worker-client.ts @@ -10,13 +10,51 @@ import { assert, type WorkerCommand, type WorkerResponse, type WorkerResponseDat // @ts-expect-error An esbuild plugin handles this, TypeScript doesn't need to understand import createWorker from './codec.worker'; -let workerPromise: Promise | null; +type ExtendedWorker = Worker & { + ref?: () => void; + unref?: () => void; +}; + +let workerPromise: Promise | null; let nextMessageId = 0; const pendingMessages = new Map void; reject: (reason?: unknown) => void; }>(); +let refCount = 0; +let keepAliveInterval: ReturnType | null = null; + +export const refWorker = async () => { + refCount++; + if (refCount === 1) { + keepAliveInterval = setInterval(() => {}, 2 ** 31 - 1); + const worker = await ensureWorker(); + worker.ref?.(); + } +}; + +export const unrefWorker = async () => { + refCount--; + if (refCount === 0) { + if (keepAliveInterval !== null) { + clearInterval(keepAliveInterval); + keepAliveInterval = null; + } + + const worker = await workerPromise; + if (worker) { + if (worker.unref) { + worker.unref(); // If we don't do this, then the Node process never terminates by itself + } else if (typeof window === 'undefined') { + // Non-browser environment without unref - terminate instead + worker.terminate(); + workerPromise = null; + } + } + } +}; + export const sendCommand = async ( command: WorkerCommand & { type: T }, transferables?: Transferable[], @@ -41,7 +79,8 @@ export const sendCommand = async ( const ensureWorker = () => { return workerPromise ??= (async () => { // eslint-disable-next-line @typescript-eslint/no-unsafe-call - const worker = (await createWorker()) as Worker; + const worker = (await createWorker()) as ExtendedWorker; + worker.unref?.(); // Start unreffed const onMessage = (data: WorkerResponse) => { const pending = pendingMessages.get(data.id); diff --git a/scripts/esbuild/inlined-workers.ts b/scripts/esbuild/inlined-workers.ts index 0f08ba9..f93098d 100644 --- a/scripts/esbuild/inlined-workers.ts +++ b/scripts/esbuild/inlined-workers.ts @@ -28,7 +28,7 @@ export default function Worker() { const inlineWorkerFunctionCode = ` export default async function inlineWorker(scriptText) { if (typeof Worker !== 'undefined' && typeof Bun === 'undefined') { - // Browser, Deno + // Browser, Deno (Deno can't do dynamic import of worker_threads) const blob = new Blob([scriptText], { type: "text/javascript" }); const url = URL.createObjectURL(blob); @@ -39,10 +39,10 @@ export default async function inlineWorker(scriptText) { // Node, Bun (Bun's Worker is flaky, worker_threads works much better) let Worker; + const workerModule = 'node:worker_threads'; try { - Worker = (await import('worker_threads')).Worker; + Worker = (await import(workerModule)).Worker; } catch { - const workerModule = 'worker_threads'; Worker = require(workerModule).Worker; }