Zum Inhalt springen
Deutsch

Server-WebSocket

Serverseitiges WebSocket ist Teil des HTTP-Route-DSL. Die websocket()-Direktive upgradet passende Requests und verbindet jede Verbindung mit einem einzelnen WebsocketServerActor — dem Hub — den du implementierst. Für ausgehende Client-Verbindungen siehe WebsocketClientActor.

import {
ActorSystem,
} from 'actor-ts';
import {
HttpExtensionId,
WebsocketServerActor,
websocket,
type WebsocketConnection,
} from 'actor-ts/http';
import { match } from 'ts-pattern';
type SetNameMessage = { kind: 'setName'; name: string };
type SayMessage = { kind: 'say'; text: string };
type ClientMessage = SetNameMessage | SayMessage;
type ServerMessage = { kind: 'system'; text: string } | { kind: 'chat'; from: string; text: string };
class ChatRoom extends WebsocketServerActor<ServerMessage, ClientMessage> {
private readonly names = new Map<string, string>();
onMessage(message: ClientMessage): void {
match(message)
.with({ kind: 'setName' }, (m) => this.onSetName(m))
.with({ kind: 'say' }, (m) => this.onSay(m))
.exhaustive();
}
private onSetName(m: SetNameMessage): void {
this.names.set(this.connection.id, m.name);
this.reply({ kind: 'system', text: `hi ${m.name}` });
}
private onSay(m: SayMessage): void {
this.broadcast({ kind: 'chat', from: this.names.get(this.connection.id) ?? 'anon', text: m.text });
}
override onClientDisconnected(c: WebsocketConnection<ServerMessage>): void {
this.names.delete(c.id);
}
}
const system = ActorSystem.create('chat');
const chat = system.spawn(ChatRoom, 'chat');
await system.extension(HttpExtensionId)
.newServerAt('0.0.0.0', 8080)
.bind(websocket('/ws', chat));

Die Typparameter des Hubs sind aus Sicht des Servers zu lesen: TOut = die Nachrichten, die der Server sendet, TIn = die dekodierten Nachrichten, die er empfängt.

websocket() erzeugt eine Route und komponiert daher mit dem restlichen Route-DSL — path(), concat() und withMiddleware(). Zwei Formen:

import { websocket, path, concat } from 'actor-ts/http';
// 1. Bare directive — mount it under a path yourself:
path('ws', websocket(chat));
// 2. Path sugar — equivalent to path(p, websocket(target)):
websocket('/ws', chat);
// Mixed with normal HTTP routes on the same server:
concat(
path('api', apiRoutes),
websocket('/ws', chat),
);
type WebsocketRouteOptions<TOut, TIn> = {
codec?: WebsocketCodec<TOut, TIn>; // default jsonCodec()
maxFrameBytes?: number; // default 1 MiB
onOversizeFrame?: 'close' | 'drop'; // default 'close' (1009)
onInvalidMessage?: 'close' | 'drop' | 'hook'; // default 'close' (1003)
maxBufferedBytes?: number; // default 4 MiB
onBackpressure?: 'drop' | 'close'; // default 'drop'
allowedOrigins?: string[]; // CSWSH-Schutz — siehe unten
maxConnections?: number; // Limit gleichzeitiger Verbindungen; Default unbegrenzt
maxPreAttachFrames?: number; // Frame-Limit im Aufbaufenster; Default 256
maxPreAttachBytes?: number; // Byte-Limit im Aufbaufenster; Default 4 MiB
acceptTimeoutMs?: number; // Aufbau-Frist; Default 10 000, Infinity schaltet ab
};
const webSocketRouteOptions = WebsocketRouteOptions.create()
.withMaxFrameBytes(256 * 1024)
.withOnOversizeFrame('close') // close with 1009 Message Too Big
.withOnInvalidMessage('close');
websocket('/ws', chat, webSocketRouteOptions); // close with 1003 Unsupported Data

Das Frame-Größenlimit wird am rohen Frame vor dem Dekodieren durchgesetzt. onInvalidMessage: 'hook' leitet Dekodier-Fehler an den onInvalidMessage-Hook des Actors weiter, statt zu schließen.

