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';
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 dropped with a warning
onInvalidMessage?: 'drop' | 'hook' | 'disconnect'; // default 'drop'
pingIntervalMs?: number; // application-level ping; default off
}

Da die Settings BrokerCommonOptionsType erweitern, gelten die Blöcke reconnect ({ initialDelayMs, maxDelayMs, factor, maxAttempts }), 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';
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.

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';
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';
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.

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

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
maxFrameBytes = 1048576
reconnect { initialDelayMs = 500, maxDelayMs = 30000, factor = 2.0 } // 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.