Worker pool & tasks
Tasks cross the thread boundary via postMessage; the data they touch stays in shared memory.
One method, one file
A serviceMethod unit holds a method's wire schemas and its worker implementation together. The schemas type run's parameters, and validation happens inside the worker — the trust boundary is the message.
import { serviceMethod } from '@atolljs/core';
import { z } from 'zod';
export const queryIncidents = serviceMethod({
def: { argsSchema: z.tuple([queryArgsSchema]), resultSchema: queryResultSchema },
run(q) { // q: QueryArgs — inferred from the schema
// scan/sort shared list rows in place, return only the visible page
return runQuery(q);
},
});Worker side — defineWorker owns the method list
import { defineWorker } from '@atolljs/core';
import { incidentsMemory } from './incidents.memory';
import { queryIncidents } from './service/queryIncidents';
// Wires INIT_MEMORY / EXECUTE_TASK and registers each method under its name.
export const incidentsWorker = defineWorker({
sharedMemory: incidentsMemory,
methods: { queryIncidents, ping: () => 'pong' }, // units or plain functions
// services: { pricing: { reprice } } → client.pricing.reprice(), taskId "pricing.reprice"
});
export type IncidentsWorker = typeof incidentsWorker; // all main ever importsMain thread — connectWorker, type-only import
import { connectWorker } from '@atolljs/core';
import { incidentsMemory } from './incidents.memory';
import type { IncidentsWorker } from './incidents.worker'; // zero worker code in this bundle
export const incidents = connectWorker<IncidentsWorker>({
sharedMemory: incidentsMemory,
worker: () => new Worker(new URL('./incidents.worker.ts', import.meta.url), { type: 'module' }),
poolSize: 'auto', // 'auto' = navigator.hardwareConcurrency ?? 4
});
// A Proxy typed by the worker's signatures — pool spawns on first call:
const page = await incidents.queryIncidents({ offset: 0, limit: 50 });
incidents.terminate(); // next call re-spawns
// Already hold a pool (e.g. Nest's @InjectAtollPool)? Wrap it directly:
import { workerClient } from '@atolljs/core';
const api = workerClient<IncidentsWorker>(pool);connectWorker config
| Option | Meaning |
|---|---|
worker | Factory () => new Worker(new URL(...)) — required for esbuild/webpack/turbopack to detect the entry. A URL also works where the bundler emits one (Vite). |
sharedMemory | Optional. A contract to bind on the main thread and ship to workers — type-checked against the worker's declaration. Omit for a message-only pool: no SharedArrayBuffer, no COOP/COEP headers required. |
poolSize | A number, or 'auto' (default) for navigator.hardwareConcurrency ?? 4. |
concurrency | Max in-flight calls per worker (default 1). Dispatch is least-busy; when every worker is at the cap, calls queue FIFO. |
maxQueue | Queue bound (default unbounded). A full queue rejects immediately with PoolQueueFullError — backpressure instead of unbounded growth. |
taskTimeout | Default per-call timeout in ms, measured from enqueue (queue wait + run). Rejects with TaskTimeoutError. |
respawn | Default true — a crashed worker is replaced and its in-flight calls reject with WorkerCrashedError. false shrinks the pool instead. |
memory | Buffer growth: maximumPages (default 16384 = 1 GB), growthFactor. |
lazy | Default true — spawn on first call (SSR-safe import). false spawns at construction. |
Cancellation, timeouts, backpressure
// Per-call controls ride on .with() — the same typed surface
const ctrl = new AbortController();
const page = incidents.with({ signal: ctrl.signal, timeout: 2_000 }).queryIncidents(q);
ctrl.abort(); // rejects with TaskAbortedError — queued → dropped;
// in-flight → rejects now, the worker's late reply is discarded
// Observe the pool
incidents.pool?.stats();
// { workers, idle, inFlight, queued, completed, failed, aborted,
// waitMs: { count, mean, max }, runMs: { count, mean, max } }
await incidents.pool?.close(); // drain the queue, then terminateA signal abort rejects with TaskAbortedError (exported from @atolljs/core, alongside PoolQueueFullError, TaskTimeoutError, and WorkerCrashedError). JavaScript can't interrupt a running function, so an in-flight abort rejects the caller and holds the worker's slot until its reply arrives — the next call goes to a genuinely free worker. For a hard stop, call terminate().
Client members start(), terminate(), with(), pool, sharedMemory are reserved — a worker method by those names is a compile error. Under the hood connectWorker builds a WorkerPool; the explicit-contract API (TaskContract + TaskRegistry.register + WorkerPool's tasks) remains available when both threads need the contract object at runtime.
Explicit service contracts (advanced)
defineWorker/connectWorker build on a lower, framework-neutral layer — the same one the NestJS binding dispatches through. Reach for it when both threads need the contract object at runtime (a hand-built WorkerPool, a DI provider, a test stub):
| Export | What it does |
|---|---|
defineService(name, methods) | Declares the contract bundle once — { method: { argsSchema?, resultSchema? } }. Wire ids derive as service.method; signatures infer from the schemas (argsSchema: z.tuple(...) → args, resultSchema → return type). |
implementService(service, handlers) | Worker-side registration; throws at bind time on a missing method. |
createClient(service, runner) | Typed proxy over any TaskRunner — a pool, a SharedWorker client, or a test stub. |
service.tasks | A TaskMap you can feed straight to new WorkerPool({ tasks }). |
rpc<A, R>({ taskId? }) | Escape hatch for schema-less methods or explicit wire ids. |
defineTask — latest-wins runners
defineTask(fn) wraps an async call into an AsyncTask: latest-wins (stale results dropped), plus an observable snapshot { data, pending, settled, elapsedMs, error } the framework bindings render. runOnce() makes init-style tasks remount/StrictMode-safe.
// Bindings accept any async function and wrap it per call site:
const page = useTask(incidents.queryIncidents); // React — taskState() in Angular/Svelte
page.run({ offset: 0, limit: 50 }); // latest-wins
page.data; page.pending; page.elapsedMs;
// The primitive underneath, if you need it outside a framework:
import { defineTask } from '@atolljs/core';
const queryTask = defineTask((q: QueryArgs) => incidents.queryIncidents(q));
queryTask.subscribe((snap) => render(snap));Node workers
WorkerPool/connectWorker run on node:worker_threads unchanged — @atolljs/node adapts Node's Worker to the DOM surface the pool expects. createNodePool({ workerFile | worker | createWorker, … }) is WorkerPool with the adapter baked in; its worker: factory may return a node:worker_threads.Worker directly (adapted internally), which keeps the bundler-detectable new Worker(new URL('./x.worker.ts', import.meta.url)) literal usable on Node. For connectWorker or new WorkerPool directly, wrap the spawn yourself with createNodeWorker.
A Node worker entry imports @atolljs/node/shim first — it binds self = parentPort before defineWorker's bootstrap evaluates. SharedArrayBuffer works in Node with no headers: isolation is a browser-only requirement.