Zum Inhalt springen
Deutsch

TCP

Zwei Actors, ein Protokoll:

ActorRichtungBesitzt
TcpSocketActorausgehend — wählt einen Remote-Hosteine Verbindung
TcpServerActoreingehend — bindet einen lokalen Portden Listener + jede angenommene Verbindung

Beide erben von BrokerActor, teilen sich also Lifecycle, Reconnect-Policy und die BrokerConnected- / BrokerDisconnected-Events, und beide schneiden eingehende Bytes mit denselben Framing-Strategien.

import { ActorSystem } from 'actor-ts';
import { TcpSocketActor, TcpSocketOptions } from 'actor-ts/io';
const tcpSocketOptions = TcpSocketOptions.create()
.withHost('metrics-collector.example.com')
.withPort(8125)
.withTarget(protocolHandler); // erforderlich: wohin eingehende Frames gehen
const tcp = system.spawn(() => new TcpSocketActor(tcpSocketOptions), 'tcp-client');
// Rohe Bytes senden:
tcp.tell({ kind: 'send', payload: new Uint8Array([0x01, 0x02, 0x03]) });
// Oder ein String (UTF-8-kodiert für dich):
tcp.tell({ kind: 'send', payload: 'PING\n' });
interface TcpSocketOptionsType extends BrokerCommonOptionsType {
host?: string; // erforderlich
port?: number; // erforderlich
framing?: TcpFraming; // Frame-Extraktion; Default { kind: 'bytes' }
target?: ActorRef<unknown>; // erforderlich: wohin eingehende Frames gehen
idleTimeoutMs?: number; // Lese-Idle-Deadline; Default 0 (aus)
connectTimeoutMs?: number; // Connect-Deadline; Default 0 (aus)
keepAliveMs?: number; // OS-Keepalive-Verzögerung; Default 45_000, 0 = aus
}

Eingehende Frames werden direkt an den target-Actor gepusht, den du über .withTarget(..) verdrahtest — es gibt kein subscribe-Kommando und keinen Envelope. Jede Nachricht ist der Frame: ein Uint8Array bei bytes- / length-prefixed- Framing, ein string bei lines.

class ProtocolHandler extends Actor<Uint8Array> {
override onReceive(frame: Uint8Array): void {
this.handleBytes(frame);
}
}

Ohne Framer ({ kind: 'bytes' }, der Default) ist jede Nachricht ein roher Chunk Bytes — KEINE logische Nachricht. TCP ist ein Byte-Stream; Framing ist dein Job.

Der Verbindungs-Lifecycle wird nicht an das Target geliefert — er wird auf system.eventStream als BrokerConnected- / BrokerDisconnected-Events veröffentlicht, die alle Broker-Actors teilen:

import { BrokerConnected, BrokerDisconnected } from 'actor-ts/io';
system.eventStream.subscribe(monitorRef, BrokerConnected);
system.eventStream.subscribe(monitorRef, BrokerDisconnected);

close und error sind die einzigen Events, die ein Socket auslöst — und ein Peer, der ohne FIN oder RST verschwindet, löst keines von beiden aus. Der Socket bleibt offen, der Actor bleibt connected, und jedes Senden landet in einem Kernel-Puffer, den nie jemand leert. Drei Stellschrauben adressieren das:

const tcpSocketOptions = TcpSocketOptions.create()
.withHost(host)
.withPort(port)
.withTarget(protocolHandler)
.withIdleTimeoutMs(90_000) // no inbound bytes for 90 s → reconnect
.withConnectTimeoutMs(10_000); // handshake stalled for 10 s → fail the attempt
  • keepAliveMs ist die einzige, die standardmäßig an ist (45 s). Sie schaltet TCP-Keepalive auf OS-Ebene ein; dessen Probes beantwortet der Kernel des Peers, unabhängig davon, ob dessen Anwendung etwas zu sagen hat — sie kann sich über eine gesunde Verbindung also nie irren, sondern höchstens langsam sein: wie lange eine tote nach der Verzögerung braucht, um aufzufallen, entscheidet das OS (Linux: neun Probes im Abstand von 75 s). 0 schaltet sie aus.
  • idleTimeoutMs ist eine Lese-Deadline, zurückgesetzt von eingehenden Bytes und ausdrücklich nicht von ausgehenden — ein Client, der in ein schwarzes Loch schreibt, ist genau der Fall, für den sie existiert. Standardmäßig aus: nur du weißt, wie lange dein Peer still sein darf, und ein Wert unterhalb seines eigenen Heartbeat-Intervalls kappt gesunde Verbindungen in Endlosschleife.
  • connectTimeoutMs begrenzt einen einzelnen Verbindungsversuch. Ohne sie hält ein Peer, der den TCP-Handshake abschließt und dann stehen bleibt, den Actor beliebig lange in connecting.

