From e691b3a5da762c7d34cd8949313f3d11e0ec2bde Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Fri, 10 Apr 2026 13:04:51 -0400 Subject: [PATCH] feat: streaming pipeline chaining and task-arg support Co-Authored-By: Claude Opus 4.6 (1M context) --- src/execute.ts | 11 ++++++++++- test/stream-context.test.ts | 29 +++++++++++++++++++++++++++++ test/stream-pipeline.test.ts | 35 +++++++++++++++++++++++++++++++++++ 3 files changed, 74 insertions(+), 1 deletion(-) create mode 100644 test/stream-context.test.ts create mode 100644 test/stream-pipeline.test.ts diff --git a/src/execute.ts b/src/execute.ts index cd636c2..44b7e88 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -5,6 +5,8 @@ import { freezeModule } from './registry.ts'; import { serializeArg, deserializeArg } from './shared/reconstruct.ts'; import { extractTransferables, collectTransferables } from './transfer.ts'; import { Task } from './task.ts'; +import { StreamTask } from './stream-task.ts'; +import { runStreamOnDedicated } from './dedicated-runner.ts'; import { STREAM } from './stream.ts'; import type { StreamOptions } from './stream.ts'; @@ -69,10 +71,17 @@ function pipeToPort(iterable: AsyncIterable, port: MessagePort, highWat function prepareArg(arg: unknown): unknown { if (typeof arg === 'object' && arg !== null && STREAM in arg) { const data = (arg as any)[STREAM]; + let iterable = data.iterable; const highWater = data.options?.highWaterMark ?? DEFAULT_HIGH_WATER; + + // If the iterable is a StreamTask, dispatch it to its dedicated worker first + if (iterable instanceof StreamTask) { + iterable = runStreamOnDedicated(iterable.id, iterable.args); + } + const { port1, port2 } = new MessageChannel(); port1.unref(); - pipeToPort(data.iterable, port1, highWater); + pipeToPort(iterable, port1, highWater); streamPortStack[streamPortStack.length - 1].push(port2); return port2; } diff --git a/test/stream-context.test.ts b/test/stream-context.test.ts new file mode 100644 index 0000000..d1a4f97 --- /dev/null +++ b/test/stream-context.test.ts @@ -0,0 +1,29 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { mo, workers } from 'moroutine'; + +const makeMultiplier = mo(import.meta, (factor: number): number => { + return factor; +}); + +const streamMultiplied = mo(import.meta, async function* (factor: number, count: number) { + for (let i = 0; i < count; i++) { + yield i * factor; + } +}); + +describe('streaming with task-args', () => { + it('resolves task-args before streaming', async () => { + const factor = makeMultiplier(3); + const run = workers(1); + try { + const results: number[] = []; + for await (const value of run(streamMultiplied(factor, 4))) { + results.push(value); + } + assert.deepEqual(results, [0, 3, 6, 9]); + } finally { + run[Symbol.dispose](); + } + }); +}); diff --git a/test/stream-pipeline.test.ts b/test/stream-pipeline.test.ts new file mode 100644 index 0000000..e79d8df --- /dev/null +++ b/test/stream-pipeline.test.ts @@ -0,0 +1,35 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { mo, stream } from 'moroutine'; + +const generate = mo(import.meta, async function* (n: number) { + for (let i = 1; i <= n; i++) yield i; +}); + +const double = mo(import.meta, async function* (input: AsyncIterable) { + for await (const n of input) yield n * 2; +}); + +const square = mo(import.meta, async function* (input: AsyncIterable) { + for await (const n of input) yield n * n; +}); + +describe('streaming pipeline', () => { + it('chains two streaming moroutines', async () => { + const results: number[] = []; + for await (const value of double(stream(generate(3)))) { + results.push(value); + } + assert.deepEqual(results, [2, 4, 6]); + }); + + it('chains three streaming moroutines', async () => { + const results: number[] = []; + const doubled = double(stream(generate(3))); + const squared = square(stream(doubled)); + for await (const value of squared) { + results.push(value); + } + assert.deepEqual(results, [4, 16, 36]); + }); +}); -- 2.51.2