maxFrameBytes wird zweimal geprüft, zu zwei verschiedenen Zeitpunkten:

  • Im Transport, bei deinem aufgelösten maxFrameBytes — das Framework reicht es beim Binden an den Socket der Laufzeit weiter. Ein größerer Frame wird abgewiesen, während er eintrifft; ein feindlicher Peer kann den Prozess also nicht dazu bringen, 16 MiB (Bun) oder 100 MiB (ws) zu puffern, nur damit das Ergebnis anschließend verworfen wird. Weil die Laufzeit dabei auflegt, statt einen Policy-Close zu senden, sieht der Client hier einen abnormalen Close (1006) und keine 1009. Zwei Laufzeit/Backend-Paare haben diese Ebene nicht — siehe den Warnhinweis weiter unten.
  • Im Connection-Actor, bei demselben maxFrameBytes — am vollständig empfangenen Frame, vor dem Dekodieren, mit Close 1009.

Nur die zweite Ebene ist eine Zusage. Sie läuft überall und hält den zu großen Frame von deinem Actor fern; das Transport-Limit ist eine Optimierung darüber, die dem Prozess die Pufferung erspart.

Die Zahl ist in beide Richtungen deine: maxFrameBytes zu senken verkleinert auch das Puffer-Fenster, und es über den 1-MiB-Default anzuheben, lässt Frames dieser Größe tatsächlich durch. Aufgelöst wird wie üblich — Route-Options, dann actor-ts.http.websocket.maxFrameBytes, dann der eingebaute Default —, und zwar beim Binden des Servers; ein ungültiger Wert ist damit ein Fehler in bind() und nicht erst beim ersten Upgrade.

Das Transport-Limit gilt pro Server, nicht pro Route. Ein Server hat genau einen Transport für alle seine WebSocket-Routen (ein WebSocketServer auf Express, eine Plugin-Registrierung auf Fastify, ein Bun.serve). Wo mehrere Routen sich unterscheiden, nimmt der Transport deshalb das größte ihrer Limits. Was eine strengere Route annimmt, ändert das nicht — der Connection-Actor weist ihre zu großen Frames weiterhin mit 1009 ab —, sie bekommt nur nicht das engere Puffer-Fenster, das sie allein gehabt hätte.

Das ist eine Entscheidung und keine Lücke. Fastify und Bun.serve können immer nur ein einziges Limit halten; Express könnte strukturell eines pro Route halten, und es dort zu tun hieße, dass dieselbe Konfiguration je nach Backend etwas anderes bedeutet — bei einem Limit, dessen ganze Aufgabe es ist zu begrenzen, was eine feindliche Gegenstelle allozieren kann, ist eine Form überall mehr wert als ein engeres Fenster auf einem Backend von dreien.

Zwischen dem Abschluss des Handshakes und dem Moment, in dem der Connection-Actor seine Listener anhängt, liegt ein kurzes Fenster — zwei Mailbox-Schritte —, in dem der Socket lebt und niemand mitliest. Frames, die darin ankommen, werden gehalten und in dem Moment nachgespielt, in dem der Actor anhängt; verloren geht also nichts. Zwei Stellschrauben begrenzen, was dieses Fenster kosten darf, und beide gibt es, weil sein Ende nicht garantiert ist: den Actor spawnt der Hub, und ein Hub, der gestoppt wurde, hängt oder umkonfiguriert wurde, liefert unter Umständen nie einen.

  • maxPreAttachFrames / maxPreAttachBytes begrenzen, was gehalten wird. Jenseits davon wird der Socket mit 1013 (“try again later”) geschlossen statt weiter gepuffert. Ein legitimer Client schickt in dieses Fenster nichts oder eine Begrüßung, die Defaults — 256 Frames, 4 MiB — liegen also um Größenordnungen über normalem Verkehr; abgewiesen wird ein Peer, der in eine noch gar nicht angenommene Verbindung hineinstreamt. Der erste Frame wird immer angenommen, unabhängig von seiner Größe (maxFrameBytes begrenzt ihn ohnehin bereits) — eine Route, die maxFrameBytes über 4 MiB hebt, muss maxPreAttachBytes also nur mitheben, wenn ihre Clients mehrere solcher Frames auf einmal senden.
  • acceptTimeoutMs begrenzt, wie lange das Fenster offen bleiben darf. Hat sich innerhalb dieser Frist kein Connection-Actor angehängt, wird der Socket mit 1013 geschlossen und sein maxConnections-Platz freigegeben. Standardmäßig zehn Sekunden: ein gesunder Hub hängt sich in Mikrosekunden an, das hier ist also ein Auffangnetz und keine Liveness-Regel — und ein Actor, der nach der Frist doch noch auftaucht, bekommt das Close statt eines Sockets, den das Framework längst geschlossen hat. Infinity schaltet es ab.

