Zum Inhalt springen
Deutsch

WebsocketClientActor

WebsocketClientActor öffnet eine clientseitige WebSocket-Verbindung zu einem Remote-Endpoint. Eingehende Frames werden in typisierte Nachrichten dekodiert, die du in onMessage behandelst; ausgehende Nachrichten werden kodiert und gesendet (während der Trennung gepuffert, nach dem Reconnect erneut gesendet). Zum Annehmen eingehender WS-Verbindungen (Serverseite) siehe Server-WebSocket.

import { match } from 'ts-pattern';
import { WebsocketClientActor, WebsocketClientOptions, websocketSend } from 'actor-ts/http';
type ClientMessage = { kind: 'ping'; n: number };
type ServerMessage = { kind: 'pong'; n: number } | { kind: 'event'; data: unknown };
class Feed extends WebsocketClientActor<ClientMessage, ServerMessage> {
constructor() {
const webSocketClientOptions = WebsocketClientOptions.create<ClientMessage, ServerMessage>().withUrl('ws://localhost:8080/ws');
super(webSocketClientOptions);
}
override onConnected(): void {
this.send({ kind: 'ping', n: 1 });
}
onMessage(message: ServerMessage): void {
// decoded server message
match(message)
.with({ kind: 'pong' }, (m) => this.onPong(m))
.exhaustive();
}
private onPong(message: PongMessage): void {
this.log.info(`pong ${message.n}`);
}
}
const feed = system.spawn(Feed, 'feed');

Die beiden Typparameter sind aus Sicht des Clients zu lesen: TOut = die Nachrichten, die dieser Client sendet, TIn = die dekodierten Server-Nachrichten, die er empfängt.

WebsocketClientActor erweitert die gemeinsame BrokerActor-Basis und erbt damit Reconnect-mit-Backoff, den Outbound-Buffer über Reconnects hinweg und den Circuit-Breaker kostenlos.

Der Konstruktor nimmt Partial<WebsocketClientOptionsType<TOut, TIn>>. Nur url ist erforderlich; alles andere hat einen Default.

interface WebsocketClientOptionsType<TOut, TIn> extends BrokerCommonOptionsType {
url: string; // required — ws:// or wss://
protocols?: string | string[]; // subprotocols
codec?: WebsocketCodec<TOut, TIn>; // default jsonCodec()
maxFrameBytes?: number; // default 1 MiB; oversize inbound closes with 1009
onInvalidMessage?: 'drop' | 'hook' | 'disconnect'; // default 'drop'
pingIntervalMs?: number; // application-level ping; default off
idleTimeoutMs?: number; // read-idle deadline; default 0 (off)
connectTimeoutMs?: number; // connect deadline; default 0 (off)
}

Da die Settings BrokerCommonOptionsType erweitern, gelten die Blöcke reconnect ({ initialDelayMs, maxDelayMs, factor, maxAttempts, randomFactor }), outboundBuffer und circuitBreaker allesamt — ihre Semantik steht in der BrokerActor-Basis.

Innerhalb des Actors kodiert this.send(message) die Nachricht über den Codec und stellt sie in die Queue; während der Trennung wird die Nachricht gepuffert und nach dem Reconnect erneut gesendet. Eingehende Frames werden dekodiert und an onMessage zugestellt.

class Chat extends WebsocketClientActor<ClientMessage, ServerMessage> {
constructor() {
const webSocketClientOptions = WebsocketClientOptions.create<ClientMessage, ServerMessage>().withUrl('wss://chat.example.com/ws');
super(webSocketClientOptions);
}
override onConnected(): void {
this.send({ kind: 'setName', name: 'alice' }); // encode + enqueue → true
}
onMessage(message: ServerMessage): void {
// one call per decoded server frame, in frame order
}
}

send(message: TOut): boolean gibt zurück, ob die Nachricht angenommen wurde (sie kann abgelehnt werden, wenn der Outbound-Buffer voll ist). Für einen rohen Frame — unter Umgehung des Codecs — nutze sendRaw(frame: WebsocketFrame): boolean.

Andere Actors rufen send nicht direkt auf; sie pushen einen typisierten Send über die Ref mit dem websocketSend-Helper:

import { websocketSend } from 'actor-ts/http';
feed.tell(websocketSend({ kind: 'ping', n: 42 }));

onMessage(message: TIn) ist abstrakt — du musst es implementieren. Der Rest sind optionale Overrides:

HookWann
onConnected()Die Verbindung ist (wieder) aufgebaut. Guter Ort für einen Handshake / Re-Subscribe.
onDisconnected(cause?: Error)Die Verbindung ist abgerissen; cause ist bei fehlergetriebenen Abbrüchen gesetzt.
onInvalidMessage(error: WebsocketDecodeError)Ein Frame ließ sich nicht dekodieren (nur bei onInvalidMessage: 'hook').
onSelfMessage(message: TSelf)Behandelt den optionalen dritten Typparameter — Nachrichten, die der Actor an sich selbst sendet.
onTerminated(signal: Terminated)Ein Aktor, den du mit context.watch beobachtest, hat gestoppt. Von BrokerActor geerbt.

