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.

service/queryIncidents.ts
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

incidents.worker.ts
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 imports

Main thread — connectWorker, type-only import

incidents.ts
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

OptionMeaning
workerFactory () => new Worker(new URL(...)) — required for esbuild/webpack/turbopack to detect the entry. A URL also works where the bundler emits one (Vite).
sharedMemoryOptional. 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.
poolSizeA number, or 'auto' (default) for navigator.hardwareConcurrency ?? 4.
concurrencyMax in-flight calls per worker (default 1). Dispatch is least-busy; when every worker is at the cap, calls queue FIFO.
maxQueueQueue bound (default unbounded). A full queue rejects immediately with PoolQueueFullError — backpressure instead of unbounded growth.
taskTimeoutDefault per-call timeout in ms, measured from enqueue (queue wait + run). Rejects with TaskTimeoutError.
respawnDefault true — a crashed worker is replaced and its in-flight calls reject with WorkerCrashedError. false shrinks the pool instead.
memoryBuffer growth: maximumPages (default 16384 = 1 GB), growthFactor.
lazyDefault 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 terminate

A 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):

ExportWhat 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.tasksA 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.