Beides ist nicht versehentlich erreichbar. Beides beantwortet dieselbe Störung: ein aufgewerteter Socket ohne Actor dahinter, der vorher für die Lebensdauer des Prozesses offen blieb, Frames sammelte, die niemand je abholt, und einen Platz belegte, den niemand je zurückgibt.

Browser hängen die Cookies des Nutzers automatisch an einen WebSocket-Upgrade an — eine WS-Route mit ambienter Authentifizierung (Session-Cookie oder IpAllowlist) ist damit anfällig für Cross-Site-WebSocket-Hijacking: eine fremde Seite öffnet new WebSocket('wss://dein-host/ws') und nutzt die Credentials des Opfers.

Mit allowedOrigins wird der Handshake über den Origin-Header des Browsers abgesichert:

const wsOptions = WebsocketRouteOptions.create()
.withAllowedOrigins(['https://app.example.com']);
websocket('/ws', chat, wsOptions);

Ein Upgrade, dessen Origin vorhanden, aber nicht gelistet ist, wird auf allen drei Backends vor dem Handshake mit 403 abgelehnt. Ein fehlender Origin (Nicht-Browser-Client — natives WebSocket, Server-zu-Server) wird zugelassen, da CSWSH ein reiner Browser-Angriff ist. Der Vergleich ist case-insensitiv.

BearerTokenAuth ist bereits geschützt (Browser können beim WS-Handshake keinen Authorization-Header setzen); allowedOrigins ist v. a. bei Cookie- oder IP-basierter Auth relevant.

withMiddleware() komponiert mit websocket(), aber die Middleware läuft einmal, zur Upgrade-Zeit, gegen den HTTP-Upgrade-Request — nicht pro WebSocket-Nachricht. Eine ablehnende Middleware gibt einen normalen HTTP-Fehler zurück und das Upgrade kommt nie zustande, sodass Auth-Middleware wie BearerTokenAuth oder IpAllowlist den Handshake absichert:

import { websocket, withMiddleware } from 'actor-ts/http';
withMiddleware(BearerTokenAuth({ /* ... */ }), websocket('/ws', chat));
// A bad token → HTTP 401, no upgrade. A good token → the socket opens.

Die Middleware-Sammlung steht in der HTTP-Übersicht.

Antwort-dekorierende Middleware funktioniert ebenfalls

Abschnitt betitelt „Antwort-dekorierende Middleware funktioniert ebenfalls“

Middleware, die next() aufruft und eine dekorierte Kopie des Ergebnisses zurückgibt — securityHeaders(), contentSecurityPolicy(), strictTransportSecurity(), requestId(), csrfProtection() —, komponiert ebenso über einer websocket()-Route. Ein einziges Wrap kann damit einen ganzen Baum härten, ohne den Socket herauszuschneiden:

import { requestId, securityHeaders, websocket, withMiddleware } from 'actor-ts/http';
withMiddleware(securityHeaders(), withMiddleware(requestId(), websocket('/ws', chat)));
// Both run at upgrade time. The socket opens.

Die Header, die solche Middleware setzt, erreichen einen erfolgreichen Handshake nicht — diese 101-Antwort schreibt das Backend selbst, und eine WebSocket-Verbindung ist kein Dokument, auf das ein Browser X-Frame-Options oder eine CSP anwenden würde. Auf einer Ablehnung fahren sie dagegen mit: Ein 401 einer inneren BearerTokenAuth kommt mit den Security-Headern zurück — und genau das ist der Fall, auf den es ankommt, denn diese Antwort rendert ein Browser tatsächlich.

withMiddleware() ist der portable Weg — auf jedem Backend identisch —, aber ein Handshake ist auch für das darunterliegende Framework ein ganz normaler Request, und dessen eigene Middleware läuft ebenfalls:

