diff --git a/docs/superpowers/specs/2026-04-14-load-balancing-design.md b/docs/superpowers/specs/2026-04-14-load-balancing-design.md new file mode 100644 index 0000000..65955ca --- /dev/null +++ b/docs/superpowers/specs/2026-04-14-load-balancing-design.md @@ -0,0 +1,91 @@ +# Extensible Load Balancing + +## Goal + +Support pluggable load balancing strategies for worker pools, with built-in round-robin and least-busy implementations. Custom balancers can track their own metrics using the exposed worker thread and active count. + +## `Balancer` Interface + +```ts +interface Balancer { + select(workers: readonly WorkerHandle[], task: Task | StreamTask): WorkerHandle; + [Symbol.dispose]?(): void; + [Symbol.asyncDispose]?(): Promise; +} +``` + +`select()` is called synchronously on every dispatch to choose which worker runs the task. It receives the full handles array and the task being dispatched. Pinned tasks (via `assign()`) bypass the balancer. + +`Symbol.dispose` and `Symbol.asyncDispose` are optional cleanup hooks. The pool calls the appropriate one during disposal — `Symbol.dispose` on sync dispose, `Symbol.asyncDispose` on async dispose. If only one is defined, the pool falls back to whichever is available. + +## `WorkerHandle` Additions + +```ts +interface WorkerHandle { + exec(task: Task): Promise; + exec(task: StreamTask, opts?: ChannelOptions): AsyncIterable; + readonly thread: Worker; // underlying worker_threads.Worker + readonly activeCount: number; // in-flight tasks on this worker +} +``` + +`thread` exposes the raw `worker_threads.Worker` for advanced use cases (event loop utilization via `worker.performance`, custom messaging, etc.). + +`activeCount` is the number of currently in-flight tasks on this worker. Incremented on dispatch, decremented on task settle or stream completion. This is per-handle tracking, separate from the pool-level `inflight` set. + +## `WorkerOptions` Update + +```ts +interface WorkerOptions { + shutdownTimeout?: number; + balance?: Balancer; +} +``` + +## `workers()` Overloads + +```ts +workers(): Runner +workers(size: number): Runner +workers(opts: WorkerOptions): Runner +workers(size: number, opts: WorkerOptions): Runner +``` + +When the first argument is an object (not a number), it's treated as `WorkerOptions` with default size. Default balancer is `roundRobin`. + +## Built-in Balancers + +Exported as factory functions from `moroutine`. Each call returns a fresh `Balancer` instance, safe for use with a single pool. + +```ts +import { workers, roundRobin, leastBusy } from 'moroutine'; + +workers(4) // round-robin by default +workers(4, { balance: leastBusy() }) // least-busy +workers({ balance: leastBusy() }) // least-busy, default size +workers(4, { balance: myCustomBalancer }) // custom +``` + +**`roundRobin()`** — cycles through workers in order. Maintains an internal index. + +**`leastBusy()`** — picks the worker with the lowest `activeCount`. Ties broken by index (first wins). Stateless — reads `activeCount` on each call. + +## Dispatch Flow + +1. If `task.worker` is set (pinned via `assign()`), dispatch to that worker. Balancer is not called. +2. Otherwise, call `balancer.select(handles, task)` to choose a worker. +3. Dispatch to the chosen worker. + +## Dispose Flow + +On sync dispose: call `balancer[Symbol.dispose]?.()`, falling back to `balancer[Symbol.asyncDispose]?.()` (fire-and-forget if async). + +On async dispose: call `balancer[Symbol.asyncDispose]?.()`, falling back to `balancer[Symbol.dispose]?.()`. If async, await it. + +## Files Changed + +- `src/runner.ts` — add `Balancer` interface, update `WorkerHandle` with `thread` and `activeCount`, update `WorkerOptions` +- `src/balancers.ts` — new file, `roundRobin` and `leastBusy` implementations +- `src/worker-pool.ts` — overloaded `workers()` signature, integrate balancer into dispatch, activeCount tracking per handle, balancer dispose +- `src/index.ts` — export `Balancer`, `roundRobin`, `leastBusy` +- Tests