diff --git a/src/balancers.ts b/src/balancers.ts index 6c02493..bb31c30 100644 --- a/src/balancers.ts +++ b/src/balancers.ts @@ -1,23 +1,27 @@ +import { uint32atomic } from './shared/descriptors.ts'; +import type { Uint32Atomic } from './shared/uint32-atomic.ts'; import type { Balancer, WorkerHandle } from './runner.ts'; /** * Creates a round-robin balancer that cycles through workers in order. + * The cursor lives in shared memory so all scheduling threads advance the + * same sequence. (The uint32 cursor wraps at 2^32; a wrap may skip one + * position, which round-robin semantics tolerate.) * @returns A fresh Balancer instance. */ -export function roundRobin(): Balancer { - let next = 0; +export function roundRobin(): Balancer<{ cursor: Uint32Atomic }> { return { - select(workers: readonly WorkerHandle[]): WorkerHandle { - const worker = workers[next % workers.length]; - next++; - return worker; + initialState: () => ({ cursor: uint32atomic() }), + select(workers: readonly WorkerHandle[], _task, { cursor }): WorkerHandle { + return workers[cursor.add(1) % workers.length]; }, }; } /** * Creates a least-busy balancer that picks the worker with the lowest activeCount. - * Ties are broken by index (first wins). + * Ties are broken by index (first wins). Reads the pool's shared active counts + * via the handles, so it needs no state of its own. * @returns A fresh Balancer instance. */ export function leastBusy(): Balancer { diff --git a/src/runner.ts b/src/runner.ts index fb18748..e6c3b45 100644 --- a/src/runner.ts +++ b/src/runner.ts @@ -22,9 +22,14 @@ type TaskResults[]> = { }; /** A load balancing strategy for choosing which worker runs a task. */ -export interface Balancer { +export interface Balancer { + /** Produces the balancer's state once at pool/runtime boot. May contain + * shared memory (which crosses threads by reference); plain values ride + * along by serialization. Closure state is safe only where scheduling is + * single-threaded (a `workers()` pool); shared, declared state works everywhere. */ + initialState?(): S; /** Choose a worker for the given task. Called synchronously on every dispatch. */ - select(workers: readonly WorkerHandle[], task: Task): WorkerHandle; + select(workers: readonly WorkerHandle[], task: Task, state: S): WorkerHandle; /** Optional cleanup on sync dispose. */ [Symbol.dispose]?(): void; /** Optional cleanup on async dispose. */ @@ -35,8 +40,8 @@ export interface Balancer { 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; + /** Load balancing strategy. Defaults to least-busy. */ + balance?: Balancer; } /** A handle to a specific worker in a pool. Uniform across threads. */ diff --git a/src/worker-pool.ts b/src/worker-pool.ts index c9954e0..b0a08ba 100644 --- a/src/worker-pool.ts +++ b/src/worker-pool.ts @@ -3,7 +3,7 @@ import { Worker } from 'node:worker_threads'; import { availableParallelism } from 'node:os'; import { setupWorker, execute, dispatchStream } from './execute.ts'; import { AsyncIterableTask } from './stream-task.ts'; -import { roundRobin } from './balancers.ts'; +import { leastBusy } from './balancers.ts'; import { ActiveCounts } from './active-counts.ts'; import type { ChannelOptions } from './channel.ts'; import type { Task, Balancer, Runner, WorkerHandle, WorkerOptions } from './runner.ts'; @@ -31,7 +31,8 @@ export function workers(sizeOrOpts?: number | WorkerOptions, opts?: WorkerOption size = sizeOrOpts ?? availableParallelism(); } - const balancer: Balancer = opts?.balance ?? roundRobin(); + const balancer: Balancer = opts?.balance ?? leastBusy(); + const balancerState: unknown = balancer.initialState?.(); const pool: Worker[] = []; for (let i = 0; i < size; i++) { @@ -99,7 +100,7 @@ export function workers(sizeOrOpts?: number | WorkerOptions, opts?: WorkerOption const idx = workerHandles.indexOf(task.worker); if (idx !== -1) return { worker: pool[idx], idx }; } - const handle = balancer.select(workerHandles, task); + const handle = balancer.select(workerHandles, task, balancerState); const idx = workerHandles.indexOf(handle); if (idx === -1) throw new Error('Balancer returned a handle that is not in this pool'); return { worker: pool[idx], idx }; diff --git a/test/balancers.test.ts b/test/balancers.test.ts index 103ce8d..beeb5a8 100644 --- a/test/balancers.test.ts +++ b/test/balancers.test.ts @@ -10,12 +10,23 @@ function mockHandle(activeCount: number, index = 0): WorkerHandle { describe('roundRobin()', () => { it('cycles through workers in order', () => { const b = roundRobin(); - const handles = [mockHandle(0), mockHandle(0), mockHandle(0)]; + const state = b.initialState!(); + const handles = [mockHandle(0, 0), mockHandle(0, 1), mockHandle(0, 2)]; 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]); + assert.equal(b.select(handles, task, state), handles[0]); + assert.equal(b.select(handles, task, state), handles[1]); + assert.equal(b.select(handles, task, state), handles[2]); + assert.equal(b.select(handles, task, state), handles[0]); + }); + + it('state is shared-memory backed (two instances over the same state agree)', () => { + const b1 = roundRobin(); + const state = b1.initialState!(); + const b2 = roundRobin(); + const handles = [mockHandle(0, 0), mockHandle(0, 1)]; + const task = { id: 'test', args: [], uid: 0 } as any; + assert.equal(b1.select(handles, task, state), handles[0]); + assert.equal(b2.select(handles, task, state), handles[1]); }); }); @@ -24,13 +35,13 @@ describe('leastBusy()', () => { 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]); + assert.equal(b.select(handles, task, undefined), 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]); + assert.equal(b.select(handles, task, undefined), handles[0]); }); }); diff --git a/test/load-balancing.test.ts b/test/load-balancing.test.ts index e87af5a..e4dec09 100644 --- a/test/load-balancing.test.ts +++ b/test/load-balancing.test.ts @@ -5,7 +5,7 @@ import type { Balancer } from 'moroutine'; import { identity, slow } from './fixtures/load-balancing.ts'; describe('load balancing', () => { - it('defaults to round-robin', async () => { + it('defaults to least-busy', async () => { const run = workers(2); try { const result = await run(identity(42));