onReceive und onCommand sind versiegelt — die Basisklassen dispatchen an die Hooks oben; überschreibe sie nicht. Ein Terminated aus einer eigenen Watch landet bei onTerminated, nicht bei onSelfMessage.

Der Codec wandelt TOut-Nachrichten in Wire-Frames und Wire-Frames zurück in TIn-Nachrichten. Der Default ist jsonCodec():

import { jsonCodec } from 'actor-ts/http';
new Feed(); // uses jsonCodec() — text frames <-> JSON
// With runtime validation (zod etc.):
class Validated extends WebsocketClientActor<ClientMessage, ServerMessage> {
constructor() {
const webSocketClientOptions = WebsocketClientOptions.create<ClientMessage, ServerMessage>()
.withUrl('wss://...')
.withCodec(jsonCodec<ClientMessage, ServerMessage>({
validate: (value: unknown): ServerMessage => ServerMessageSchema.parse(value),
}));
super(
webSocketClientOptions,
);
}
onMessage(message: ServerMessage): void { /* already validated */ }
}

Für Binärprotokolle ist rawCodec() die Ausweichluke — TOut = TIn = WebsocketFrame, sodass du rohe Text-/Binär-Frames selbst behandelst:

import { match } from 'ts-pattern';
import { rawCodec, type WebsocketFrame } from 'actor-ts/http';
class Binary extends WebsocketClientActor<WebsocketFrame, WebsocketFrame> {
constructor() {
const webSocketClientOptions = WebsocketClientOptions.create<WebsocketFrame, WebsocketFrame>()
.withUrl('wss://...')
.withCodec(rawCodec());
super(webSocketClientOptions);
}
onMessage(frame: WebsocketFrame): void {
match(frame)
.with({ kind: 'binary' }, (f) => this.onBinary(f))
.otherwise(() => {});
}
private onBinary(frame: BinaryFrame): void {
this.handleBytes(frame.data); // Uint8Array
}
}

Ein WebsocketFrame ist { kind: 'text'; data: string } oder { kind: 'binary'; data: Uint8Array }. Dekodier-Fehler werfen einen WebsocketDecodeError; das Frame-Größenlimit (maxFrameBytes) wird am rohen Frame vor dem Dekodieren durchgesetzt.

data ist auf jeder Laufzeitumgebung ein schlichtes Uint8Array. Der Client setzt beim Verbindungsaufbau binaryType des Sockets auf 'arraybuffer', denn die Laufzeitumgebungen sind sich über den Standardwert uneins — Bun liefert einen Node-Buffer, Node und Deno liefern ein Blob —, und ein Blob gibt seine Bytes nur asynchron heraus, was der Eingangspfad nicht bezahlen kann: Darauf zu warten würde die maxFrameBytes-Prüfung hinter einen Microtask schieben und die Garantie aufgeben, dass Frames in Ankunftsreihenfolge bei onMessage landen. Zwei Folgen sind wissenswert. Die Bytes gehören immer dem Frame selbst und sind nie ein Ausschnitt aus einem gepoolten Puffer, frame.data.buffer darfst du also weiterreichen. Und der Wert ist auch auf Bun ein Uint8Array und kein Buffer — für toString('hex') und Verwandte greifst du also zu Buffer.from(frame.data).

Ein zu großer Frame schließt die Verbindung mit Code 1009 (“Message Too Big”) und startet den üblichen Reconnect, statt an Ort und Stelle verworfen zu werden. Sei dir im Klaren darüber, was das bringt und was nicht: Die Allokation verhindert es nicht. Anders als die Server-Backends hat der Client kein Transport-Limit, an das er die Grenze durchreichen könnte — kein natives WebSocket der unterstützten Runtimes beachtet ein Payload-Limit, das man seinem Konstruktor übergibt. Wenn maxFrameBytes geprüft wird, hat die Runtime den Frame also längst vollständig auf dem Heap zusammengesetzt. Was das Schließen bringt: Das Gegenüber kann diesen Heap nicht erneut auf derselben Verbindung ausgeben — ohne den Close ließe sich eine einzige Verbindung dazu bringen, die Allokation für jeden beliebigen weiteren Frame zu wiederholen. Eine weitere Runde kostet das Gegenüber jetzt einen vollständigen Reconnect, den der reconnect-Backoff und der Circuit Breaker bereits drosseln. Siehe Untrusted Input begrenzen.

