diff --git a/docs/superpowers/plans/2026-04-14-load-balancing.md b/docs/superpowers/plans/2026-04-14-load-balancing.md new file mode 100644 index 0000000..e10c1bf --- /dev/null +++ b/docs/superpowers/plans/2026-04-14-load-balancing.md @@ -0,0 +1,637 @@ +# Extensible Load Balancing Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Support pluggable load balancing strategies for worker pools, with built-in round-robin and least-busy implementations. + +**Architecture:** A `Balancer` interface with a `select()` method chooses which worker runs each task. `WorkerHandle` exposes `thread` and `activeCount` for balancers to read. Built-in `roundRobin()` and `leastBusy()` are exported as factories. The `workers()` function accepts overloaded signatures so opts can be passed without size. + +**Tech Stack:** TypeScript (erasable syntax only), worker_threads, node:test. + +--- + +### Task 1: Update Types — Balancer, WorkerHandle, WorkerOptions, workers() overloads + +**Files:** +- Modify: `src/runner.ts` + +- [ ] **Step 1: Update runner.ts** + +Replace the entire contents of `src/runner.ts` with: + +```ts +import type { Worker } from 'node:worker_threads'; +import type { Task } from './task.ts'; +import type { StreamTask } from './stream-task.ts'; +import type { ChannelOptions } from './channel.ts'; + +type TaskResults[]> = { [K in keyof T]: T[K] extends Task ? R : never }; + +/** A load balancing strategy for choosing which worker runs a task. */ +export interface Balancer { + /** Choose a worker for the given task. Called synchronously on every dispatch. */ + select(workers: readonly WorkerHandle[], task: Task | StreamTask): WorkerHandle; + /** Optional cleanup on sync dispose. */ + [Symbol.dispose]?(): void; + /** Optional cleanup on async dispose. */ + [Symbol.asyncDispose]?(): Promise; +} + +/** Options for configuring a worker pool. */ +export interface WorkerOptions { + /** Maximum time in ms to wait for in-flight tasks during async dispose. If exceeded, workers are force-terminated. */ + shutdownTimeout?: number; + /** Load balancing strategy. Defaults to round-robin. */ + balance?: Balancer; +} + +/** A handle to a specific worker in a pool. */ +export interface WorkerHandle { + /** Dispatches a task pinned to this worker. */ + exec(task: Task): Promise; + /** Dispatches a streaming task pinned to this worker. */ + exec(task: StreamTask, opts?: ChannelOptions): AsyncIterable; + /** The underlying worker thread. */ + readonly thread: Worker; + /** Number of currently in-flight tasks on this worker. */ + readonly activeCount: number; +} + +/** + * 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; + /** AbortSignal that fires when the pool starts disposing. */ + readonly signal: AbortSignal; + /** Read-only array of worker handles, one per pool worker. */ + readonly workers: readonly WorkerHandle[]; + /** Terminates all workers immediately. */ + [Symbol.dispose](): void; + /** Aborts signal, waits for in-flight tasks to settle, then terminates workers. */ + [Symbol.asyncDispose](): Promise; +}; +``` + +- [ ] **Step 2: Run type check** + +Run: `pnpm tsc --noEmit 2>&1` +Expected: Errors in `worker-pool.ts` because `WorkerHandle` now requires `thread` and `activeCount`. That's expected — Task 3 fixes it. + +- [ ] **Step 3: Commit** + +```bash +git add src/runner.ts +git commit -m "feat: add Balancer interface, thread and activeCount to WorkerHandle" +``` + +--- + +### Task 2: Create Built-in Balancers + +**Files:** +- Create: `src/balancers.ts` +- Modify: `src/index.ts` +- Create: `test/balancers.test.ts` + +- [ ] **Step 1: Write failing tests** + +Create `test/balancers.test.ts`: + +```ts +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { roundRobin, leastBusy } from 'moroutine'; +import type { WorkerHandle } from 'moroutine'; + +function mockHandle(activeCount: number): WorkerHandle { + return { activeCount, thread: {} as any, exec: {} as any }; +} + +describe('roundRobin()', () => { + it('cycles through workers in order', () => { + const b = roundRobin(); + const handles = [mockHandle(0), mockHandle(0), mockHandle(0)]; + const task = { id: 'test', args: [], uid: 0 } as any; + assert.equal(b.select(handles, task), handles[0]); + assert.equal(b.select(handles, task), handles[1]); + assert.equal(b.select(handles, task), handles[2]); + assert.equal(b.select(handles, task), handles[0]); + }); +}); + +describe('leastBusy()', () => { + it('picks worker with lowest activeCount', () => { + const b = leastBusy(); + const handles = [mockHandle(3), mockHandle(1), mockHandle(2)]; + const task = { id: 'test', args: [], uid: 0 } as any; + assert.equal(b.select(handles, task), handles[1]); + }); + + it('breaks ties by index', () => { + const b = leastBusy(); + const handles = [mockHandle(1), mockHandle(1), mockHandle(1)]; + const task = { id: 'test', args: [], uid: 0 } as any; + assert.equal(b.select(handles, task), handles[0]); + }); +}); +``` + +- [ ] **Step 2: Run test to verify it fails** + +Run: `node --no-warnings --test test/balancers.test.ts 2>&1 | tail -10` +Expected: FAIL — `roundRobin` and `leastBusy` not exported. + +- [ ] **Step 3: Implement balancers** + +Create `src/balancers.ts`: + +```ts +import type { Balancer, WorkerHandle } from './runner.ts'; +import type { Task } from './task.ts'; +import type { StreamTask } from './stream-task.ts'; + +/** + * Creates a round-robin balancer that cycles through workers in order. + * @returns A fresh Balancer instance. + */ +export function roundRobin(): Balancer { + let next = 0; + return { + select(workers: readonly WorkerHandle[]): WorkerHandle { + const worker = workers[next % workers.length]; + next++; + return worker; + }, + }; +} + +/** + * Creates a least-busy balancer that picks the worker with the lowest activeCount. + * Ties are broken by index (first wins). + * @returns A fresh Balancer instance. + */ +export function leastBusy(): Balancer { + return { + select(workers: readonly WorkerHandle[]): WorkerHandle { + let best = workers[0]; + for (let i = 1; i < workers.length; i++) { + if (workers[i].activeCount < best.activeCount) { + best = workers[i]; + } + } + return best; + }, + }; +} +``` + +- [ ] **Step 4: Export from index.ts** + +Add to `src/index.ts`, after the `assign` export: + +```ts +export { roundRobin, leastBusy } from './balancers.ts'; +``` + +Also add `Balancer` to the type export: + +Change: +```ts +export type { Runner, WorkerHandle, WorkerOptions } from './runner.ts'; +``` +to: +```ts +export type { Balancer, Runner, WorkerHandle, WorkerOptions } from './runner.ts'; +``` + +- [ ] **Step 5: Run tests** + +Run: `node --no-warnings --test test/balancers.test.ts 2>&1 | tail -10` +Expected: All tests pass. + +- [ ] **Step 6: Commit** + +```bash +git add src/balancers.ts src/index.ts test/balancers.test.ts +git commit -m "feat: roundRobin() and leastBusy() balancer factories" +``` + +--- + +### Task 3: Integrate Balancer into worker-pool.ts + +Wire balancers into dispatch, add `thread` and `activeCount` to WorkerHandle, overload `workers()` signature, and dispose the balancer. + +**Files:** +- Modify: `src/worker-pool.ts` +- Create: `test/load-balancing.test.ts` +- Create: `test/fixtures/load-balancing.ts` + +- [ ] **Step 1: Create test fixture** + +Create `test/fixtures/load-balancing.ts`: + +```ts +import { setTimeout } from 'node:timers/promises'; +import { mo } from 'moroutine'; + +export const identity = mo(import.meta, (n: number): number => n); + +export const slow = mo(import.meta, async (ms: number): Promise => { + await setTimeout(ms); + return 'done'; +}); +``` + +- [ ] **Step 2: Write tests** + +Create `test/load-balancing.test.ts`: + +```ts +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { workers, leastBusy } from 'moroutine'; +import type { Balancer, WorkerHandle } from 'moroutine'; +import { identity, slow } from './fixtures/load-balancing.ts'; + +describe('load balancing', () => { + it('defaults to round-robin', async () => { + const run = workers(2); + try { + const result = await run(identity(42)); + assert.equal(result, 42); + } finally { + run[Symbol.dispose](); + } + }); + + it('accepts a balancer option', async () => { + const run = workers(2, { balance: leastBusy() }); + try { + const result = await run(identity(42)); + assert.equal(result, 42); + } finally { + run[Symbol.dispose](); + } + }); + + it('accepts opts without size', async () => { + const run = workers({ balance: leastBusy() }); + try { + assert.ok(run.workers.length > 0); + const result = await run(identity(42)); + assert.equal(result, 42); + } finally { + run[Symbol.dispose](); + } + }); + + it('exposes thread on WorkerHandle', () => { + const run = workers(1); + try { + assert.ok(run.workers[0].thread); + assert.equal(typeof run.workers[0].thread.threadId, 'number'); + } finally { + run[Symbol.dispose](); + } + }); + + it('tracks activeCount on WorkerHandle', async () => { + const run = workers(1); + try { + assert.equal(run.workers[0].activeCount, 0); + const promise = run(slow(100)); + assert.equal(run.workers[0].activeCount, 1); + await promise; + assert.equal(run.workers[0].activeCount, 0); + } finally { + run[Symbol.dispose](); + } + }); + + it('custom balancer receives task', async () => { + const seen: string[] = []; + const custom: Balancer = { + select(workers, task) { + seen.push(task.id); + return workers[0]; + }, + }; + const run = workers(1, { balance: custom }); + try { + await run(identity(1)); + assert.equal(seen.length, 1); + assert.ok(seen[0].includes('#')); + } finally { + run[Symbol.dispose](); + } + }); + + it('balancer dispose is called on sync dispose', () => { + let disposed = false; + const custom: Balancer = { + select(workers) { return workers[0]; }, + [Symbol.dispose]() { disposed = true; }, + }; + const run = workers(1, { balance: custom }); + run[Symbol.dispose](); + assert.ok(disposed); + }); + + it('balancer asyncDispose is called on async dispose', async () => { + let disposed = false; + const custom: Balancer = { + select(workers) { return workers[0]; }, + async [Symbol.asyncDispose]() { disposed = true; }, + }; + const run = workers(1, { balance: custom }); + await run[Symbol.asyncDispose](); + assert.ok(disposed); + }); + + it('pinned tasks bypass balancer', async () => { + let called = false; + const custom: Balancer = { + select(workers) { called = true; return workers[0]; }, + }; + const run = workers(2, { balance: custom }); + try { + const { assign } = await import('moroutine'); + await run(assign(run.workers[1], identity(42))); + assert.ok(!called); + } finally { + run[Symbol.dispose](); + } + }); +}); +``` + +- [ ] **Step 3: Run tests to verify they fail** + +Run: `node --no-warnings --test test/load-balancing.test.ts 2>&1 | tail -10` +Expected: FAIL — `workers()` doesn't accept opts as first arg, `thread`/`activeCount` not on handle. + +- [ ] **Step 4: Update worker-pool.ts** + +Replace `src/worker-pool.ts` with: + +```ts +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 { Task } from './task.ts'; +import { StreamTask } from './stream-task.ts'; +import { roundRobin } from './balancers.ts'; +import type { ChannelOptions } from './channel.ts'; +import type { Balancer, Runner, WorkerHandle, WorkerOptions } from './runner.ts'; + +const workerEntryUrl = new URL('./worker-entry.ts', import.meta.url); + +/** + * Creates a pool of worker threads with configurable load balancing. + * @param size - Number of worker threads, or options object. Defaults to `os.availableParallelism()`. + * @param opts - Optional configuration including shutdown timeout and balancer. + * @returns A disposable {@link Runner} for dispatching tasks. + */ +export function workers(sizeOrOpts?: number | WorkerOptions, opts?: WorkerOptions): Runner { + let size: number; + if (typeof sizeOrOpts === 'object') { + opts = sizeOrOpts; + size = availableParallelism(); + } else { + size = sizeOrOpts ?? availableParallelism(); + } + + const balancer: Balancer = opts?.balance ?? roundRobin(); + + const pool: Worker[] = []; + for (let i = 0; i < size; i++) { + const worker = new Worker(workerEntryUrl); + setupWorker(worker); + worker.unref(); + pool.push(worker); + } + + let disposed = false; + const ac = new AbortController(); + const inflight = new Set>(); + const activeCounts = new Map(); + + function track(handle: WorkerHandle, promise: Promise): Promise { + inflight.add(promise); + activeCounts.set(handle, (activeCounts.get(handle) ?? 0) + 1); + promise.finally(() => { + inflight.delete(promise); + activeCounts.set(handle, (activeCounts.get(handle) ?? 1) - 1); + }).catch(() => {}); + return promise; + } + + function terminateAll(): void { + for (const worker of pool) { + worker.terminate(); + } + pool.length = 0; + } + + function disposeBalancer(): void { + if (Symbol.dispose in balancer) { + balancer[Symbol.dispose]!(); + } else if (Symbol.asyncDispose in balancer) { + (balancer[Symbol.asyncDispose]! as () => Promise)(); + } + } + + async function asyncDisposeBalancer(): Promise { + if (Symbol.asyncDispose in balancer) { + await balancer[Symbol.asyncDispose]!(); + } else if (Symbol.dispose in balancer) { + balancer[Symbol.dispose]!(); + } + } + + function resolveWorkerAndHandle(task: Task | StreamTask): { worker: Worker; handle: WorkerHandle } { + if (task.worker != null) { + const idx = workerHandles.indexOf(task.worker as WorkerHandle); + if (idx !== -1) return { worker: pool[idx], handle: workerHandles[idx] }; + } + const handle = balancer.select(workerHandles, task); + const idx = workerHandles.indexOf(handle); + return { worker: pool[idx], handle }; + } + + function dispatch(task: Task): Promise { + if (disposed) return Promise.reject(new Error('Worker pool is disposed')); + const { worker, handle } = resolveWorkerAndHandle(task); + return track(handle, execute(worker, task.id, task.args)); + } + + function makeWorkerHandle(worker: Worker, idx: number): WorkerHandle { + let handle: WorkerHandle; + handle = { + exec(task: Task | StreamTask, channelOpts?: ChannelOptions): any { + if (task instanceof StreamTask) { + if (disposed) throw new Error('Worker pool is disposed'); + const { iterable, done } = dispatchStream(worker, task.id, task.args, channelOpts); + track(handle, done); + return iterable; + } + if (disposed) return Promise.reject(new Error('Worker pool is disposed')); + return track(handle, execute(worker, task.id, task.args)); + }, + get thread() { + return worker; + }, + get activeCount() { + return activeCounts.get(handle) ?? 0; + }, + }; + return handle; + } + + 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, handle } = resolveWorkerAndHandle(taskOrTasks); + const { iterable, done } = dispatchStream(worker, taskOrTasks.id, taskOrTasks.args, channelOpts); + track(handle, done); + return iterable; + } + if (Array.isArray(taskOrTasks)) { + return Promise.all(taskOrTasks.map((t) => dispatch(t))); + } + return dispatch(taskOrTasks); + }, + { + get signal() { + return ac.signal; + }, + get workers() { + return workerHandles; + }, + [Symbol.dispose]() { + disposed = true; + ac.abort(); + disposeBalancer(); + terminateAll(); + }, + async [Symbol.asyncDispose]() { + disposed = true; + ac.abort(); + const settle = Promise.allSettled(inflight); + if (opts?.shutdownTimeout != null) { + const timeoutAc = new AbortController(); + await Promise.race([ + settle.finally(() => timeoutAc.abort()), + setTimeout(opts.shutdownTimeout, undefined, { signal: timeoutAc.signal }).catch(() => {}), + ]); + } else { + await settle; + } + await asyncDisposeBalancer(); + terminateAll(); + }, + }, + ); + + return run; +} +``` + +- [ ] **Step 5: Run tests** + +Run: `node --no-warnings --test test/load-balancing.test.ts 2>&1 | tail -15` +Expected: All tests pass. + +Run: `pnpm test 2>&1 | tail -10` +Expected: All existing tests pass. + +Run: `pnpm tsc --noEmit 2>&1` +Expected: No errors. + +- [ ] **Step 6: Commit** + +```bash +git add src/worker-pool.ts test/load-balancing.test.ts test/fixtures/load-balancing.ts +git commit -m "feat: integrate balancer into worker pool with activeCount and thread" +``` + +--- + +### Task 4: Update README + +**Files:** +- Modify: `README.md` + +- [ ] **Step 1: Update workers() docs** + +Read `README.md`. Find the `### \`workers(size)\`` section. Update the heading and description, and add a load balancing subsection. + +Change the heading from: +```markdown +### `workers(size)` +``` +to: +```markdown +### `workers(size?, opts?)` +``` + +After the existing Graceful Shutdown subsection, add: + +```markdown +#### Load Balancing + +The pool uses round-robin scheduling by default. Pass a `balance` option to change the strategy: + +\`\`\`ts +import { workers, leastBusy } from 'moroutine'; + +{ + using run = workers(4, { balance: leastBusy() }); + // tasks dispatched to whichever worker has the fewest in-flight tasks +} +\`\`\` + +Built-in balancers: +- `roundRobin()` — cycles through workers in order (default) +- `leastBusy()` — picks the worker with the lowest active task count + +Custom balancers implement the `Balancer` interface: + +\`\`\`ts +import type { Balancer } from 'moroutine'; + +const myBalancer: Balancer = { + select(workers, task) { + return workers[0]; // always use first worker + }, +}; +\`\`\` + +Each worker handle exposes `thread` (the underlying `worker_threads.Worker`) and `activeCount` for building custom strategies. +``` + +- [ ] **Step 2: Run type check** + +Run: `pnpm tsc --noEmit 2>&1` +Expected: No errors. + +- [ ] **Step 3: Commit** + +```bash +git add README.md +git commit -m "docs: document load balancing in README" +```