BackendWas beim Handshake zusätzlich läuft
FastifyRoute-Hooks, inklusive preValidation
Expressalles, was mit app.use(...) registriert wurde
Honoalles, was mit app.use(...) registriert wurde

Eine App, die /ws mit app.use(requireLogin) absichert, ist also wirklich abgesichert, und eine native Middleware, die den Request beantwortet, lehnt das Upgrade ab. Der Guard des Frameworks läuft nach allen anderen — withMiddleware() und allowedOrigins haben damit immer das letzte Wort.

Nimm withMiddleware(), wenn sich eine Route über alle Backends hinweg gleich verhalten soll; greif zu nativer Middleware, wenn du ein bestehendes Ökosystem weiterverwendest (Sessions, Rate Limiting, Observability).

Ein Actor pro Route — der Hub. Er sieht die Events jeder Verbindung, serialisiert:

abstract class WebsocketServerActor<TOut, TIn, TSelf = never> {
// You implement:
abstract onMessage(message: TIn): void | Promise<void>;
// Optional overrides:
protected onClientConnected(client: WebsocketConnection<TOut>): void;
protected onClientDisconnected(client: WebsocketConnection<TOut>, info: WebsocketCloseInfo): void;
protected onInvalidMessage(client: WebsocketConnection<TOut>, error: WebsocketDecodeError): void;
protected onSelfMessage(message: TSelf): void;
}

Innerhalb von onMessage und den Hooks stehen zur Verfügung:

MemberWas er tut
this.connectionDie WebsocketConnection<TOut>, deren Event gerade verarbeitet wird.
this.reply(message)Sendet message an die aktuelle Verbindung.
this.broadcast(message, filter?)Sendet an jede Verbindung (optional per Prädikat gefiltert).
this.clientsReadonlyMap<string, WebsocketConnection<TOut>> aller aktiven Verbindungen.
this.closeAll(code?, reason?)Schließt jede Verbindung.
onMessage(message: ClientMessage): void {
this.reply({ kind: 'system', text: 'got it' }); // → the sender
this.broadcast({ kind: 'chat', from: 'x', text: 'hi' }); // → everyone
this.broadcast(notice, (c) => c.id !== this.connection.id); // → everyone else
}

Jedes Event einer gegebenen Verbindung wird über den einen Hub-Actor in dieser Reihenfolge serialisiert:

onClientConnected → onMessage* (in frame order) → onClientDisconnected

onClientConnected läuft einmal, dann null oder mehr onMessage-Aufrufe in Frame-Reihenfolge, dann genau ein onClientDisconnected. Da alles auf einem einzelnen Actor läuft, bekommst du die übliche Garantie des Actor-Modells — keine nebenläufige Handler-Ausführung, keine Locks.

Ein Hub trägt jeden eingehenden Frame jeder Verbindung der Route und ist damit der Musterfall für eine begrenzte Mailbox: ein Actor, der einem Produzenten ausgesetzt ist, den du nicht kontrollierst. Auf ihm läuft aber auch der verbindungsinterne Verkehr des Frameworks. Daraus folgen zwei Dinge.

Das Kommando, das den Actor einer Verbindung spawnt, kann nicht verworfen werden. Es nimmt dieselbe Spur wie ein Terminated der Death-Watch — es wird wie alles andere am Ende der Schlange eingereiht, ist aber von jeder Overflow-Policy ausgenommen, sodass drop-head, drop-new, reject und auch eine eigene PriorityMailbox es unangetastet lassen. Das ist wichtig, weil dieses Kommando genau einmal gesendet wird — aus dem Upgrade-Callback des Backends — und die einzige Referenz auf einen Socket trägt, dessen Handshake bereits abgeschlossen ist. Es zu verlieren würde diesen Socket aufgewertet zurücklassen, ohne dass etwas daran hängt: keine Listener, eingehende Frames stapeln sich ungelesen, und sein maxConnections-Platz bleibt belegt, bis der Client aufgibt. Einen Hub, der aus einem anderen Grund nicht antworten kann — er ist gestoppt, oder er kommt nie zu dem Kommando —, fängt stattdessen das Aufbaufenster ab: es schließt den Socket und gibt den Platz zurück.

