Zum Inhalt springen
Deutsch

MultiNodeSpec

MultiNodeSpec lässt N Cluster-Nodes innerhalb eines Test-Prozesses laufen. Jeder ist ein echtes ActorSystem; sie kommunizieren via einen In-Process-Transport (kein TCP). Nützlich zum Testen von Cluster-Verhalten — Sharding, Singletons, Gossip- Konvergenz, Failover — ohne Docker.

import { MultiNodeSpec, TestProbe } from 'actor-ts/testkit';
// Jeder Rollenname ist ein Node — und zugleich dessen System-Name.
const spec = new MultiNodeSpec({ roles: ['a', 'b', 'c'] });
await spec.start();
await spec.awaitMembers('a', 3); // warten, bis der Cluster konvergiert
// Das ActorSystem jedes Nodes via systemFor(role) erreichen:
const node2 = spec.systemFor('b');
// Actor auf einem bestimmten Node spawnen und mit einer Probe asserten:
const probe = new TestProbe(spec.systemFor('a'));
const remote = node2.spawnAnonymous(Worker);
remote.tell({ kind: 'do', replyTo: probe });
await probe.expectMessage({ kind: 'done' });
await spec.stop();

Drei primäre Fälle:

  1. Cluster-Verhalten testen — Sharding-Rebalances, Singleton-Failover, Gossip-Konvergenz — diese ohne echtes Netzwerk-Setup verifizieren.
  2. Verteilte Bugs reproduzieren — leichter zu isolieren, wenn der gesamte Cluster in einem Prozess läuft.
  3. CI-freundliche Cluster-Tests — schnell (sub-Sekunde), kein Docker, keine Ports.

Für Tests, die echte Parallelität brauchen (echte OS-Threads — nebenläufige Journal-Schreibvorgänge, Scheduler-Interleaving), nutze ParallelMultiNodeSpec. Für echtes TCP/TLS oder Cross-Host-Latenz nutze einen externen Docker-Compose-Cluster.

type MultiNodeSpecOptionsType = {
roles: ReadonlyArray<string>; // ein Node pro Rolle; müssen eindeutig sein
seedRoles?: ReadonlyArray<string>; // Bootstrap-Seeds — Default [roles[0]]
addresses?: Record<string, { host: string; port: number }>;
failureDetector?: ClusterOptionsType['failureDetector'];
gossipIntervalMs?: number; // Default 100 (vs. 1 s in Produktion)
awaitTimeoutMs?: number; // Default 10_000 — wie lange await*-Helfer warten
logLevel?: LogLevel; // Default: ein leiser NoopLogger
downing?: (role: string) => DowningProvider | undefined;
};
FeldZweck
rolesNode-Liste — jeder Eintrag ist ein Node, und der String ist zugleich dessen System-Name. Müssen eindeutig sein.
seedRolesWelche Rollen als Bootstrap-Seeds dienen. Default ist die erste Rolle.
addressesPer-Rolle Host/Port-Overrides. Automatisch vergeben, wenn ausgelassen.
failureDetectorFailure-Detector-Overrides — Tests ziehen diese meist enger, damit Crashes schnell erkannt werden.
gossipIntervalMsGossip-Runden-Intervall. Default 100 ms.
awaitTimeoutMsDefault-Timeout für die await*-Helfer. Default 10 s.
logLevelLog-Level — standardmäßig leise.
downingPer-Rolle Split-Brain-Resolver-Factory.

Übergib ein einfaches Objekt wie gezeigt, oder baue es fluent mit MultiNodeSpecOptions.create().withRoles(['a', 'b', 'c']).

Jeder Eintrag in roles ist ein eigener Node — der String ist die Identität dieses Nodes und sein ActorSystem-Name, muss also eindeutig sein (der Konstruktor wirft bei Duplikaten oder leerer Liste). Standardmäßig ist die erste Rolle der einzige Bootstrap-Seed; überschreibe das mit seedRoles:

