Zum Inhalt springen
Deutsch

BrokerActor-Basisklasse

Jeder Protokoll-Actor in actor-ts/io/broker — KafkaActor, MqttActor, NatsActor usw. — erbt von BrokerActor. Die Basisklasse besitzt den gemeinsamen Lifecycle: Verbindungszustandsautomat, Reconnect-mit-Backoff, Outbound-Buffer, Subscriber-Fan-out, Publishen von Lifecycle-Events.

BrokerActor (die abstrakte Basis) besitzt:

  • Lifecycle-Zustandsautomat — disconnected ↔ connecting ↔ connected ↔ disconnecting.
  • Outbound-Buffer — Nachrichten, die vor stehender Verbindung gesendet werden.
  • Reconnect-Schleife — exponentielles Backoff bei Verbindungsverlust.
  • Subscriber-Tracking — Fan-out für eingehende Events.
  • Gewünschte Subscriptions — Protokoll-Subscriptions, die jede einzelne Verbindung überleben und bei jedem Reconnect neu aufgebaut werden.

Subklassen implementieren drei Protokoll-Hooks:

HookWann aufgerufen
connectImplementationDie protokollspezifische Verbindung öffnen.
disconnectImplementationSie schließen — beim Stoppen und vor jedem Re-Connect-Versuch.
dispatchOutgoing(envelope)Eine einzelne gepufferte Nachricht auf die Leitung senden.

Subklassen implementieren die drei Protokoll-Hooks; den Rest macht die Basisklasse. Diese Seite dokumentiert das, was geteilt wird. Für protokollspezifische Details siehe die Seiten pro Protokoll.

connect

ok

Fehler — Reconnect-Zyklus

stop

disconnected

connecting

connected

disconnecting

Vier Zustände:

  • disconnected — initial; aktuell nicht verbunden.
  • connecting — connectImplementation läuft.
  • connected — Verbindung steht; Nachrichten fließen.
  • disconnecting — disconnectImplementation läuft.

Ein disconnected → connecting → Fehler löst eine Reconnect-Schleife aus: Backoff + Retry bis Erfolg oder bis maxAttempts erschöpft sind.

Jeder Re-Connect-Versuch startet auf sauberem Grund: die Basisklasse ruft zuerst disconnectImplementation auf, wenn ein früherer Versuch irgendetwas geöffnet hat — eine Subklasse baut also nie eine neue Verbindung auf den Handles der toten auf. Damit ist disconnectImplementation ein idempotenter Kontrakt: er darf auf einer bereits toten Verbindung aufgerufen werden und muss ausschließlich die lebenden Handles verwerfen, niemals die gewünschten Subscriptions (siehe unten).

