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.
Settings
Abschnitt betitelt „Settings“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.
Senden und Empfangen
Abschnitt betitelt „Senden und Empfangen“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.
Von einem anderen Actor senden
Abschnitt betitelt „Von einem anderen Actor senden“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:
| Hook | Wann |
|---|---|
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.
Reconnect und Buffering
Abschnitt betitelt „Reconnect und Buffering“Eine unerwartete Trennung übernimmt der BrokerActor-Lifecycle:
- Der Actor wechselt in
disconnectedund löstBrokerDisconnectedauf dem Event-Stream aus. - Der Reconnect-Zyklus startet, mit exponentiellem Backoff gemäß
den
reconnect-Settings (initialDelayMs,maxDelayMs,factor,maxAttempts) — gejittert umrandomFactor(Default ±20 %), damit eine Flotte nicht im Gleichschritt reconnectet. - Nach einem erfolgreichen Reconnect läuft
onConnected()undBrokerConnectedwird 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.
Ping / Keep-Alive
Abschnitt betitelt „Ping / Keep-Alive“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.
Eine still gekappte Verbindung erkennen
Abschnitt betitelt „Eine still gekappte Verbindung erkennen“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 attemptLä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.
HOCON-Config
Abschnitt betitelt „HOCON-Config“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.
Wann einsetzen
Abschnitt betitelt „Wann einsetzen“Drei primäre Einsatzfälle:
- Echtzeit-Datenfeeds subscribieren — Aktien-Ticker, Krypto-Börsen, Social-Media-Streams.
- Ausgehende WebSocket-Clients zu Drittservices (Slack RTM, Streaming-LLM-APIs, IoT-Vendor-APIs).
- 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.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- 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.
