diff --git a/examples/pipeline/main.ts b/examples/pipeline/main.ts new file mode 100644 index 0000000..76c331d --- /dev/null +++ b/examples/pipeline/main.ts @@ -0,0 +1,22 @@ +// Pipeline of streaming moroutines, each step on its own dedicated worker. +// Data flows: generate → double → square → toString, all via MessageChannel. +// Requires Node v24+. +// +// Run: node examples/pipeline/main.ts + +import { stream } from '../../src/index.ts'; +import { generate, double, square, toString } from './steps.ts'; + +const numbers = generate(5); +const doubled = double(stream(numbers)); +const squared = square(stream(doubled)); +const labels = toString(stream(squared)); + +for await (const label of labels) { + console.log(label); +} +// => 4 +// => 16 +// => 36 +// => 64 +// => 100 diff --git a/examples/pipeline/steps.ts b/examples/pipeline/steps.ts new file mode 100644 index 0000000..34d05a0 --- /dev/null +++ b/examples/pipeline/steps.ts @@ -0,0 +1,25 @@ +import { mo } from '../../src/index.ts'; + +export const generate = mo(import.meta, async function* (count: number) { + for (let i = 1; i <= count; i++) { + yield i; + } +}); + +export const double = mo(import.meta, async function* (input: AsyncIterable) { + for await (const n of input) { + yield n * 2; + } +}); + +export const square = mo(import.meta, async function* (input: AsyncIterable) { + for await (const n of input) { + yield n * n; + } +}); + +export const toString = mo(import.meta, async function* (input: AsyncIterable) { + for await (const n of input) { + yield `=> ${n}`; + } +}); diff --git a/src/dedicated-runner.ts b/src/dedicated-runner.ts index 5928c29..8b11058 100644 --- a/src/dedicated-runner.ts +++ b/src/dedicated-runner.ts @@ -8,8 +8,8 @@ function getWorker(id: string): Worker { let worker = workers.get(id); if (!worker) { worker = new Worker(workerEntryUrl); - worker.unref(); setupWorker(worker); + worker.unref(); workers.set(id, worker); } return worker; diff --git a/src/execute.ts b/src/execute.ts index 44b7e88..2250aa8 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -19,6 +19,8 @@ export function setupWorker(worker: Worker): void { const call = pending.get(msg.callId); if (!call) return; pending.delete(msg.callId); + // Unref the worker now that this call is done (re-ref happened in execute()) + worker.unref(); if (msg.error !== undefined) { call.reject(new Error(msg.error)); } else { @@ -95,6 +97,8 @@ export function execute(worker: Worker, id: string, args: unknown[]): Promise const url = id.slice(0, id.lastIndexOf('#')); freezeModule(url); const callId = nextCallId++; + // Ref the worker for the duration of this call so the event loop stays alive + worker.ref(); return new Promise((resolve, reject) => { pending.set(callId, { resolve, reject }); const extracted = extractTransferables(args); @@ -112,13 +116,14 @@ export function dispatchStream(worker: Worker, id: string, args: unknown[], o const highWater = opts?.highWaterMark ?? DEFAULT_HIGH_WATER; const { port1, port2 } = new MessageChannel(); - port1.unref(); const extracted = extractTransferables(args); streamPortStack.push([]); const preparedArgs = extracted.args.map(prepareArg); const ports = streamPortStack.pop()!; const msg = { id, args: preparedArgs, port: port2 }; + // Ref the worker for the duration of the stream so the event loop stays alive + worker.ref(); worker.postMessage(msg, [...extracted.transfer, ...ports, port2] as any[]); const queue: T[] = []; @@ -126,18 +131,25 @@ export function dispatchStream(worker: Worker, id: string, args: unknown[], o let error: Error | null = null; let paused = false; let waiting: (() => void) | null = null; + let workerUnrefed = false; + + function unrefWorkerOnce() { + if (!workerUnrefed) { workerUnrefed = true; worker.unref(); } + } port1.on('message', (msg: { value?: unknown; done?: boolean; error?: string }) => { if (msg.error) { error = new Error(msg.error); done = true; port1.close(); + unrefWorkerOnce(); if (waiting) { waiting(); waiting = null; } return; } if (msg.done) { done = true; port1.close(); + unrefWorkerOnce(); if (waiting) { waiting(); waiting = null; } return; } @@ -148,6 +160,7 @@ export function dispatchStream(worker: Worker, id: string, args: unknown[], o port1.postMessage('pause'); } }); + port1.unref(); return { [Symbol.asyncIterator]() { @@ -169,6 +182,7 @@ export function dispatchStream(worker: Worker, id: string, args: unknown[], o }, async return(): Promise> { port1.close(); + unrefWorkerOnce(); return { done: true, value: undefined }; }, }; diff --git a/src/worker-pool.ts b/src/worker-pool.ts index 4526310..b7eccd8 100644 --- a/src/worker-pool.ts +++ b/src/worker-pool.ts @@ -17,8 +17,8 @@ export function workers(size: number = availableParallelism()): Runner { const pool: Worker[] = []; for (let i = 0; i < size; i++) { const worker = new Worker(workerEntryUrl); - worker.unref(); setupWorker(worker); + worker.unref(); pool.push(worker); } diff --git a/test/fixtures/exit-dedicated-main.ts b/test/fixtures/exit-dedicated-main.ts new file mode 100644 index 0000000..de0cc7b --- /dev/null +++ b/test/fixtures/exit-dedicated-main.ts @@ -0,0 +1,4 @@ +import { gen } from './exit-gen.ts'; + +for await (const v of gen()) {} +console.log('DONE'); diff --git a/test/fixtures/exit-gen.ts b/test/fixtures/exit-gen.ts new file mode 100644 index 0000000..619e797 --- /dev/null +++ b/test/fixtures/exit-gen.ts @@ -0,0 +1,3 @@ +import { mo } from 'moroutine'; + +export const gen = mo(import.meta, async function* () { yield 1; yield 2; }); diff --git a/test/fixtures/exit-pool-main.ts b/test/fixtures/exit-pool-main.ts new file mode 100644 index 0000000..a3bd2d8 --- /dev/null +++ b/test/fixtures/exit-pool-main.ts @@ -0,0 +1,7 @@ +import { workers } from 'moroutine'; +import { gen } from './exit-gen.ts'; + +const run = workers(1); +for await (const v of run(gen())) {} +run[Symbol.dispose](); +console.log('DONE'); diff --git a/test/stream-exit.test.ts b/test/stream-exit.test.ts new file mode 100644 index 0000000..5e9df20 --- /dev/null +++ b/test/stream-exit.test.ts @@ -0,0 +1,29 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { execFile } from 'node:child_process'; +import { promisify } from 'node:util'; +import { fileURLToPath } from 'node:url'; +import { join } from 'node:path'; + +const exec = promisify(execFile); +const fixturesDir = join(fileURLToPath(import.meta.url), '..', 'fixtures'); + +describe('streaming process exit', () => { + it('process exits after streaming on dedicated worker', async () => { + const { stdout } = await exec(process.execPath, [ + '--no-warnings', + '--experimental-strip-types', + join(fixturesDir, 'exit-dedicated-main.ts'), + ], { timeout: 5000 }); + assert.ok(stdout.includes('DONE')); + }); + + it('process exits after streaming with pool', async () => { + const { stdout } = await exec(process.execPath, [ + '--no-warnings', + '--experimental-strip-types', + join(fixturesDir, 'exit-pool-main.ts'), + ], { timeout: 5000 }); + assert.ok(stdout.includes('DONE')); + }); +});