diff --git a/README.md b/README.md index d0ac526..d346f4f 100644 --- a/README.md +++ b/README.md @@ -219,6 +219,90 @@ const pos = shared({ x: int32, y: int32 }); console.log(pos.load()); // { x: 1, y: 1 } ``` +## Streaming + +### Streaming Moroutines + +Wrap an `async function*` with `mo()` to create a streaming moroutine. Values are streamed between threads via `MessageChannel` with pause/resume backpressure. + +```ts +// count.ts +import { mo } from 'moroutine'; + +export const countUp = mo(import.meta, async function* (n: number) { + for (let i = 0; i < n; i++) { + yield i; + } +}); +``` + +Iterate directly (dedicated worker) or dispatch via a pool: + +```ts +import { workers } from 'moroutine'; +import { countUp } from './count.ts'; + +// Dedicated worker +for await (const n of countUp(5)) { + console.log(n); // 0, 1, 2, 3, 4 +} + +// Worker pool +{ + using run = workers(2); + for await (const n of run(countUp(5))) { + console.log(n); // 0, 1, 2, 3, 4 + } +} +``` + +### `channel()` and Fan-out + +When you pass the same `AsyncIterable` or `StreamTask` argument to multiple tasks, each task gets its own copy of the data. Use `channel()` to share a single source across multiple workers — each item goes to exactly one consumer (work stealing). + +```ts +import { workers, channel, mo } from 'moroutine'; + +const generate = mo(import.meta, async function* (n: number) { + for (let i = 0; i < n; i++) yield i; +}); + +const process = mo(import.meta, async (input: AsyncIterable): Promise => { + const items: number[] = []; + for await (const n of input) items.push(n); + return items; +}); +``` + +```ts +const data = channel(generate(100)); + +{ + using run = workers(4); + const [a, b, c, d] = await run([ + process(data), + process(data), + process(data), + process(data), + ]); + // Items distributed across workers — no duplicates, no gaps +} +``` + +Without `channel()`, `AsyncIterable` and `StreamTask` arguments are auto-detected and streamed to a single consumer. `channel()` is only needed for fan-out. + +### Pipelines + +Chain streaming moroutines by passing one as an argument to the next. Each stage runs on its own dedicated worker. + +```ts +const doubled = double(generate(5)); +const squared = square(doubled); +for await (const n of squared) { + console.log(n); +} +``` + ## Transfers Use `transfer()` for zero-copy movement of `ArrayBuffer`, `TypedArray`, `MessagePort`, or streams. diff --git a/src/channel.ts b/src/channel.ts index 6120184..9ace508 100644 --- a/src/channel.ts +++ b/src/channel.ts @@ -7,7 +7,9 @@ import { runStreamOnDedicated } from './dedicated-runner.ts'; const CHANNEL = Symbol.for('moroutine.channel'); +/** Options for configuring a channel. */ export interface ChannelOptions { + /** Maximum number of items buffered before signaling backpressure. Defaults to 16. */ highWaterMark?: number; } diff --git a/src/mo.ts b/src/mo.ts index 2f493da..587563c 100644 --- a/src/mo.ts +++ b/src/mo.ts @@ -8,8 +8,9 @@ const counters = new Map(); * Wraps a function to run on a worker thread. Must be called at module scope. * @param importMeta - The `import.meta` of the calling module, used to identify the source file. * @param fn - The function to offload to a worker thread. - * @returns A function that creates a {@link Task} when called. + * @returns A callable that creates a {@link Task} (or {@link StreamTask} for `async function*`) when invoked. */ + /** A value or a Task that resolves to that value on the worker. */ export type Arg = T | Task; diff --git a/src/runner.ts b/src/runner.ts index b99cc27..dec649f 100644 --- a/src/runner.ts +++ b/src/runner.ts @@ -4,10 +4,19 @@ import type { ChannelOptions } from './channel.ts'; type TaskResults[]> = { [K in keyof T]: T[K] extends Task ? R : never }; -/** A callable that dispatches tasks to a worker pool. Disposable via `using` or `[Symbol.dispose]()`. */ +/** + * A callable that dispatches tasks to a worker pool. Disposable via `using` or `[Symbol.dispose]()`. + * + * @param task - A single {@link Task} to run on a worker. + * @returns `Promise` for a single task, `Promise<[...results]>` for a batch, or `AsyncIterable` for a streaming task. + */ export type Runner = { + /** Dispatches a single task and returns its result. */ (task: Task): Promise; + /** Dispatches a batch of tasks in parallel and returns all results. */ []>(tasks: [...T]): Promise>; + /** Dispatches a streaming task and returns an async iterable of yielded values. */ (task: StreamTask, opts?: ChannelOptions): AsyncIterable; + /** Terminates all workers in the pool. */ [Symbol.dispose](): void; }; diff --git a/src/stream-task.ts b/src/stream-task.ts index 432bf14..a9c186d 100644 --- a/src/stream-task.ts +++ b/src/stream-task.ts @@ -2,18 +2,25 @@ import { runStreamOnDedicated } from './dedicated-runner.ts'; let nextUid = 0; -/** A deferred streaming computation. When dispatched, returns an AsyncIterable instead of a Promise. */ +/** + * A deferred streaming computation. When dispatched via a {@link Runner} or iterated directly, + * returns an `AsyncIterable` of yielded values instead of a `Promise`. + * Created by calling an `async function*` wrapped with {@link mo}. + */ export class StreamTask { readonly uid: number; readonly id: string; readonly args: unknown[]; + /** @param id - The moroutine identifier (module URL + index). + * @param args - The arguments to pass to the worker generator function. */ constructor(id: string, args: unknown[]) { this.uid = nextUid++; this.id = id; this.args = args; } + /** Enables `for await...of` by dispatching to a dedicated worker. @returns An iterator of yielded values. */ [Symbol.asyncIterator](): AsyncIterator { const iterable = runStreamOnDedicated(this.id, this.args); return iterable[Symbol.asyncIterator](); diff --git a/src/task.ts b/src/task.ts index 2057237..445ef41 100644 --- a/src/task.ts +++ b/src/task.ts @@ -2,18 +2,24 @@ import { runOnDedicated } from './dedicated-runner.ts'; let nextUid = 0; -/** A deferred computation that runs on a worker thread when awaited. */ +/** + * A deferred computation that runs on a worker thread when awaited or dispatched via a {@link Runner}. + * Created by calling a function wrapped with {@link mo}. + */ export class Task { readonly uid: number; readonly id: string; readonly args: unknown[]; + /** @param id - The moroutine identifier (module URL + index). + * @param args - The arguments to pass to the worker function. */ constructor(id: string, args: unknown[]) { this.uid = nextUid++; this.id = id; this.args = args; } + /** Enables `await task` by dispatching to a dedicated worker. @returns The worker function's result. */ then( onfulfilled?: ((value: T) => T1 | PromiseLike) | null, onrejected?: ((reason: any) => T2 | PromiseLike) | null,