Beide Deadlines laufen in handleConnectionLost, sodass Reconnect-Policy, Outbound-Puffer und Circuit Breaker sich exakt so verhalten wie bei einem beobachteten Abbruch — siehe BrokerActor-Basis.

import { TcpSocketActor, TcpSocketOptions } from 'actor-ts/io';
const tcpSocketOptions = TcpSocketOptions.create()
.withHost(host)
.withPort(port)
.withFraming({ kind: 'length-prefixed' })
.withTarget(subscriber);
new TcpSocketActor(tcpSocketOptions);

Framing wählst du über .withFraming(..) mit einer TcpFraming-Config-Union — es gibt keine Framer-Klassen zum Instanziieren. Drei Strategien werden mitgeliefert:

type TcpFraming =
| { kind: 'bytes' } // Default — rohe Chunks
| { kind: 'lines'; delimiter?: string; maxLineLen?: number }
| { kind: 'length-prefixed'; maxFrameLen?: number };
  • bytes (Default) — jeder Chunk wird roh geliefert; das Target behandelt die Byte-Stream-Semantik selbst.
  • lines — split am delimiter (Default '\n'); jeder Frame kommt als string an. maxLineLen begrenzt eine Zeile in Bytes, die nicht terminierte eingeschlossen — diese zweite Hälfte ist es, die den Re-Assembly-Buffer gegen einen Peer begrenzt, der nie einen Delimiter schickt. Bytes, keine Zeichen: ein CJK-Zeichen kostet drei davon.
  • length-prefixed — die ersten 4 Bytes (Big-Endian uint32) tragen die Payload-Größe; die Prefix-Breite ist fest, nur maxFrameLen ist konfigurierbar.

Mit lines oder length-prefixed ist jede Nachricht, die das Target empfängt, ein vollständiger Frame, kein beliebiger Chunk.

Ein Frame über seinem Limit wird nie gepuffert: Die noch offenen Bytes werden verworfen und die Verbindung mit ihnen. Der Client zerstört seinen Socket und fällt auf seine Reconnect-Policy zurück, der Listener schließt nur die betroffene Verbindung.

Beide Actors validieren framing beim Start, ein falscher Wert wirft also einen OptionsError, bevor es einen Socket gibt — nicht erst beim ersten eingehenden Byte. Ein Limit, das keine positive Ganzzahl ist, wird abgelehnt, und ein leerer delimiter ebenso: Er passt an jeder Position, ohne etwas zu verbrauchen, und das ist kein langsamer Framer, sondern ein festgefahrener Prozess.

Beide Actors lesen dieselben Strategien aus demselben Code — auf TcpServerActor wird das gewählte Framing pro angenommener Verbindung angewandt, jede mit eigenem Re-Assembly-Buffer, sodass zwei Clients mitten im Frame nie ineinanderlaufen.

TcpServerActor bindet einen Port und bedient jede eingehende Verbindung. Er baut auf derselben runtime-übergreifenden TCP-Schicht auf wie der Cluster-Transport — Bun, Node und Deno, und TLS, kommen also aus einer Hand.

import { ActorSystem } from 'actor-ts';
import { TcpServerActor, TcpServerOptions } from 'actor-ts/io';
const tcpServerOptions = TcpServerOptions.create()
.withBindHost('0.0.0.0')
.withBindPort(9000)
.withFraming({ kind: 'lines' })
.withTarget(connectionHandler); // erforderlich: wohin Events + Frames gehen
const server = system.spawn(() => new TcpServerActor(tcpServerOptions), 'tcp-server');
// An genau eine Verbindung schreiben, adressiert über die erhaltene Id:
server.tell({ kind: 'send', connectionId, payload: 'PONG\n' });
// Eine Verbindung auflegen. Der Listener bedient den Rest weiter:
server.tell({ kind: 'close', connectionId });
interface TcpServerOptionsType extends BrokerCommonOptionsType {
bindHost?: string; // Default '0.0.0.0'
bindPort?: number; // erforderlich; 0 = das OS wählt
framing?: TcpFraming; // pro Verbindung; Default { kind: 'bytes' }
target?: ActorRef<TcpServerMessage>; // erforderlich
tls?: TlsTransportOptionsType; // TLS statt Klartext ausliefern
maxConnections?: number; // Aufnahme-Limit; Default Infinity
}

Mit bindPort: 0 wählt das OS den Port — lies ihn über boundPort des Actors zurück, sobald gebunden ist. connectionCount meldet die aktiven Verbindungen.

Das target empfängt eine kind-getaggte Union, keine nackten Frames — ein Listener hat viele Verbindungen, also benennt jede Nachricht die, zu der sie gehört:

type TcpServerMessage =
| { kind: 'connectionOpened'; connectionId: string; remoteAddress?: string }
| { kind: 'frame'; connectionId: string; payload: Uint8Array | string }
| { kind: 'connectionClosed'; connectionId: string };

