跳转到内容
简体中文

TCP

此内容尚不支持你的语言。

Two actors, one protocol:

ActorDirectionOwns
TcpSocketActoroutbound — dials a remote hostone connection
TcpServerActorinbound — binds a local portthe listener + every connection it accepts

Both extend BrokerActor, so they share the lifecycle, the reconnect policy and the BrokerConnected / BrokerDisconnected events, and both cut inbound bytes with the same framing strategies.

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); // required: where inbound frames go
const tcp = system.spawn(() => new TcpSocketActor(tcpSocketOptions), 'tcp-client');
// Send raw bytes:
tcp.tell({ kind: 'send', payload: new Uint8Array([0x01, 0x02, 0x03]) });
// Or a string (UTF-8 encoded for you):
tcp.tell({ kind: 'send', payload: 'PING\n' });
interface TcpSocketOptionsType extends BrokerCommonOptionsType {
host?: string; // required
port?: number; // required
framing?: TcpFraming; // frame extraction; default { kind: 'bytes' }
target?: ActorRef<unknown>; // required: where inbound frames go
idleTimeoutMs?: number; // read-idle deadline; default 0 (off)
connectTimeoutMs?: number; // connect deadline; default 0 (off)
keepAliveMs?: number; // OS keepalive delay; default 45_000, 0 = off
}

Inbound frames are pushed straight to the target actor you wire in through .withTarget(..) — there is no subscribe command and no envelope. Each message is the frame: a Uint8Array for bytes / length-prefixed framing, a string for lines.

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

Without a framer ({ kind: 'bytes' }, the default) each message is a raw chunk of bytes — NOT a logical message. TCP is a byte stream; framing is your job.

Connection lifecycle is not delivered to the target — it is published on system.eventStream as BrokerConnected / BrokerDisconnected events, shared by every broker actor:

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

close and error are the only events a socket raises, and a peer that disappears without FIN or RST raises neither. The socket stays open, the actor stays connected, and every send goes into a kernel buffer nothing will ever drain. Three knobs address that:

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 is the only one that is on by default (45 s). It turns on OS-level TCP keepalive, whose probes are answered by the peer’s kernel whether or not its application has anything to say — so it can never be wrong about a healthy connection, only slow: how long a dead one takes to surface after the delay belongs to the OS (Linux: nine probes, 75 s apart). 0 turns it off.
  • idleTimeoutMs is a read deadline, reset by inbound bytes and not by outbound ones — a client writing into a black hole is the case it exists for. Off by default: only you know how long your peer is allowed to stay quiet, and a value below its own heartbeat interval severs healthy connections in a loop.
  • connectTimeoutMs bounds one connect attempt. Without it, a peer that completes the TCP handshake and then stalls keeps the actor in connecting for as long as it likes.

Either deadline routes into handleConnectionLost, so the reconnect policy, the outbound buffer and the circuit breaker behave exactly as they do for an observed drop — see BrokerActor base.

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 is chosen through .withFraming(..) with a TcpFraming config union — there are no framer classes to instantiate. Three strategies ship:

type TcpFraming =
| { kind: 'bytes' } // default — raw chunks
| { kind: 'lines'; delimiter?: string; maxLineLen?: number }
| { kind: 'length-prefixed'; maxFrameLen?: number };
  • bytes (default) — every chunk delivered raw; the target handles byte-stream semantics itself.
  • lines — split on delimiter (default '\n'); each frame arrives as a string. maxLineLen caps a line in bytes, the un-terminated one included — that second half is what bounds the re-assembly buffer against a peer that never sends a delimiter. Bytes, not characters: one CJK character costs three of them.
  • length-prefixed — the first 4 bytes (big-endian uint32) carry the payload size; the prefix width is fixed, only maxFrameLen is configurable.

With lines or length-prefixed, each message the target receives is one full frame, not an arbitrary chunk.

A frame past its cap is never buffered: the pending bytes are dropped, and the connection with them. The client destroys its socket and falls back to its reconnect policy; the listener closes only the offending connection.