const spec = new MultiNodeSpec({
roles: ['seed', 'worker-1', 'worker-2'],
seedRoles: ['seed'],
});
await spec.start();

Ports werden automatisch auf 127.0.0.1 vergeben; fixiere bestimmte via addresses, wenn ein Test feste Endpunkte braucht:

const spec = new MultiNodeSpec({
roles: ['a', 'b'],
addresses: { a: { host: '127.0.0.1', port: 2551 } },
});

Alle Nodes teilen sich:

  • Einen gemeinsamen In-Process-Nachrichtenbus — sie können miteinander reden.
  • Dasselbe Gossip-Protokoll — Membership konvergiert.
  • Unabhängige Actor-Systeme — separate Dispatcher, Scheduler, Supervisor-Bäume.

Das heißt:

  • Echte Cluster-Semantik — Mitglieder kommen hoch, Gossip konvergiert, Failure-Detector beobachtet, Sharding rebalancet.
  • Keine Serialisierung — Nachrichten zwischen „Nodes” werden per Referenz übergeben (in-process). Kein Fit zum Testen serialisierungs-abhängigen Verhaltens.

Zwei Wege, einen Node herunterzunehmen: crash(role) reißt dessen Transport abrupt weg (Peers erkennen das über den Failure-Detector), und leave(role) führt einen sauberen Cluster-Austritt durch.

// Singleton-Failover verifizieren:
const spec = new MultiNodeSpec({
roles: ['a', 'b', 'c'],
failureDetector: { heartbeatIntervalMs: 50, unreachableAfterMs: 200, downAfterMs: 400 },
});
await spec.start();
await Promise.all([
spec.awaitMembers('a', 3),
spec.awaitMembers('b', 3),
spec.awaitMembers('c', 3),
]);
// ... den Singleton auf jedem Node starten ...
// Den aktuellen Host ermitteln, dann crashen:
const host = spec.clusterFor('a').leader().toNullable();
await spec.crash('a');
// Warten, bis die Überlebenden 'a' downen und neu konvergieren:
await spec.awaitMembers('b', 2);
await spec.awaitMemberStatus('b', 'a', 'down');
// ... asserten, dass der Singleton zu einem Überlebenden umgezogen ist ...
await spec.stop();

crash(role) entfernt den Node abrupt; die anderen beobachten den Unreachable-Status, gossipen die Änderung und lösen dann Downing + Failover aus. Die await*-Helfer ersetzen fragile setTimeout-Schlafphasen — sie pollen, bis die Bedingung gilt (oder werfen nach awaitTimeoutMs).

MultiNodeSpec hat keine Probe-Factory — konstruiere eine TestProbe gegen das System des Nodes, den du beobachten willst:

const probeA = new TestProbe(spec.systemFor('a')); // Probe auf Node 'a'
const probeB = new TestProbe(spec.systemFor('b')); // Probe auf Node 'b'

Jede Probe ist an das Actor-System eines Nodes gebunden. Nützlich beim Testen von Routing — verifizieren, dass eine Nachricht auf dem erwarteten Node landet.

Starte eine Sharding-Region auf jedem Node über dessen Cluster — erreiche jede via clusterFor(role):

const spec = new MultiNodeSpec({ roles: ['a', 'b', 'c'] });
await spec.start();
await spec.awaitMembers('a', 3);
const startShardingOptions = StartShardingOptions.create<Command>()
.withTypeName('entity')
.withEntityActor(Entity)
.withExtractEntityId((message) => message.id);
const regions = spec.allRoles().map((role) =>
spec.clusterFor(role).sharding.start(startShardingOptions)
);
// Entities spawnen; verifizieren, dass sie sich über Nodes verteilen:
for (const id of ['e1', 'e2', 'e3', 'e4', 'e5']) {
regions[0].tell({ id, kind: 'wake-up' });
}
await spec.stop();

Der Koordinator des Frameworks verteilt Shards über die Nodes wie in einem echten Cluster.