ProducerController
Это содержимое пока не доступно на вашем языке.
ProducerController wraps a sender’s outgoing messages,
assigning sequence numbers + holding them until acked by
the consumer.
import { ProducerController, ProducerControllerOptions } from 'actor-ts/delivery';
const producerControllerOptions = ProducerControllerOptions.create<OrderEvent>() .withProducerId('orders') // stable identity .withConsumer(consumerRef) .withWindowSize(16) // flow-control window .withResendTimeout(500);const producer = system.spawn( () => new ProducerController<OrderEvent>(producerControllerOptions), 'producer',);
producer.tell({ kind: 'reliable-delivery.send', body: { orderId: 'o-1' } });Internally:
- Assigns seq 1 to the first message, 2 to the next, …
- Sends each via wrapped tell to the consumer.
- Holds in a buffer until acked.
- Retransmits if no ack within
resendTimeout.
Configuration
Section titled “Configuration”type ProducerControllerOptionsType<T> = { consumer: ActorRef<Delivery<T>>; resendTimeout?: number; // default 500 windowSize?: number; // default 16 producerId?: string; // randomly generated if omitted};| Field | Purpose |
|---|---|
consumer | The ConsumerController ref to deliver to. |
resendTimeout | How long to wait for an ack before retransmitting. |
windowSize | Flow-control window — pauses queueing after N unacked messages in flight. |
producerId | Stable identifier the consumer keys its dedup state on. Randomly generated if omitted; pin it when you want one identity across restarts in logs, metrics and the consumer’s map. It does not carry a dedup window across a restart — see Restart and the producer incarnation. |
A generated producerId is random, not sequential. It used to be a module
counter — producer-1, producer-2, … — which was wrong twice over. An
Acknowledgment names a producerId and a seq, and the seq is a small
integer by construction, so a counter left neither half of that pair to guess.
The counter was also shared: two processes running the same service each minted
producer-1, stamped one identity onto their deliveries, and then kept
resetting each other’s dedup window at a consumer they had in common. The
framework now draws 16 hex characters instead.
That id is minted fresh on every construction, so leaving producerId unset
means there is no identity to be stable about — set it yourself whenever
something downstream is supposed to recognise this producer after a restart.
Backpressure
Section titled “Backpressure”windowSize: 16The producer holds up to 16 unacked messages in flight. Beyond that,
incoming { kind: 'reliable-delivery.send' } messages are queued — the
caller’s tell succeeds, but actual sending pauses until the
consumer’s acks free up the buffer.
There is no signal back to the caller when this queueing happens: a
reliable-delivery.send is fire-and-forget — the producer never replies
to it, so there is no ask-based “I sent it” acknowledgment to await.
For senders that need to know a message got through, pass a confirm
callback on the send. It fires once the consumer acks that message —
with null on success, or an Error if the producer stops while the
message is still queued:
producer.tell({ kind: 'reliable-delivery.send', body: { orderId: 'o-1' }, confirm: (err) => { if (err) console.error('delivery failed', err); else console.log('delivered and acked'); },});To turn that into real backpressure, pace your sends against these confirmations — hold off on the next batch until earlier messages confirm.
Retransmission
Section titled “Retransmission”When an ack doesn’t arrive within resendTimeout:
seq 5 sent at t=0seq 5 ack never arrivesat t=500 → retransmit seq 5at t=1000 → retransmit seq 5... until ack arrives or producer stopsRetransmits use the same seq as the original, and carry the same producer incarnation, so the consumer recognises them as the same message.
Restart and the producer incarnation
Section titled “Restart and the producer incarnation”Every construction of a ProducerController mints a fresh, unguessable
incarnation token and stamps it onto each Delivery. Nothing configures
it; it exists because producerId alone cannot answer two questions.
Is this a retransmit or a restarted producer? The sequence counter is
in-memory, so a restarted producer numbers from 1 again while producerId
stays as configured. A consumer keyed on producerId alone would see the
whole post-restart prefix as sequence numbers it had already handled, absorb
them as duplicates, and — because an absorbed delivery is still acked — report
each one back as successfully delivered. The consumer keys on (producerId, incarnation) instead, so a new incarnation starts a fresh dedup window and
its messages reach the handler.
Did this acknowledgment come from a real delivery? producerId and seq
are both guessable, so a three-field ack could be manufactured by anything able
to address the producer, cancelling its retransmit and firing your confirm
with success. The incarnation travels only on deliveries the producer actually
sent, so echoing it back is the evidence the producer requires.
An ack that names the wrong incarnation — forged, or a straggler from the
previous incarnation of the same producerId — is ignored.
// Without persistence:producer crashes → buffer lost → unacked messages are gone// On restart, send resumes from seq 1 under a new incarnation, so the// consumer treats those sends as new rather than as duplicates
// With ProducerController + PersistentActor wrapper:producer crashes → recovery rebuilds buffer + last-acked-seq// On restart, send resumes from the right seqFor full durability, wrap the producer in a PersistentActor
pattern — persist outgoing messages before sending; on
restart, replay any unacked.
The framework’s plain ProducerController is in-memory —
sufficient for ephemeral streams. For durable streams, layer
persistence on top.
Multiple consumers
Section titled “Multiple consumers”// One producer per consumer:const p1 = system.spawnAnonymous(() => new ProducerController({ consumer: c1Ref, ... }));const p2 = system.spawnAnonymous(() => new ProducerController({ consumer: c2Ref, ... }));Each ProducerController is 1:1 with one ConsumerController. For fan-out, multiple producers send to different consumers.
For routing-based fan-out (one logical stream split across N consumers), you’d write a custom router on top.
Cluster-aware
Section titled “Cluster-aware”ProducerController’s consumer ref can point at a remote
consumer (different cluster node). The cluster transport
serializes envelopes; retransmissions work the same.
The producer holds the buffer locally — on producer-host crash, those messages are lost unless persisted.
Comparison with raw tell
Section titled “Comparison with raw tell”ref.tell(msg): ProducerController:- No seq, no ack, no retransmit - Reliable delivery contract- Lost on dead recipient - Survives transient failures- Sub-microsecond cost - ~50µs per message overheadUse the controller only for streams where loss is unacceptable;
raw tell for everything else.
Where to next
Section titled “Where to next”- Delivery overview — the bigger picture.
- ConsumerController — the receiver side.
- Ack semantics — when acks fire and what they mean.
- PersistentActor — for durable producer state.