Was eine Begrenzung dagegen weiterhin verwirft, ist der Rest. Eingehende Frames und die Signale onClientConnected / onClientDisconnected, die die Connection-Actors melden — ein begrenzter Hub unter Last kann dir also einen Client liefern, der nie in this.clients erscheint, oder einen, der sie nie verlässt. Beides ist dein Protokoll, und nur du weißt, ob ein Verlust akzeptabel ist.

Jede Verbindung ist eine WebsocketConnection<TOut>, die ActorRef<TOut> erweitert:

interface WebsocketConnection<TOut> extends ActorRef<TOut> {
readonly id: string;
readonly remoteAddress?: string;
readonly upgrade: WebsocketUpgradeInfo;
readonly isOpen: boolean;
tell(message: TOut): void; // encode via codec + send
sendRaw(frame: WebsocketFrame): void; // bypass the codec
close(code?: number, reason?: string): void;
}

upgrade trägt den Handshake-Kontext:

type WebsocketUpgradeInfo = {
path: string;
params: Record<string, string>;
query: Record<string, string | string[] | undefined>;
headers: Record<string, string>;
remoteAddress?: string;
subprotocol?: string;
};

WebsocketCloseInfo (an onClientDisconnected übergeben) ist { code: number; reason: string; initiatedBy: 'client' | 'server' | 'error' }.

Der Route-Codec dekodiert eingehende Frames in TIn und kodiert TOut-Antworten. Der Default ist jsonCodec():

import { WebsocketRouteOptions, jsonCodec } from 'actor-ts/http';
const webSocketRouteOptions = WebsocketRouteOptions.create().withCodec(jsonCodec<ServerMessage, ClientMessage>({
validate: (v: unknown): ClientMessage => ClientMsgSchema.parse(v),
}));
websocket('/ws', chat, webSocketRouteOptions);

Für Binärprotokolle liefert rawCodec() rohe Frames (TOut = TIn = WebsocketFrame):

import { match } from 'ts-pattern';
import {
rawCodec,
WebsocketRouteOptions,
WebsocketServerActor,
type WebsocketFrame,
} from 'actor-ts/http';
const webSocketRouteOptions = WebsocketRouteOptions.create().withCodec(rawCodec());
class BinaryHub extends WebsocketServerActor<WebsocketFrame, WebsocketFrame> {
onMessage(frame: WebsocketFrame): void {
match(frame)
.with({ kind: 'binary' }, (f) => this.onBinary(f))
.otherwise(() => {});
}
private onBinary(frame: BinaryFrame): void {
this.reply({ kind: 'binary', data: process(frame.data) });
}
}
websocket('/stream', binaryHub, webSocketRouteOptions);

Dekodier-Fehler werfen WebsocketDecodeError und folgen der onInvalidMessage-Policy der Route.

websocket() funktioniert auf allen drei HTTP-Backends:

BackendPeer-DependencyRuntime-Hinweise
Fastify (Default)@fastify/websocketLäuft auf Bun und Node; auf Deno bevorzuge Hono.
Expressws—
HonoBun & Deno: in hono eingebaut. Hono-auf-Node: @hono/node-ws.Helfer pro Runtime.
Terminal-Fenster
npm install @fastify/websocket # Fastify (default)
npm install ws # Express
npm install @hono/node-ws # Hono on Node only, 1.2.0 or newer

Route-Defaults liegen unter actor-ts.http.websocket:

actor-ts.http.websocket {
maxFrameBytes = 1048576
onOversizeFrame = close
onInvalidMessage = close
maxBufferedBytes = 4194304
onBackpressure = drop
maxPreAttachFrames = 256
maxPreAttachBytes = 4M
acceptTimeoutMs = 10s
}

Präzedenz: Route-Optionen > HOCON > eingebaute Defaults.

Drei gute Einsatzfälle:

  1. Echtzeit-UIs — Updates an Browser-Clients pushen.
  2. Eigene Messaging-Protokolle — Game-Server, Chat-Backends.
  3. Binäres Streaming — nutze rawCodec() und behandle Frames direkt.

Für einseitige Streams vom Server zum Client ist SSE einfacher. Für Request/Reply-RPC ist einfaches HTTP die Norm.