BrokerActor base class
Это содержимое пока не доступно на вашем языке.
Every protocol actor in actor-ts/io/broker — KafkaActor,
MqttActor, NatsActor, etc. — extends BrokerActor. The base
class owns the shared lifecycle: connection state machine,
reconnect-with-backoff, outbound buffer, subscriber fan-out,
lifecycle event publishing.
BrokerActor (the abstract base) owns:
- Lifecycle state machine —
disconnected ↔ connecting ↔ connected ↔ disconnecting. - Outbound buffer — messages sent before the connection is up.
- Reconnect loop — exponential backoff on connection loss.
- Subscriber tracking — fan-out for incoming events.
- Desired subscriptions — protocol subscriptions that outlive any one connection and are re-established on every reconnect.
Subclasses implement three protocol hooks:
| Hook | When called |
|---|---|
connectImplementation | Open the protocol-specific connection. |
disconnectImplementation | Close it — on stop and before every re-connect attempt. |
dispatchOutgoing(envelope) | Send a single buffered message on the wire. |
Subclasses implement the three protocol hooks; the base handles the rest. This page documents what’s shared. For per-protocol specifics, see the per-protocol pages.
The state machine
Section titled “The state machine”Four states:
disconnected— initial; not currently connected.connecting—connectImplementationis running.connected— connection up; messages flow.disconnecting—disconnectImplementationis running.
A disconnected → connecting → failure triggers a reconnect
loop: backoff + retry until success or maxAttempts
exhausted.
Every re-connect attempt starts from a clean slate: the base calls
disconnectImplementation first whenever a previous attempt opened
anything, so a subclass never builds a new connection on top of the dead
one’s handles. That makes disconnectImplementation an idempotent
contract — it may be called on an already-dead connection, and it must
drop only the live handles, never the desired subscriptions below.
Stopping the actor while a reconnect attempt is in flight abandons that
connection rather than adopting it. Reconnect runs on the system scheduler,
detached from the mailbox, so an attempt already inside
connectImplementation cannot be cancelled — it resumes after postStop
has returned. The base class therefore re-checks liveness on the way out of
the handshake and, if the actor is gone, calls disconnectImplementation
once more and stays disconnected: no BrokerConnected, no buffer drain,
no further retry. Without that check a handshake that completed installed
live handles on a terminated actor, and one that failed re-armed the
backoff timer and reconnected forever, since maxAttempts defaults to
Infinity (#708).
Subclass contract
Section titled “Subclass contract”abstract class BrokerActor<S, Command, P, Subscription = never> extends Actor<Command> { // Subclasses implement: protected abstract configKey(): string; protected abstract builtInDefaultOptions(): Partial<S>; protected abstract readOptionsFromConfig(config: Config): Partial<S>; protected abstract requiredOptions(): ReadonlyArray<keyof S>; protected abstract endpointLabel(): string;
protected abstract connectImplementation(): Promise<void>; protected abstract disconnectImplementation(): Promise<void>; protected abstract dispatchOutgoing(envelope: OutboundEnvelope<P>): Promise<void>;
protected abstract onCommand(command: Command): void | Promise<void>;
// Optional — only for refs the subclass watches itself: protected onTerminated(signal: Terminated): void | Promise<void> {}}Four categories:
- Settings glue (
configKey,builtInDefaultOptions,readOptionsFromConfig,requiredOptions,endpointLabel) — describes how to assemble settings from the three layers (constructor + HOCON + defaults) and how to validate. - Protocol hooks (
connectImplementation,disconnectImplementation,dispatchOutgoing) — the protocol-specific work. - Message handling (
onCommand) — one command from the mailbox. This is where a subclass’smatch(...).exhaustive()dispatch table lives. - Death watch (
onTerminated) — optional; see Subscriber tracking.
endpointLabel is the human-readable connection identity
(“amqp://localhost:5672”, “kafka-cluster-1”) used in log lines
and lifecycle events.
Return the connection string as configured — do not redact it
yourself. The base class runs every use through
redactedEndpointLabel(), which drops the userinfo and the query
string, so amqp://svc:pw@rabbit:5671/orders reaches
BrokerConnected as amqp://rabbit:5671/orders. That matters
more than it looks: the four lifecycle events go to the
system-wide EventStream, which has no authorization concept, so
any actor in the system can read them — and a reconnect loop
republishes on every backoff tick for the length of an outage.
A joined server list keeps its shape and loses each credential
(nats://***@a:4222,nats://***@b:4222), and a label that never
carried one comes back unchanged. Call redactedEndpointLabel()
in your own log lines too.
onReceive is sealed — the base class owns it, so that it can intercept
Terminated before your dispatch table sees it. Implement onCommand
instead; it receives everything except the death-watch signal.
What the base does for you
Section titled “What the base does for you”Outbound buffering
Section titled “Outbound buffering”this.enqueueOutbound(payload);Subclasses call this.enqueueOutbound(payload) to send; it
returns true if the message was sent or buffered, false if it
was dropped. The base class:
- If connected — wraps
payloadin anOutboundEnvelopeand callsdispatchOutgoing(envelope)immediately. - If disconnected (or
connecting) — buffers up tooutboundBufferenvelopes (default1000). On reconnect, drains the buffer in order.
On overflow the base always evicts the oldest buffered
envelope (FIFO) and publishes a BrokerBufferOverflow event on
the event stream — it is never thrown. With outboundBuffer: 0
buffering is disabled: the message is dropped and a
BrokerNotConnected event is published.
Subscriber tracking
Section titled “Subscriber tracking”import { match } from 'ts-pattern';
protected override onCommand(command: MyCommand): void { match(command) .with({ kind: 'subscribe' }, (c) => this.onSubscribe(c)) .with({ kind: 'unsubscribe' }, (c) => this.onUnsubscribe(c)) .exhaustive();}
private onSubscribe(command: SubscribeCommand): void { this.subscribeRef(command.topic, command.subscriber);}
private onUnsubscribe(command: UnsubscribeCommand): void { this.unsubscribeRef(command.topic, command.subscriber);}subscribeRef(topic, ref) registers ref as interested in
topic’s inbound messages, and death-watches it — for a local ref;
see the caveat below.
unsubscribeRef matches on the ref’s path rather than the ref object, so
holding a different ref for the same actor still unsubscribes.
When the protocol pushes an inbound message, the subclass calls:
this.fanOutToTopic(topic, inboundMessage);The base delivers to every subscriber for that topic.
Desired subscriptions
Section titled “Desired subscriptions”Subscriber tracking above is inside the actor system — which local refs want which topic. A desired subscription is the other half: the subscription the actor holds on the broker, which has to be re-established every time the connection is rebuilt.
The base keeps that set separate from the live handles, so it outlives any one connection:
// Record it (and establish it now, if connected).await this.rememberSubscription('orders.new', target);// Drop it, from the broker and from the desired set.await this.forgetSubscription('orders.new');Subclasses provide three hooks:
| Hook | Purpose |
|---|---|
initialSubscriptions() | Subscriptions declared in the options. Folded into the desired set once, before the first connect. |
applySubscription(key, subscription) | Establish one subscription on the live connection. Must tolerate a key that is already live. |
revokeSubscription(key) | Tear one down. Optional — several protocols can only drop a subscription by dropping the whole consumer. |
and replay the whole set from inside connectImplementation:
protected async connectImplementation(): Promise<void> { this.connection = await MyClient.connect(this.options.url); await this.applyDesiredSubscriptions(); // configured + runtime}The replay is driven by the subclass rather than by the base class
because the right point in the handshake is protocol-specific —
kafkajs, for one, wants every subscribe in before consumer.run.
What this buys you:
- A subscription added at runtime survives a reconnect — it is not just a call on a connection that is about to die.
- A
subscribethat arrives while the actor is disconnected is remembered and applied on the next connect, instead of being dropped. - Seeding from the options is once-only, so a runtime
unsubscribeis not resurrected by the next reconnect. - Re-remembering a live key revokes it first, so a changed payload (a different target actor, say) actually takes effect.
- A subscription that cannot be established is logged as a warning and the rest of the set still goes through. One bad subject does not take the connection down — and it does not silently leave you connected-but-deaf either.
MqttActor predates this mechanism and keeps its own richer registry
(per-topic QoS, several targets per pattern, deathwatch on each), with
the same reconnect guarantees.
Reconnect-with-backoff
Section titled “Reconnect-with-backoff”reconnect: { initialDelayMs: 200, maxDelayMs: 30_000, factor: 2, maxAttempts: Infinity, // the default — retry forever randomFactor: 0.2, // the default — ±20 % jitter}Configurable per actor (or reconnect: false to disable
auto-reconnect entirely). Each attempt waits
min(initialDelayMs * factor^(attempt - 1), maxDelayMs), multiplied
by a random 1 ± randomFactor.
The jitter is what keeps a fleet from retrying in lockstep. Without
it the delay is a pure function of the attempt counter and these options,
so every broker actor that lost the same broker in the same instant wakes
in the same millisecond, wave after wave — and the herd can hold the
recovering broker down. Set randomFactor: 0 for a fully deterministic
schedule.
The same spread applies to the circuit breaker: when the breaker is
open the actor re-checks somewhere in
[resetMs, resetMs × (1 + randomFactor)] rather than exactly on the
deadline. That jitter is one-sided — a breaker is never re-tried early.
Each attempt fires BrokerReconnectAttempt on the event stream (its
delayMs is the jittered value actually scheduled); after maxAttempts
exhausted (if finite), BrokerReconnectFailed fires and the actor stays
disconnected.
For a deterministic test, inject the randomness source — random: () => 0.5 pins every delay to its un-jittered value:
reconnect: { initialDelayMs: 100, random: () => 0.5 }Liveness deadlines
Section titled “Liveness deadlines”Everything above is triggered by handleConnectionLost, and that is
only ever called when the transport says something ended — a
close event, an error, a stream that reported done. A peer that
vanishes without sending FIN or RST says none of them. A dropped NAT
entry, a container killed with SIGKILL, a route that starts
black-holing: the socket stays open, the actor stays connected, and
the reconnect machinery never runs.
Two clocks close that gap. Both are opt-in per protocol, because only the application knows how quiet its peer is allowed to be:
idleTimeoutMs— a read deadline. If nothing arrives from the peer for that long, the connection is reported lost through the ordinaryhandleConnectionLostpath, so everything on this page applies unchanged:BrokerDisconnected, backoff, buffer, breaker. It is reset by inbound bytes and not by outbound ones, deliberately: a client writing into a black hole is exactly the case that has to trip it.connectTimeoutMs— a deadline on oneconnectImplementation. Every protocol settles its connect on an event (connect,open, response headers) and none of them has a clock, so a peer that finishes the handshake and then stalls holds the actor inconnectingindefinitely. On expiry the attempt is aborted, and the failure travels the path a refused connect already takes.
TcpSocketActor, SseActor and WebsocketClientActor expose both.
Set idleTimeoutMs above whatever heartbeat the peer already
sends; below it, the deadline severs healthy connections in a loop.
0 (the default) disables either one.
Protocols with their own protocol-level liveness do not need these:
MqttActor has MQTT keepAlive, and KafkaActor runs a consumer
heartbeat.
Lifecycle events
Section titled “Lifecycle events”Published on system.eventStream:
| Event | When |
|---|---|
BrokerConnected | A connectImplementation succeeded. |
BrokerDisconnected | A live connection was lost. Not published by a graceful stop. |
BrokerReconnectAttempt | A reconnect attempt is starting. |
BrokerReconnectFailed | maxAttempts exhausted. |
BrokerBufferOverflow | The outbound buffer dropped an envelope. |
BrokerNotConnected | Sent without a connection. |
Subscribe to monitor every broker actor uniformly:
system.eventStream.subscribe(monitorRef, BrokerConnected);system.eventStream.subscribe(monitorRef, BrokerDisconnected);The events include actorPath — distinguish events from
different broker actors in the system.
Settings resolution
Section titled “Settings resolution”1. builtInDefaultOptions() ← lowest priority (always applied)2. readOptionsFromConfig() ← HOCON overrides3. Constructor argument ← highest priority (per-instance)preStart merges the three layers, validates against
requiredOptions(), and stashes the result for the rest of the
actor’s life via this.options.
Missing required settings cause an early-error throw on
preStart — the actor goes through the supervisor’s failure
path before it ever attempts to connect.
Writing a custom protocol actor
Section titled “Writing a custom protocol actor”import { match } from 'ts-pattern';import { BrokerActor, type OutboundEnvelope, type BrokerCommonOptionsType } from 'actor-ts/io';
interface MyProtocolOptionsType extends BrokerCommonOptionsType { readonly url: string;}
class MyProtocolActor extends BrokerActor<MyProtocolOptionsType, Command, MyPayload> { private connection: MyClient | null = null;
protected configKey() { return 'actor-ts.io.broker.my-protocol'; } protected builtInDefaultOptions() { return { /* ... */ }; } protected readOptionsFromConfig(c) { /* parse HOCON */ return {}; } protected requiredOptions() { return ['url'] as const; } protected endpointLabel() { return this.options.url; }
protected async connectImplementation(): Promise<void> { this.connection = await MyClient.connect(this.options.url); this.connection.onMessage((m) => this.fanOutToTopic(m.topic, m)); }
// Called on stop AND before every re-connect attempt — idempotent, // and safe on a connection that is already dead. protected async disconnectImplementation(): Promise<void> { await this.connection?.close(); this.connection = null; }
protected async dispatchOutgoing(env: OutboundEnvelope<MyPayload>): Promise<void> { await this.connection!.send(env.payload); }
// `onReceive` is sealed by the base — implement `onCommand` instead. protected override onCommand(command: Command): void { match(command) .with({ kind: 'subscribe' }, (c) => this.onSubscribe(c)) .exhaustive(); }
private onSubscribe(command: SubscribeCommand): void { this.subscribeRef(command.topic, command.subscriber); }}The base handles the rest. Most third-party clients (kafkajs,
nats.js, etc.) have an event-based message-receive API that
maps cleanly to this.fanOutToTopic(...).
Where to next
Section titled “Where to next”- I/O overview — the bigger picture: which protocols ship.
- Kafka / MQTT / NATS / etc. — per-protocol pages.
- Event stream — where lifecycle events are published.
- Backoff policy — a standalone backoff primitive (the reconnect loop rolls its own).