Die Union und jede ihrer Varianten sind importierbar, ein Handler kann also die Variante nehmen, um die es ihm geht, statt die Union in jedem Arm einzuengen:

import type { FrameMessage, TcpServerCommand, TcpServerMessage } from 'actor-ts/io';

Ein Echo-Server sind dann die naheliegenden drei Zeilen:

class EchoHandler extends Actor<TcpServerMessage> {
constructor(private readonly server: ActorRef<TcpServerCommand>) { super(); }
override onReceive(message: TcpServerMessage): void {
match(message)
.with({ kind: 'connectionOpened' }, (m) => this.onConnectionOpened(m))
.with({ kind: 'frame' }, (m) => this.onFrame(m))
.with({ kind: 'connectionClosed' }, (m) => this.onConnectionClosed(m))
.exhaustive();
}
private onFrame(message: FrameMessage): void {
this.server.tell({ kind: 'send', connectionId: message.connectionId, payload: message.payload });
}
// …
}

connectionClosed kommt bei jedem Ende: der Peer hat aufgelegt, die Verbindung ist gescheitert, du hast close geschickt, ein Frame hat sein Größenlimit gerissen, oder der Actor wurde gestoppt und hat den Port freigegeben. Ein Signal, ein Codepfad.

tls trägt das Zertifikats-Material — PEM-Inhalt oder DER-Bytes, niemals einen Pfad. Ein gesetztes ca schaltet die Prüfung von Client-Zertifikaten (mTLS) ein; requestClientCert: false erzwingt einseitiges TLS.

const tcpServerOptions = TcpServerOptions.create()
.withBindPort(9443)
.withTls({ cert: readFileSync('server.pem', 'utf8'), key: readFileSync('server.key', 'utf8') })
.withTarget(connectionHandler);

Ein halb konfiguriertes Credential — ein cert ohne key, oder ein ca allein — wird beim Start des Actors abgelehnt, nicht still im Klartext gebunden. Für tls gibt es bewusst kein HOCON-Leaf: eine Config-Datei ist der falsche Ort für einen privaten Schlüssel.

  • maxConnections begrenzt die gleichzeitig angenommenen Verbindungen. Eine Verbindung, die am Limit ankommt, wird sofort abgebrochen statt registriert — Abweisen an der Tür, statt einen Socket anzunehmen, aus dem niemand liest. Abgebrochen, nicht ordentlich beendet: Ein geordnetes Schließen schickt nur ein FIN, und ein Peer, der es nie beantwortet, hält den Socket samt Dateideskriptor offen. Registriert wurde dieser Socket nie, das Limit zählt ihn also auch nicht — ein Halbschluss würde nur kooperierende Peers begrenzen. Der abgewiesene Peer sieht deshalb einen Verbindungsfehler statt eines sauberen Schließens; genau so sieht auf TCP-Ebene aus, an der Kapazitätsgrenze abgewiesen zu werden.
  • Ein Frame über dem framing-Limit trennt nur diese Verbindung. Der Client-Actor beantwortet denselben Verstoß damit, seine Verbindung fallen zu lassen — dort ist die eine Verbindung der gesamte Transport — und richtet sich danach nach seiner Reconnect-Policy; ein Listener, der das genauso hielte, ließe jeden einzelnen Client den Dienst für alle abschalten. In beiden Fällen gehen die noch gepufferten Bytes mit dem Socket: Sinn eines Limits ist nicht, den Verstoß zu protokollieren, sondern aufzuhören, das festzuhalten, was ihn ausgelöst hat.
  • outboundBuffer ist hier 0 — anders als bei jedem anderen Broker-Actor. Puffern im getrennten Zustand existiert, damit eine Nachricht einen Reconnect überlebt — aber “getrennt” heißt beim Listener, dass der Port unten ist, und damit ist jede angenommene Verbindung ohnehin weg: ein gepufferter Write benennt eine Id, die nie wiederkommt. Ein BrokerNotConnected-Event sagt das, ein Replay nicht.

Drei legitime Einsatzfälle:

  1. Mit Legacy-Protokollen sprechen — proprietäre Protokolle, für die es keine höherrangigen Wrapper gibt.
  2. Eigene Binärprotokolle — Game-Server, Metric-Collector mit eigenem Wire-Format.
  3. Brücken zu Nicht-HTTP-Services — Message-Queues mit proprietärem Wire (z. B. einige Finanzprotokolle).

Für neue Anwendungsprotokolle gilt: greife nicht zuerst zu raw TCP. HTTP, WebSocket oder gRPC sind für fast alles die besseren Startpunkte.

  • I/O-Übersicht — das große Bild.
  • BrokerActor-Basis — der gemeinsame Lifecycle.
  • UDP — die verbindungslose Alternative.
  • gRPC — typisierter RPC über HTTP/2 — meist die bessere Wahl als raw TCP.