Worker mesh
이 콘텐츠는 아직 번역되지 않았습니다.
JavaScript is single-threaded per actor system. For
parallelism within one OS process, the framework’s
worker mesh runs multiple ActorSystems — one per worker
thread — all participating in the same cluster via a
MessageChannel transport.
Each is a separate cluster node to the cluster’s view — gossip + membership + sharding all apply. Communication between them goes via in-process MessageChannel (no serialization to bytes, no TCP).
When to use it
Section titled “When to use it”Two main scenarios:
- CPU-bound parallelism in one process — actor-ts is single-threaded per system; multi-threading needs multiple systems. Worker mesh distributes them.
- Isolation within one process — a worker failing doesn’t
take down the main system. That covers an uncaught throw, an
unhandled rejection, and a bootstrap that fails to load: the
framework subscribes the worker’s
errorevent, so the failure reachesrestartPolicyinstead of the host’s own crash path. See Failure containment for what is and isn’t guaranteed per runtime.
For multi-process parallelism (separate OS processes), use regular cluster + TCP transport. Worker mesh is specifically for the in-process case.
High-level API: WorkerCluster
Section titled “High-level API: WorkerCluster”Wiring the broker, channels, and per-worker handshake by hand (see
Under the hood below) is exactly what
WorkerCluster automates. It spawns a pool of workers from one
entrypoint module, runs the hello/init/ready handshake, registers each
worker with a shared WorkerBroker, and restarts crashed workers per a
restartPolicy. The underlying worker primitive is picked per runtime
— Web Workers on Bun/Deno, node:worker_threads on Node — so the same
code runs everywhere.
// main.ts — main threadimport { WorkerCluster, WorkerClusterOptions } from 'actor-ts/worker';
const workerClusterOptions = WorkerClusterOptions.create() .withWorkers(4) .withBootstrap(new URL('./worker-node.js', import.meta.url)) .withSystemName('multi-core') .withBasePort(2552);const cluster = await WorkerCluster.spawn(workerClusterOptions);
console.log(`Spawned ${cluster.size} workers:`);for (const address of cluster.addresses) console.log(' -', address.toString());
// ...later, on shutdown:await cluster.terminate();Each worker runs the bootstrap module. It calls WorkerNode.join()
to complete the handshake and receive its address, system name,
transport, and initData, then joins the cluster like any other node:
// worker-node.ts — runs inside each workerimport { ActorSystem } from 'actor-ts';import { Cluster, ClusterOptions } from 'actor-ts/cluster';import { WorkerNode } from 'actor-ts/worker';
async function main(): Promise<void> { const context = await WorkerNode.join<{ seedAddress?: string }>(); const system = ActorSystem.create(context.systemName); const clusterOptions = ClusterOptions.create() .withHost(context.self.host) .withPort(context.self.port) .withSeeds(context.initData.seedAddress ? [context.initData.seedAddress] : []) .withTransport(context.transport); await Cluster.join(system, clusterOptions);
system.spawn(MyActor, 'worker'); context.ready(); // tell the main thread this node is up}
void main();| Option | What |
|---|---|
withBootstrap(url) | Worker entrypoint module. Required. |
withWorkers(n | 'auto') | Pool size; 'auto' uses hardware concurrency. Default 'auto'. |
withSystemName(name) | ActorSystem name each worker hosts. Default 'worker-cluster'. |
withHostname(host) | Hostname component of each worker’s address. Default 'worker'. |
withBasePort(port) | Port of the first worker; each subsequent worker increments. Default 1. |
withInitData(data) | Payload delivered to every worker’s join() context (initData). Default null. |
withRestartPolicy(p) | 'always' / 'on-failure' / 'never'. Default 'on-failure'. |
withReadyTimeoutMs(ms) | Per-worker handshake timeout. Default 10000. |
withRestartMinBackoffMs(ms) | Delay before the first respawn of a crashed slot. Default 200. |
withRestartMaxBackoffMs(ms) | Ceiling for the respawn delay, which doubles per attempt. Default 10000. |
withRestartRandomFactor(f) | ± jitter fraction on each respawn delay, in [0, 1]. Default 0.2. |
withMaxRestarts(n) | Restarts granted per slot within the window before it is retired; -1 for unlimited. Default 10. |
withRestartWindowMs(ms) | Sliding window the restart budget counts over; 0 never resets it. Default 60000. |
withOnWorkerPermanentlyDown(fn) | Called once per slot whose restart budget is spent. Default: a console.error line. |
withBackend(backend) | Spawn through this WorkerBackend instead of the detected one — a runtime auto-detection does not know, or an in-memory fake in a test. Default: detected. |
The returned WorkerCluster exposes size, addresses
(the NodeAddress of each worker), the shared broker, and
terminate(). Misconfigured options (a bad basePort, a
non-positive workers) throw OptionsError at spawn.
spawn() leaves no live workers behind when it fails: a worker whose
handshake times out is terminated before the rejection propagates, and
so is every sibling that did come up — the instance that owns them is
never returned, so nothing else could reach them. terminate() waits
for the threads to actually go rather than firing and forgetting, within
the bound described below.
Under the hood: manual wiring
Section titled “Under the hood: manual wiring”The WorkerCluster above does this wiring for you. Reach for the
manual path only when you need full control — a custom topology, a
bespoke handshake, or extra per-node setup.
// main.ts — main threadimport { Worker } from 'node:worker_threads';import { ActorSystem } from 'actor-ts';import { Cluster, ClusterOptions, MessageChannelTransport, NodeAddress } from 'actor-ts/cluster';
const channel = new MessageChannel();
const w1 = new Worker('./worker.js', { workerData: { mainPort: channel.port2 }, transferList: [channel.port2],});
const mainAddress = new NodeAddress('main', 'main', 0);const transport = new MessageChannelTransport(mainAddress, channel.port1);
const system = ActorSystem.create('main');const clusterOptions = ClusterOptions.create() .withHost('main') .withPort(0) .withSeeds(['main']) .withTransport(transport);await Cluster.join( system, clusterOptions,);
// worker.js — runs in the worker threadimport { parentPort, workerData } from 'node:worker_threads';import { ActorSystem } from 'actor-ts';import { Cluster, ClusterOptions, MessageChannelTransport, NodeAddress } from 'actor-ts/cluster';
const workerAddress = new NodeAddress('w1', 'w1', 0);const transport = new MessageChannelTransport(workerAddress, workerData.mainPort);
const system = ActorSystem.create('w1');const cluster2Options = ClusterOptions.create() .withHost('w1') .withPort(0) .withSeeds(['main']) .withTransport(transport);await Cluster.join( system, cluster2Options,);
// From here on, w1 is just another cluster nodeThe mesh shape
Section titled “The mesh shape”The transport is a star, not a full mesh. A single
WorkerBroker on the main thread is the hub: it holds one
MessagePort per node and forwards each frame to the node named
by the envelope’s to address. Every node’s
MessageChannelTransport holds a single port — its own end
of one channel to the broker, never an array:
import { MessageChannelTransport, NodeAddress } from 'actor-ts/cluster';import { WorkerBroker } from 'actor-ts/worker';
// Main thread: one broker is the hub for every node.const broker = new WorkerBroker();
// One MessageChannel per worker — the broker keeps one end and// the worker receives the other (transferred via workerData).const workerAddress = new NodeAddress('w1', 'w1', 0);const channel = new MessageChannel();broker.register(workerAddress, channel.port1);
// Inside that worker, its transport holds the single far end:const transport = new MessageChannelTransport(workerAddress, channel.port2);So N workers need N channels (one per worker) — not binomial(N,2). The broker relays every hop, which keeps the wiring linear but makes the main thread the relay for all cross-node traffic.
How it differs from TCP cluster
Section titled “How it differs from TCP cluster”TCP transport: MessageChannelTransport:- Sockets, framing - postMessage between threads- Serialized bytes - Structured cloning (no JSON)- Network latency - Sub-microsecond- Cross-host - Same process onlyMessages between worker systems go through structured clone — faster than framing and re-parsing text. It carries the same rich types (Map, Set, Date, …) the TCP wire’s tagged JSON tree carries, without the text round-trip.
Use cases
Section titled “Use cases”Sharding across cores
Section titled “Sharding across cores”// 4-worker mesh; sharding distributes entities across them:const startShardingOptions = StartShardingOptions.create() .withTypeName('order') .withEntityActor(...) .withExtractEntityId((message) => message.id) .withNumShards(16);sharding.start( startShardingOptions,);The coordinator (on main) allocates shards to the 4 workers. CPU-bound entity work parallelizes across cores.
Per-worker isolation
Section titled “Per-worker isolation”// Worker that handles GPU-bound jobs:system.spawn(GpuJobActor, 'gpu-jobs');
// Crashes within this worker stay isolated from main + other workersA worker crashing doesn’t take down the main system — separate event loops.
Failure containment
Section titled “Failure containment”What “a worker failing doesn’t take down the main system” actually covers, and what it costs:
- An uncaught throw, an unhandled rejection, or a bootstrap that
won’t load is contained. The framework subscribes the worker’s
errorevent on every runtime; without that subscription Node re-raises the worker’s error on the host and Deno rejects an internal promise, and both exit1beforerestartPolicyis ever consulted. - One failure produces one respawn. Node and Bun emit
errorand then the exit event for a single throw, and Deno emits onlyerror— whichever arrives first drives the restart and the other is dropped. - Respawns back off and are budgeted. Each slot waits
restartMinBackoffMs, doubling torestartMaxBackoffMswithrestartRandomFactorjitter, and getsmaxRestartsattempts insiderestartWindowMs. A slot that spends its budget is retired: no further worker is started for it, andonWorkerPermanentlyDownis called once with the slot’s index, address and restart count. A replacement that never becomes ready counts against the same budget, so a permanently broken bootstrap stops instead of looping. - A failed respawn degrades the mesh by one worker rather than
killing the host.
sizedrops, the slot’s address is unregistered from the broker, and the failure is reported.
Two runtime asymmetries are worth knowing:
Malformed traffic is dropped, not fatal: the broker validates a frame’s
to and from addresses before routing it, so one bad
postMessage from a worker no longer throws inside the host’s message
listener. The frame’s payload is validated by the receiving
transport, not by the broker.
Sender identity
Section titled “Sender identity”The broker re-addresses every frame to the port it arrived on. A
frame’s from is the address that port was registered under — not the
address the sending worker wrote into it:
// Worker 1 posts a frame claiming to be worker 2 …port1.postMessage({ from: worker2.toJSON(), to: worker3.toJSON(), payload });
// … and worker 3's transport still reports worker 1 as the sender.The host assigns those addresses — WorkerCluster mints one per slot,
and manual wiring hands one to broker.register — so a receiving
node’s peer identity is derived from the channel rather than from the
payload. That identity is what its failure detector credits, what an
envelope’s sender is taken to be, and what every “may this peer say
that?” rule keys on, so without the correction one worker could keep a
dead sibling looking alive and have its own messages attributed to
that sibling.
Only the system@host:port slot is corrected. The optional
incarnation an address may carry is passed through as the sender
wrote it, because nothing in the cluster keys on it yet.
When NOT to use it
Section titled “When NOT to use it”Where to next
Section titled “Where to next”- Cluster overview — the cluster model worker-mesh participates in.
- Transports — the transport interface MessageChannelTransport implements.
- Sharding — the primary consumer of mesh parallelism.