Both actors validate framing while starting, so a bad setting throws an OptionsError before a socket exists rather than at the first inbound byte. A cap that is not a positive integer is rejected, and so is an empty delimiter — it matches at every offset without consuming anything, which is not a slow framer but a wedged process.

Both actors read the same strategies from the same code — on TcpServerActor the chosen framing is applied per accepted connection, each with its own re-assembly buffer, so two clients mid-frame never splice into one another.

TcpServerActor binds a port and serves every connection that arrives. It is built on the same cross-runtime TCP layer the cluster transport uses, so Bun, Node and Deno — and TLS — come from one place.

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); // required: where events + frames go
const server = system.spawn(() => new TcpServerActor(tcpServerOptions), 'tcp-server');
// Write to one connection, addressed by the id you were handed:
server.tell({ kind: 'send', connectionId, payload: 'PONG\n' });
// Hang up on one connection. The listener keeps serving the rest:
server.tell({ kind: 'close', connectionId });
interface TcpServerOptionsType extends BrokerCommonOptionsType {
bindHost?: string; // default '0.0.0.0'
bindPort?: number; // required; 0 = let the OS pick
framing?: TcpFraming; // per connection; default { kind: 'bytes' }
target?: ActorRef<TcpServerMessage>; // required
tls?: TlsTransportOptionsType; // serve TLS instead of plaintext
maxConnections?: number; // admission cap; default Infinity
}

With bindPort: 0 the OS picks the port — read it back from the actor’s boundPort once it is bound. connectionCount reports the live connections.

The target receives a kind-tagged union, not bare frames — a listener has many connections, so every message names the one it belongs to:

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

The union and each of its variants are importable, so a handler can take the variant it is about rather than narrowing the union at every arm:

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

An echo server is then the obvious three lines:

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 arrives for every ending: the peer hung up, the connection errored, you sent close, a frame breached its size cap, or the actor stopped and unbound the port. One signal, one code path.

tls carries the certificate material — PEM contents or DER bytes, never a path. Supplying ca turns on client-certificate verification (mTLS); set requestClientCert: false for one-way TLS.

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

A half-configured credential — a cert with no key, or a ca alone — is rejected when the actor starts, not quietly bound in plaintext. There is deliberately no HOCON leaf for tls: a config file is the wrong place for a private key.

  • maxConnections caps simultaneously accepted connections. A connection arriving at the cap is aborted immediately instead of being registered — refusing at the door, rather than accepting a socket nobody reads from. Aborted, not ended: an orderly close only sends a FIN, and a peer that never answers it keeps the socket, and its file descriptor, alive. Since that socket was never registered the cap does not count it either, so a half-close would bound only peers that cooperate. The refused peer therefore sees a connection error rather than a clean close — which is what being turned away at capacity looks like at the TCP level.
  • A frame past its framing cap drops only that connection. The client actor answers the same breach by dropping its connection — there the one connection is the whole transport — and then follows its reconnect policy; a listener that did the same would let any single client take the service down for everyone. Either way the pending bytes go with the socket: the point of a cap is not to log the breach but to stop holding what caused it.
  • outboundBuffer defaults to 0 here, unlike every other broker actor. Buffering while disconnected exists so a message survives a reconnect — but “disconnected” for a listener means the port is down, which means every connection it accepted is already gone, so a buffered write names an id that can never come back. A BrokerNotConnected event says that; a replay would not.

Three legitimate uses:

  1. Talking to legacy protocols — proprietary protocols that don’t have higher-level wrappers.
  2. Custom binary protocols — game servers, metric collectors with custom wire formats.
  3. Bridging to non-HTTP services — message queues with proprietary wire (some financial protocols, e.g.).

For new application protocols, don’t reach for raw TCP first. HTTP, WebSocket, or gRPC are better starting points for almost everything.

  • I/O overview — the bigger picture.
  • BrokerActor base — the shared lifecycle.
  • UDP — the connectionless alternative.
  • gRPC — typed RPC over HTTP/2 — usually the better choice than raw TCP.