From 655d5ae6425f663fb66594ba010650cce871cb46 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Mon, 13 Apr 2026 23:12:57 -0400 Subject: [PATCH] feat: WorkerHandle, run.workers, and assign() dispatch support Co-Authored-By: Claude Opus 4.6 (1M context) --- src/worker-pool.ts | 20 +++++++--- test/fixtures/worker-handle.ts | 7 ++++ test/worker-handle.test.ts | 68 ++++++++++++++++++++++++++++++++++ 3 files changed, 89 insertions(+), 6 deletions(-) create mode 100644 test/fixtures/worker-handle.ts create mode 100644 test/worker-handle.test.ts diff --git a/src/worker-pool.ts b/src/worker-pool.ts index b0fa65a..2c3e374 100644 --- a/src/worker-pool.ts +++ b/src/worker-pool.ts @@ -2,7 +2,7 @@ import { setTimeout } from 'node:timers/promises'; import { Worker } from 'node:worker_threads'; import { availableParallelism } from 'node:os'; import { setupWorker, execute, dispatchStream } from './execute.ts'; -import type { Task } from './task.ts'; +import { Task } from './task.ts'; import { StreamTask } from './stream-task.ts'; import type { ChannelOptions } from './channel.ts'; import type { Runner, WorkerHandle, WorkerOptions } from './runner.ts'; @@ -42,10 +42,19 @@ export function workers(size: number = availableParallelism(), opts?: WorkerOpti pool.length = 0; } - function dispatch(task: Task): Promise { - if (disposed) return Promise.reject(new Error('Worker pool is disposed')); + function resolveWorker(task: Task | StreamTask): Worker { + if (task.worker != null) { + const idx = workerHandles.indexOf(task.worker as WorkerHandle); + if (idx !== -1) return pool[idx]; + } const worker = pool[next % pool.length]; next++; + return worker; + } + + function dispatch(task: Task): Promise { + if (disposed) return Promise.reject(new Error('Worker pool is disposed')); + const worker = resolveWorker(task); return track(execute(worker, task.id, task.args)); } @@ -64,14 +73,13 @@ export function workers(size: number = availableParallelism(), opts?: WorkerOpti }; } - const workerHandles: readonly WorkerHandle[] = pool.map(makeWorkerHandle); + const workerHandles: readonly WorkerHandle[] = Object.freeze(pool.map(makeWorkerHandle)); const run: Runner = Object.assign( (taskOrTasks: Task | Task[] | StreamTask, channelOpts?: ChannelOptions): any => { if (taskOrTasks instanceof StreamTask) { if (disposed) throw new Error('Worker pool is disposed'); - const worker = pool[next % pool.length]; - next++; + const worker = resolveWorker(taskOrTasks); const { iterable, done } = dispatchStream(worker, taskOrTasks.id, taskOrTasks.args, channelOpts); track(done); return iterable; diff --git a/test/fixtures/worker-handle.ts b/test/fixtures/worker-handle.ts new file mode 100644 index 0000000..a72de2d --- /dev/null +++ b/test/fixtures/worker-handle.ts @@ -0,0 +1,7 @@ +import { mo } from 'moroutine'; + +export const identity = mo(import.meta, (n: number): number => n); + +export const countUp = mo(import.meta, async function* (n: number) { + for (let i = 0; i < n; i++) yield i; +}); diff --git a/test/worker-handle.test.ts b/test/worker-handle.test.ts new file mode 100644 index 0000000..a41e842 --- /dev/null +++ b/test/worker-handle.test.ts @@ -0,0 +1,68 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { workers, assign, channel } from 'moroutine'; +import { identity, countUp } from './fixtures/worker-handle.ts'; + +describe('WorkerHandle', () => { + it('run.workers is a frozen array matching pool size', () => { + const run = workers(3); + try { + assert.equal(run.workers.length, 3); + assert.ok(Object.isFrozen(run.workers)); + } finally { + run[Symbol.dispose](); + } + }); + + it('w.exec() dispatches a task to the specific worker', async () => { + const run = workers(2); + try { + const result = await run.workers[0].exec(identity(42)); + assert.equal(result, 42); + } finally { + run[Symbol.dispose](); + } + }); + + it('w.exec() dispatches a streaming task', async () => { + const run = workers(1); + try { + const results: number[] = []; + for await (const n of run.workers[0].exec(countUp(3))) { + results.push(n); + } + assert.deepEqual(results, [0, 1, 2]); + } finally { + run[Symbol.dispose](); + } + }); + + it('assign() pins task to a specific worker via run()', async () => { + const run = workers(2); + try { + const result = await run(assign(run.workers[1], identity(99))); + assert.equal(result, 99); + } finally { + run[Symbol.dispose](); + } + }); + + it('assign() works in a batch', async () => { + const run = workers(2); + try { + const results = await run([ + assign(run.workers[0], identity(1)), + assign(run.workers[1], identity(2)), + ]); + assert.deepEqual(results, [1, 2]); + } finally { + run[Symbol.dispose](); + } + }); + + it('w.exec() rejects after dispose', async () => { + const run = workers(1); + run[Symbol.dispose](); + await assert.rejects(() => run.workers[0].exec(identity(1)), { message: /disposed/ }); + }); +}); -- 2.51.2