コンテンツにスキップ
日本語

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.

Main process (single OS process)

ActorSystem 'main'

main thread

ActorSystem 'w1'

Worker thread 1

ActorSystem 'w2'

Worker thread 2

ActorSystem 'w3'

Worker thread 3

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).

Two main scenarios:

  1. CPU-bound parallelism in one process — actor-ts is single-threaded per system; multi-threading needs multiple systems. Worker mesh distributes them.
  2. 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 error event, so the failure reaches restartPolicy instead 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.

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 thread
import { 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 worker
import { 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();
OptionWhat
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.

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 thread
import { 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 thread
import { 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 node

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.

TCP transport: MessageChannelTransport:
- Sockets, framing - postMessage between threads
- Serialized bytes - Structured cloning (no JSON)
- Network latency - Sub-microsecond
- Cross-host - Same process only

Messages 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.

// 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.

// Worker that handles GPU-bound jobs:
system.spawn(GpuJobActor, 'gpu-jobs');
// Crashes within this worker stay isolated from main + other workers

A worker crashing doesn’t take down the main system — separate event loops.

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 error event on every runtime; without that subscription Node re-raises the worker’s error on the host and Deno rejects an internal promise, and both exit 1 before restartPolicy is ever consulted.
  • One failure produces one respawn. Node and Bun emit error and then the exit event for a single throw, and Deno emits only error — whichever arrives first drives the restart and the other is dropped.
  • Respawns back off and are budgeted. Each slot waits restartMinBackoffMs, doubling to restartMaxBackoffMs with restartRandomFactor jitter, and gets maxRestarts attempts inside restartWindowMs. A slot that spends its budget is retired: no further worker is started for it, and onWorkerPermanentlyDown is 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. size drops, 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.

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.

  • Cluster overview — the cluster model worker-mesh participates in.
  • Transports — the transport interface MessageChannelTransport implements.
  • Sharding — the primary consumer of mesh parallelism.