Begleitet wird der Close von einer Warnzeile. Sie benennt die Verbindung mit einem redigierten Label (wss://feed.example.com/ws/orders) statt mit url: Userinfo und Query-String fallen weg, denn ein WebSocket-Endpunkt wird üblicherweise mit einem ?token=… authentifiziert, und der Pfad bleibt erhalten, weil er zwei Verbindungen zum selben Host unterscheidet. Dieselbe Maskierung steht deinem eigenen Code als redactedUrlLabel zur Verfügung.

Eine unerwartete Trennung übernimmt der BrokerActor-Lifecycle:

  • Der Actor wechselt in disconnected und löst BrokerDisconnected auf dem Event-Stream aus.
  • Der Reconnect-Zyklus startet, mit exponentiellem Backoff gemäß den reconnect-Settings (initialDelayMs, maxDelayMs, factor, maxAttempts) — gejittert um randomFactor (Default ±20 %), damit eine Flotte nicht im Gleichschritt reconnectet.
  • Nach einem erfolgreichen Reconnect läuft onConnected() und BrokerConnected wird ausgelöst.

Ausgehende Nachrichten, die während der Trennung gesendet werden, werden gemäß der Outbound-Buffer-Policy gepuffert und nach dem Reconnect erneut gesendet. Der geerbte Circuit-Breaker löst nach wiederholten Verbindungsfehlern aus, damit ein toter Endpoint die Reconnect-Schleife nicht endlos drehen lässt.

class KeepAlive extends WebsocketClientActor<ClientMessage, ServerMessage> {
constructor() {
const webSocketClientOptions = WebsocketClientOptions.create<ClientMessage, ServerMessage>()
.withUrl('wss://...')
.withPingIntervalMs(30_000);
super(webSocketClientOptions);
}
onMessage(message: ServerMessage): void {}
}

Ein Ping auf Anwendungsebene alle pingIntervalMs verhindert, dass zwischengeschaltete Proxies / Load-Balancer Idle-Verbindungen schließen — die meisten produktiven Setups wollen das. Der WebSocket-Ping auf Protokollebene (Control-Frame) ist separat und wird von der Runtime übernommen; der Ping auf Anwendungsebene ist für Proxies, die Payloads inspizieren, keine Frames.

Ein Ping verhindert, dass ein Proxy die Verbindung schließt. Er sagt dir nicht, wann eine trotzdem gekappt wurde — ein verfallener NAT-Eintrag, ein Load-Balancer, der den Flow vergessen hat, ein mit SIGKILL beendeter Peer. Keine Seite wird benachrichtigt, also feuern close und error nie, onDisconnected läuft nie, und der Actor meldet connected, während jedes send ins Leere geht.

idleTimeoutMs ist die Deadline, die diesen Zustand beendet:

const webSocketClientOptions = WebsocketClientOptions.create<ClientMessage, ServerMessage>()
.withUrl('wss://realtime.example.com/feed')
.withIdleTimeoutMs(90_000) // no inbound frame for 90 s → reconnect
.withConnectTimeoutMs(10_000); // handshake stalled for 10 s → fail the attempt

Läuft sie ab, wird der Socket geschlossen, onDisconnected läuft mit einer idle timeout-Ursache, und der gewöhnliche Reconnect-Zyklus startet — derselbe Pfad, den ein beobachtetes close nimmt.

Gezählt werden Anwendungsnachrichten, keine Pongs. Ein Pong auf Protokollebene wird auf keiner unterstützten Runtime als message-Event zugestellt, pingIntervalMs frischt diese Deadline also nicht auf. Bemiss idleTimeoutMs am Heartbeat-Intervall des Servers, nicht an deinem Ping-Intervall; darunter kappt die Deadline gesunde Verbindungen in Endlosschleife. Genau deshalb ist sie standardmäßig aus.

connectTimeoutMs deckt das andere Ende ab: der Handshake endet nur über open oder error, also hält ein Server, der die TCP-Verbindung annimmt und den Upgrade nie abschließt, den Actor unbegrenzt in connecting.

Client-Defaults lassen sich unter actor-ts.io.broker.websocket setzen:

actor-ts.io.broker.websocket {
url = "wss://realtime.example.com/feed"
protocols = ["my-proto-v1"]
pingIntervalMs = 30000
idleTimeoutMs = 90000
connectTimeoutMs = 10000
maxFrameBytes = 1048576
reconnect { initialDelayMs = 500, maxDelayMs = 30000, factor = 2.0, randomFactor = 0.2 } // maxAttempts omitted = unlimited (the default)
circuitBreaker { /* ... */ }
outboundBuffer { /* ... */ }
}

Die Präzedenz folgt der Broker-Konvention: Konstruktor-Argument > HOCON > eingebaute Defaults.

Drei primäre Einsatzfälle:

  1. Echtzeit-Datenfeeds subscribieren — Aktien-Ticker, Krypto-Börsen, Social-Media-Streams.
  2. Ausgehende WebSocket-Clients zu Drittservices (Slack RTM, Streaming-LLM-APIs, IoT-Vendor-APIs).
  3. Eigene WS-basierte Protokolle, die darauf aufsetzen — nutze rawCodec() für binäre Wire-Formate.

Für serverseitiges WebSocket (Client-Verbindungen annehmen) siehe die Server-WebSocket-Seite.

  • I/O-Übersicht — das große Bild.
  • Server-WebSocket — die websocket()-Direktive + WebsocketServerActor.
  • SSE — Server-Sent Events als einseitige Streaming-Alternative.
  • BrokerActor-Basis — der gemeinsame Reconnect-/Buffer-/Circuit-Breaker-Lifecycle.