Wird der Aktor gestoppt, während ein Reconnect-Versuch läuft, wird diese Verbindung verworfen statt übernommen. Reconnect läuft auf dem System-Scheduler, entkoppelt von der Mailbox — ein Versuch, der schon in connectImplementation steckt, lässt sich nicht mehr abbrechen und läuft weiter, nachdem postStop zurückgekehrt ist. Die Basisklasse prüft daher beim Verlassen des Handshakes erneut, ob der Aktor noch lebt, und ruft andernfalls ein weiteres Mal disconnectImplementation auf und bleibt disconnected: kein BrokerConnected, kein Leeren des Puffers, kein weiterer Retry. Ohne diese Prüfung installierte ein erfolgreicher Handshake lebende Handles auf einem beendeten Aktor, und ein gescheiterter stellte den Backoff-Timer neu und verband sich endlos weiter, denn maxAttempts ist standardmäßig Infinity (#708).

abstract class BrokerActor<S, Command, P, Subscription = never> extends Actor<Command> {
// Subklassen implementieren:
protected abstract configKey(): string;
protected abstract builtInDefaultOptions(): Partial<S>;
protected abstract readOptionsFromConfig(config: Config): Partial<S>;
protected abstract requiredOptions(): ReadonlyArray<keyof S>;
protected abstract endpointLabel(): string;
protected abstract connectImplementation(): Promise<void>;
protected abstract disconnectImplementation(): Promise<void>;
protected abstract dispatchOutgoing(envelope: OutboundEnvelope<P>): Promise<void>;
protected abstract onCommand(command: Command): void | Promise<void>;
// Optional — nur für Refs, die die Subklasse selbst beobachtet:
protected onTerminated(signal: Terminated): void | Promise<void> {}
}

Vier Kategorien:

  • Settings-Glue (configKey, builtInDefaultOptions, readOptionsFromConfig, requiredOptions, endpointLabel) — beschreibt, wie sich die Settings aus den drei Ebenen zusammensetzen (Konstruktor + HOCON + Defaults) und wie validiert wird.
  • Protokoll-Hooks (connectImplementation, disconnectImplementation, dispatchOutgoing) — die protokollspezifische Arbeit.
  • Nachrichtenverarbeitung (onCommand) — ein Command aus der Mailbox. Hier lebt die match(...).exhaustive()-Dispatch-Tabelle einer Subklasse.
  • Death-Watch (onTerminated) — optional; siehe Subscriber-Tracking.

endpointLabel ist die menschenlesbare Verbindungs-Identität (“amqp://localhost:5672”, “kafka-cluster-1”), die in Log-Zeilen und Lifecycle-Events auftaucht.

Gib den Verbindungs-String so zurück, wie er konfiguriert ist — redigiere ihn nicht selbst. Die Basisklasse schickt jede Verwendung durch redactedEndpointLabel(), was Userinfo und Query-String entfernt; amqp://svc:pw@rabbit:5671/orders erreicht BrokerConnected also als amqp://rabbit:5671/orders. Das wiegt schwerer, als es aussieht: Die vier Lifecycle-Events landen auf dem systemweiten EventStream, der keinen Autorisierungsbegriff kennt — jeder Actor im System kann sie also lesen — und eine Reconnect-Schleife publiziert sie bei jedem Backoff-Tick erneut, über die gesamte Dauer eines Ausfalls. Eine zusammengesetzte Server-Liste behält ihre Gestalt und verliert jede Credential (nats://***@a:4222,nats://***@b:4222), ein Label ohne Credential kommt unverändert zurück. Rufe redactedEndpointLabel() auch in deinen eigenen Log-Zeilen auf.

onReceive ist versiegelt — die Basisklasse besitzt die Methode, damit sie Terminated abfangen kann, bevor deine Dispatch-Tabelle es sieht. Implementiere stattdessen onCommand; dort kommt alles außer dem Death-Watch-Signal an.

this.enqueueOutbound(payload);

Subklassen rufen this.enqueueOutbound(payload) zum Senden auf; es gibt true zurück, wenn die Nachricht gesendet oder gepuffert wurde, false, wenn sie verworfen wurde. Die Basisklasse:

  • Wenn connected — verpackt payload in ein OutboundEnvelope und ruft sofort dispatchOutgoing(envelope) auf.
  • Wenn disconnected (oder connecting) — puffert bis zu outboundBuffer Envelopes (Default 1000). Beim Reconnect wird der Buffer der Reihe nach geleert.

Bei Überlauf verwirft die Basisklasse immer das älteste gepufferte Envelope (FIFO) und publisht ein BrokerBufferOverflow- Event auf dem Event-Stream — es wird nie geworfen. Mit outboundBuffer: 0 ist das Puffern deaktiviert: die Nachricht wird verworfen und ein BrokerNotConnected-Event publisht.

import { match } from 'ts-pattern';
protected override onCommand(command: MyCommand): void {
match(command)
.with({ kind: 'subscribe' }, (c) => this.onSubscribe(c))
.with({ kind: 'unsubscribe' }, (c) => this.onUnsubscribe(c))
.exhaustive();
}
private onSubscribe(command: SubscribeCommand): void {
this.subscribeRef(command.topic, command.subscriber);
}
private onUnsubscribe(command: UnsubscribeCommand): void {
this.unsubscribeRef(command.topic, command.subscriber);
}

subscribeRef(topic, ref) registriert ref als interessiert an den eingehenden Nachrichten von topic und death-watcht die Ref — für eine lokale Ref; siehe den Hinweis weiter unten.

unsubscribeRef vergleicht über den Pfad der Ref statt über das Ref-Objekt — eine andere Ref auf denselben Actor unsubscribed also trotzdem.

Wenn das Protokoll eine eingehende Nachricht pusht, ruft die Subklasse auf:

this.fanOutToTopic(topic, inboundMessage);

Die Basisklasse stellt sie an jeden Subscriber für dieses Topic zu.

Das Subscriber-Tracking oben spielt innerhalb des Actor-Systems — welche lokalen Refs welches Topic wollen. Eine gewünschte Subscription ist die andere Hälfte: die Subscription, die der Actor auf dem Broker hält und die bei jedem Verbindungsaufbau neu gesetzt werden muss.

Die Basisklasse hält diese Menge getrennt von den lebenden Handles, damit sie jede einzelne Verbindung überlebt:

// Vormerken (und sofort setzen, falls verbunden).
await this.rememberSubscription('orders.new', target);
// Entfernen — auf dem Broker und aus der gewünschten Menge.
await this.forgetSubscription('orders.new');

Subklassen liefern drei Hooks:

HookZweck
initialSubscriptions()In den Options deklarierte Subscriptions. Einmalig vor dem ersten Connect in die gewünschte Menge übernommen.
applySubscription(key, subscription)Eine Subscription auf der lebenden Verbindung setzen. Muss einen bereits lebenden Key tolerieren.
revokeSubscription(key)Eine abbauen. Optional — mehrere Protokolle können eine einzelne Subscription nur mitsamt dem ganzen Consumer aufgeben.

und spielen die ganze Menge aus connectImplementation heraus erneut ein:

protected async connectImplementation(): Promise<void> {
this.connection = await MyClient.connect(this.options.url);
await this.applyDesiredSubscriptions(); // konfiguriert + zur Laufzeit
}

Getrieben wird das Wiedereinspielen von der Subklasse und nicht von der Basisklasse, weil der richtige Punkt im Handshake protokollspezifisch ist — kafkajs etwa will jedes subscribe vor consumer.run haben.

Was das bringt:

  • Eine zur Laufzeit hinzugefügte Subscription überlebt einen Reconnect — sie ist nicht bloß ein Aufruf auf einer Verbindung, die gerade stirbt.
  • Ein subscribe, das während einer Trennung eintrifft, wird vorgemerkt und beim nächsten Connect gesetzt, statt verworfen zu werden.
  • Das Übernehmen aus den Options passiert nur einmal, ein unsubscribe zur Laufzeit wird also nicht vom nächsten Reconnect wiederbelebt.
  • Ein erneutes rememberSubscription auf einem lebenden Key baut ihn zuerst ab, damit eine geänderte Nutzlast (etwa ein anderer Ziel-Actor) tatsächlich greift.
  • Eine Subscription, die nicht gesetzt werden kann, wird als Warnung geloggt, und der Rest der Menge geht trotzdem durch. Ein schlechtes Subject reißt die Verbindung nicht mit — hinterlässt den Actor aber auch nicht stillschweigend verbunden-und-taub.

MqttActor ist älter als dieser Mechanismus und behält seine eigene, reichhaltigere Registry (QoS pro Topic, mehrere Targets pro Pattern, Deathwatch auf jedem) — mit den gleichen Reconnect-Garantien.

reconnect: {
initialDelayMs: 200,
maxDelayMs: 30_000,
factor: 2,
maxAttempts: Infinity, // der Default — endlos wiederholen
randomFactor: 0.2, // der Default — ±20 % Jitter
}

Pro Actor konfigurierbar (oder reconnect: false, um Auto-Reconnect ganz zu deaktivieren). Jeder Versuch wartet min(initialDelayMs * factor^(attempt - 1), maxDelayMs), multipliziert mit einem zufälligen 1 ± randomFactor.

Der Jitter ist das, was eine Flotte davon abhält, im Gleichschritt zu wiederholen. Ohne ihn ist die Verzögerung eine reine Funktion des Versuchszählers und dieser Optionen — jeder Broker-Actor, der denselben Broker im selben Moment verloren hat, wacht also in derselben Millisekunde auf, Welle für Welle, und die Herde kann den sich erholenden Broker unten halten. randomFactor: 0 stellt den vollständig deterministischen Zeitplan wieder her.

Dieselbe Streuung gilt für den Circuit Breaker: Ist der Breaker offen, prüft der Actor irgendwo in [resetMs, resetMs × (1 + randomFactor)] erneut statt exakt zum Stichtag. Dieser Jitter ist einseitig — ein Breaker wird nie zu früh erneut versucht.

Jeder Versuch löst BrokerReconnectAttempt auf dem Event-Stream aus (dessen delayMs ist der tatsächlich eingeplante, gejitterte Wert); nachdem maxAttempts erschöpft sind (sofern endlich), löst BrokerReconnectFailed aus und der Actor bleibt getrennt.

Für einen deterministischen Test lässt sich die Zufallsquelle injizieren — random: () => 0.5 fixiert jede Verzögerung auf ihren ungejitterten Wert:

reconnect: { initialDelayMs: 100, random: () => 0.5 }

Alles oben Beschriebene wird von handleConnectionLost ausgelöst, und das wird nur aufgerufen, wenn der Transport ein Ende meldet — ein close-Event, ein error, ein Stream, der done gemeldet hat. Ein Peer, der ohne FIN oder RST verschwindet, meldet nichts davon. Ein verfallener NAT-Eintrag, ein mit SIGKILL beendeter Container, eine Route, die anfängt zu verschlucken: der Socket bleibt offen, der Actor bleibt connected, und die Reconnect-Maschinerie läuft nie.

Zwei Uhren schließen diese Lücke. Beide sind pro Protokoll opt-in, weil nur die Anwendung weiß, wie still ihr Peer sein darf:

  • idleTimeoutMs — eine Lese-Deadline. Kommt so lange nichts vom Peer, wird die Verbindung über den gewöhnlichen handleConnectionLost-Pfad als verloren gemeldet; alles auf dieser Seite gilt damit unverändert: BrokerDisconnected, Backoff, Puffer, Breaker. Zurückgesetzt wird sie von eingehenden Bytes und ausdrücklich nicht von ausgehenden: ein Client, der in ein schwarzes Loch schreibt, ist genau der Fall, der sie auslösen muss.
  • connectTimeoutMs — eine Deadline für einen einzelnen connectImplementation-Aufruf. Jedes Protokoll beendet seinen Connect über ein Event (connect, open, Response-Header), und keines hat eine Uhr — ein Peer, der den Handshake abschließt und dann stehen bleibt, hält den Actor unbegrenzt in connecting. Läuft die Deadline ab, wird der Versuch abgebrochen, und der Fehler nimmt den Weg, den ein abgelehnter Connect ohnehin nimmt.

TcpSocketActor, SseActor und WebsocketClientActor bieten beide an. Setze idleTimeoutMs oberhalb des Heartbeats, den der Peer ohnehin sendet; darunter kappt die Deadline gesunde Verbindungen in Endlosschleife. 0 (der Default) schaltet beide ab.

Protokolle mit eigener Liveness auf Protokollebene brauchen das nicht: MqttActor hat MQTT-keepAlive, und KafkaActor fährt einen Consumer-Heartbeat.

Publisht auf system.eventStream:

EventWann
BrokerConnectedEin connectImplementation war erfolgreich.
BrokerDisconnectedEine bestehende Verbindung ging verloren. Ein geordneter Stopp publisht es nicht.
BrokerReconnectAttemptEin Reconnect-Versuch startet gerade.
BrokerReconnectFailedmaxAttempts erschöpft.
BrokerBufferOverflowDer Outbound-Buffer hat ein Envelope verworfen.
BrokerNotConnectedGesendet ohne Verbindung.

Subscribiere, um jeden Broker-Actor einheitlich zu beobachten:

system.eventStream.subscribe(monitorRef, BrokerConnected);
system.eventStream.subscribe(monitorRef, BrokerDisconnected);

Die Events enthalten actorPath — so unterscheidest du Events verschiedener Broker-Actors im System.

1. builtInDefaultOptions() ← niedrigste Priorität (immer angewendet)
2. readOptionsFromConfig() ← HOCON-Überschreibungen
3. Konstruktor-Argument ← höchste Priorität (pro Instanz)

preStart mergt die drei Ebenen, validiert gegen requiredOptions() und legt das Ergebnis für den Rest des Actor- Lebens unter this.options ab.

Fehlende Pflicht-Settings führen zu einem frühen Fehler-Throw in preStart — der Actor durchläuft den Fehlerpfad des Supervisors, bevor er überhaupt versucht, sich zu verbinden.

import { match } from 'ts-pattern';
import { BrokerActor, type OutboundEnvelope, type BrokerCommonOptionsType } from 'actor-ts/io';
interface MyProtocolOptionsType extends BrokerCommonOptionsType {
readonly url: string;
}
class MyProtocolActor extends BrokerActor<MyProtocolOptionsType, Command, MyPayload> {
private connection: MyClient | null = null;
protected configKey() { return 'actor-ts.io.broker.my-protocol'; }
protected builtInDefaultOptions() { return { /* ... */ }; }
protected readOptionsFromConfig(c) { /* HOCON parsen */ return {}; }
protected requiredOptions() { return ['url'] as const; }
protected endpointLabel() { return this.options.url; }
protected async connectImplementation(): Promise<void> {
this.connection = await MyClient.connect(this.options.url);
this.connection.onMessage((m) => this.fanOutToTopic(m.topic, m));
}
// Wird beim Stoppen UND vor jedem Re-Connect-Versuch aufgerufen —
// idempotent und auch auf einer bereits toten Verbindung sicher.
protected async disconnectImplementation(): Promise<void> {
await this.connection?.close();
this.connection = null;
}
protected async dispatchOutgoing(env: OutboundEnvelope<MyPayload>): Promise<void> {
await this.connection!.send(env.payload);
}
// `onReceive` ist von der Basisklasse versiegelt — implementiere `onCommand`.
protected override onCommand(command: Command): void {
match(command)
.with({ kind: 'subscribe' }, (c) => this.onSubscribe(c))
.exhaustive();
}
private onSubscribe(command: SubscribeCommand): void {
this.subscribeRef(command.topic, command.subscriber);
}
}

Den Rest übernimmt die Basisklasse. Die meisten Drittanbieter- Clients (kafkajs, nats.js usw.) haben eine eventbasierte Message-Receive-API, die sich sauber auf this.fanOutToTopic(...) abbilden lässt.

  • I/O-Übersicht — das große Bild: welche Protokolle ausgeliefert werden.
  • Kafka / MQTT / NATS / usw. — Seiten pro Protokoll.
  • Event-Stream — wo Lifecycle-Events publisht werden.
  • Backoff-Policy — ein eigenständiges Backoff-Primitiv (die Reconnect-Schleife rechnet selbst).