Hexeract ships two workers in v0.2.0: OutboxWorker for the database-backed outbox and RabbitMqWorker for the AMQP bus. They are intentionally symmetric: same fluent builder shape, same run(cancel) entry point, same handler dispatch via ErasedHandler.
flowchart LR
builder["WorkerBuilder<br/>register_handler::<M, _>(handler)<br/>tuning knobs<br/>build()"]
worker["Worker"]
run["run(cancel: CancellationToken)"]
handlers["HashMap<&'static str, Arc<dyn ErasedHandler>>"]
loop["Polling / consume loop"]
builder --> worker
worker --> run
worker --> handlers
run --> loop
loop --"per envelope"--> dispatch["ErasedHandler::handle"]
dispatch --> typed["TypedHandler<M, H>"]
typed --"decode payload"--> handler["Handler<M>::handle"]
let cancel = CancellationToken::new();
let join = tokio::spawn(worker.run(cancel.clone()));
// ... business code emits events / messages ...
cancel.cancel();
join.await??;Both workers honour the CancellationToken: OutboxWorker checks between poll cycles, RabbitMqWorker selects on consumer.next() and the cancel signal. A cancelled worker drains the in-flight envelope (if any) and returns Ok(()).
OutboxWorker is a poll loop. Two knobs drive its rhythm:
poll_interval(default100 ms): sleep duration when a poll returned no rows.batch_size(default10): maximum rows fetched per poll.
A non-empty poll runs back-to-back without sleeping, so a backlog drains as fast as the handler can process. An empty poll sleeps poll_interval and tries again.
RabbitMqWorker is push-based: it calls basic_consume and reacts to deliveries the broker pushes. Three knobs:
prefetch(default16): how many unacknowledged deliveries the broker may have in flight at once.max_attempts(default5): retry budget permessage_idbefore the delivery is parked or dropped (see retry policy).max_payload_bytes(default1 MiB): cap on the size of a consumed payload, enforced before the payload is copied or deserialized. Broker bytes cross a trust boundary, so an oversize delivery follows the poison path: parked in the dead-letter queue when one is configured, dropped with a warning otherwise (see retry policy). The broker has already buffered the frame when the worker sees it, so the cap bounds the consumer's work, not the network: pair it with the broker-sidemax_message_sizeto bound ingress. Worst-case consumer buffering is roughlyprefetch x max_payload_bytes.
SchedulerWorker is a poll loop over a ScheduleStore, similar in shape to OutboxWorker but with two extra knobs for claim safety and dispatch bounding. See hexeract-scheduler for the full SchedulerBuilder surface.
poll_interval(default100 ms): sleep duration after a cycle that claimed nothing, or that failed outright.batch_size(default10): maximum number of due occurrences claimed per cycle.min_cycle_delay(default5 ms): floor delay between consecutive non-empty cycles, so a backlog still drains quickly without spinning in a tight busy loop.dispatch_timeout(default30 s): hard deadline for a single dispatch to the sink; a sink slower than this is treated as a failed attempt and follows the retry path.
A non-empty cycle waits only min_cycle_delay before the next one, so a backlog drains close to as fast as the sink can absorb it. An empty or failed cycle waits the full poll_interval before retrying.
Each cycle claims occurrences under a lease (default 30 s): the window in which the claiming worker must dispatch and settle the occurrence before another worker is allowed to reclaim it. The builder only rejects a zero lease or a zero dispatch_timeout; it does not enforce a relationship between the two. Operationally, lease should be set with headroom above dispatch_timeout (enough margin to cover the settle step that follows a timed-out dispatch), so a slow sink hits its own timeout and gets retried by the same worker instead of the lease expiring first and letting a second worker claim and dispatch the same occurrence concurrently.
The worker keeps handlers in a HashMap<&'static str, Arc<dyn ErasedHandler>> keyed by MESSAGE_TYPE / EVENT_TYPE. The user-facing trait is the typed Handler<M>; TypedHandler<M, H> is the adapter that translates from the dyn-safe ErasedHandler::handle(&envelope, &ctx) -> BoxFuture<Result<(), BusError>> to the typed H::handle(message, &ctx) -> Result<(), H::Error>.
The decoding step (envelope.decode::<M>()) lives in TypedHandler, so the worker's dispatch loop never needs to know the concrete message type. If the inbound envelope carries a message_type no handler registered for, the dispatch returns BusError::MissingHandler { message_type }, which the worker logs and (in AckMode::Manual) treats as a handler failure subject to the retry policy.
| Worker | Delivery semantics | Idempotency requirement |
|---|---|---|
OutboxWorker |
At-least-once; a crashed worker releases its SELECT ... FOR UPDATE lock and another worker picks the envelope up. |
Required. Same event_id may invoke the handler more than once. |
RabbitMqWorker |
At-least-once; redeliveries on failure, and broker reconnects can replay messages. | Required. Same message_id may invoke the handler more than once. |
Idempotency is not optional. The recommended pattern is to write the side effect plus a processed_message_id row in the same database transaction, then short-circuit on the second delivery when the processed_message_id is already present.
Both workers honour cooperative cancellation:
- Caller flips the
CancellationToken. OutboxWorkerfinishes the current poll cycle (commits the transaction or rolls back) and exits the loop.RabbitMqWorkerlets the in-flight delivery resolve through the handler, sends the ack or nack, then exits the consume loop.runreturnsOk(())and theJoinHandleresolves.
A worker that crashes (panic or unwrap) bubbles the panic to the JoinHandle. Wrap your handler logic in Result::Err mapping rather than panicking.
RabbitMqWorker never reconnects on its own. The contract of run makes the two exit paths unambiguous:
Ok(())is returned for a cooperative cancellation, and nothing else.BusError::Connectionis returned when the consumer stream ends while the token has not fired: the connection or channel is gone and the worker cannot recover it.
Recovery belongs to the caller, because a supervisor that rebuilds everything from scratch has no partial state to reconcile. Each iteration of the loop below reconnects with backoff, re-applies the application topology, and re-creates the worker, which re-declares the wait and dead-letter queues, re-enables publisher confirms and re-subscribes the consumer:
loop {
let connection =
RabbitMqConnection::connect_with_retry(&uri, 5, Duration::from_millis(500)).await?;
ensure_topology(&connection, &exchanges, &queues, &bindings).await?;
let worker = RabbitMqWorkerBuilder::new(connection)
.queue("orders.received")
.dead_letter_routing_key("orders.parked")
.register_handler::<OrderPlaced, _>(MyHandler)
.build()?;
match worker.run(cancel.clone()).await {
Ok(()) => break,
Err(err) => {
tracing::error!(error = %err, "bus worker stopped, restarting");
tokio::time::sleep(Duration::from_secs(1)).await;
}
}
}Redeliveries across reconnects are covered by the at-least-once semantics above: unacked deliveries return to the queue when the connection drops and the retry budget keeps travelling in the x-death header, so a restarted worker resumes the count instead of starting over.
On the publish side, RabbitMqTransport recovers rather than failing fast: the connection enables lapin auto-recovery, so it transparently reconnects on a broker blip and replays the topology on the new connection. A publish issued during the reconnect window is retried once across the recovered channel, so a transient outage does not surface as a publish failure. A publish still surfaces BusError::Connection or BusError::Transport when the broker is durably unreachable or the recovery times out. Wrap publishes with the outbox layer when the application needs an at-least-once guarantee that survives a process crash, not just a broker blip.