Writing a broker¶
A broker is an independent crate that implements the core traits. It depends on ruststream with
default features off, so it pulls in the trait surface and runtime without the bundled JSON codec or
any other broker:
This page is the contract. Implement the required traits, expose your own Config, add capability
traits for the features your broker supports, and prove the result with the
conformance harness. For a complete implementation on a real client, see the
worked NATS example.
The required traits¶
Broker and ConnectedBroker¶
The broker is pure lifecycle, and the lifecycle is a ladder of consuming transitions: each state is a distinct type, so out-of-order calls do not compile. The broker carries no subscriber or publisher type, so a single application can mix broker kinds.
pub trait Broker: Send + Sync + Sized {
type Error: std::error::Error + Send + Sync + 'static;
type Connected: ConnectedBroker;
async fn connect(self) -> Result<Self::Connected, Self::Error>;
}
pub trait ConnectedBroker: Send + Sync + Sized + 'static {
type Error: std::error::Error + Send + Sync + 'static;
type Closed: Send;
async fn shutdown(self) -> Result<Self::Closed, Self::Error>;
}
shutdown must never block or panic; do all fallible teardown here and return a Result. The
Closed witness has no publish or subscribe surface; carry teardown diagnostics (flush results,
drop counts) in it as plain data, or use ().
Construction is synchronous and I/O-free: new(addrs) only records configuration, all network
work happens in connect (called once at startup by the runtime), and the connected form holds
the live client directly - its own operations never check a "maybe connected" state. A broker may
additionally keep a shared cell that connect fills (or a shareable in-process state, as the
in-memory broker does) so publishers can be handed out while the app is still being assembled,
before connect runs; the cell serves those early handles, not the connected form. The
NATS example shows the cell-backed variant. This contract is what lets a
service compose with the synchronous #[ruststream::app] builder; the
conformance harness proves it end to end.
Consuming transitions make owner-side misuse unrepresentable: there is no publish or subscribe to call on a broker you already shut down. What remains a runtime rule is transport reality, not contract bookkeeping: handles that alias the connection (publishers handed out from the connected form, clones of a shareable broker) must surface an error when used after shutdown - never a silent success against a dead connection. The lifecycle check drives that path too.
Subscribe¶
Implement Subscribe on the connected form to support subscribing by name. This is what
#[subscriber("name")] uses.
pub trait Subscribe: ConnectedBroker {
type Subscriber: Subscriber;
async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error>;
}
Subscriber¶
A subscriber is a Stream of incoming messages. Back-pressure comes for free from the stream.
pub trait Subscriber: Send {
type Message: IncomingMessage;
type Error: std::error::Error + Send + Sync + 'static;
fn stream(&mut self) -> impl Stream<Item = Result<Self::Message, Self::Error>> + Send + '_;
}
stream takes &mut self, so any state buffered between polls lives behind the mutable borrow,
which keeps it cancel-safe.
IncomingMessage¶
A delivered message exposes its payload and headers, and is acked or nacked. Ack consumes self, so
double-ack is a compile error.
pub trait IncomingMessage: Send + Sync {
fn payload(&self) -> &[u8];
fn headers(&self) -> &Headers;
async fn ack(self) -> Result<(), AckError>;
async fn nack(self, requeue: bool) -> Result<(), AckError>;
// Defaulted: a plain nack(true). Override when the transport has native
// delayed redelivery (JetStream NAK with delay); handlers reach it through
// HandlerResult::retry_after.
async fn nack_after(self, delay: Duration) -> Result<(), AckError>;
// Defaulted: None. Override (with the Partitioned capability) to feed the
// runtime's keyed worker lanes, workers(n, by_key).
fn partition_key(&self) -> Option<&[u8]>;
}
The two defaulted methods are how optional broker behaviour degrades gracefully: a broker that
overrides neither still works with every runtime feature, with retry_after falling back to an
immediate requeue and keyed lanes rotating keyless messages.
Publisher¶
pub trait Publisher: Send + Sync {
type Error: std::error::Error + Send + Sync + 'static;
async fn publish(&self, msg: OutgoingMessage<'_>) -> Result<(), Self::Error>;
}
OutgoingMessage borrows its name and payload, so publishing does not force an allocation.
PublishPolicy¶
A broker publisher is a bundle of policy (an exchange, a queue timeout, a transactional id) plus
the live connection. Split it along that seam: ship a freely constructible policy type with
the builder options and no publish surface, and implement PublishPolicy to pair it with the
connected form into the live publisher. Pairing is async and fallible for brokers that do real
work when a publisher comes alive (initializing a transactional producer); for most it is a cheap
constructor call.
pub trait PublishPolicy<C: ConnectedBroker> {
type Live; // the live publisher (or live wiring form, for combinator stacks)
async fn pair(self, connected: &C) -> Result<Self::Live, PairError>;
}
The error is the type-erased PairError (wrap your broker's failure with PairError::new)
rather than C::Error: pairing runs once per publisher at startup, never on the hot path, and a
cross-broker token pairs against a broker other than the including scope's, so a broker-typed
error could not name one broker anyway.
Ship one policy/live pair per genuine publishing mode, and make mode selection a policy type
transition rather than a runtime flag: a plain policy pairs into the plain publisher, and a
transactional_id(..) builder step moves to a distinct transactional policy type whose live form
implements TransactionalPublisher - so the plain publisher has no transactional surface at all.
The in-memory broker's MemoryPublish / MemoryRequest are the minimal reference (no options, so
they are unit markers); the core's typed combinators implement PublishPolicy functorially, which
is what lets users compose codecs and transforms over your policy before it pairs.
When the plain policy is usable with its defaults (most are), also implement DefaultPublish on
the connected form to name it. That is what lets the runtime build the default reply publisher
when a publish("dest") handler is included without an explicit .publisher(..) - b.include(def)
alone compiles. Brokers whose publishers always need explicit options simply do not implement it,
and their users attach a policy at every registration.
pub trait DefaultPublish: ConnectedBroker {
type Policy: PublishPolicy<Self> + Default + Send + 'static;
}
Subscription sources¶
Subscribe covers the by-name case. When a subscription needs broker-specific options (a consumer
group, a durable name, a delivery policy), expose a descriptor type that implements
SubscriptionSource:
pub trait SubscriptionSource<C: ConnectedBroker> {
type Subscriber: Subscriber;
fn name(&self) -> &str;
fn subscribe(self, connected: &C) -> impl Future<Output = Result<Self::Subscriber, C::Error>> + Send;
}
Give the descriptor an associated constructor (OrdersStream::new(..)) rather than a free function,
so users can name it directly in the decorator: #[subscriber(OrdersStream::new("orders", "workers"))].
The macro reads the type out of the constructor call, and also accepts a builder chain on it
(#[subscriber(OrdersStream::new("orders").durable("workers"))]) as long as each method returns
Self. Because type Subscriber lives on the source, one broker can offer several subscription
kinds (pub/sub versus streams) with different subscriber types - or, as the
NATS example does, serve them all from one descriptor that branches internally.
Capability traits¶
Implement only the capabilities your broker supports; none are part of the mandatory interface.
| Trait | For brokers that support |
|---|---|
BatchSubscriber |
receiving messages in batches |
TransactionalPublisher |
begin / commit / abort around publishes on the handle |
OwnedTransactions / Transaction |
transactions whose buffer lives in a value, any number open at once per handle |
RequestReply |
native request-reply (NATS yes, Kafka no) |
Partitioned |
a partition key on outgoing messages |
Seekable / Seeker |
repositioning a live subscription in a replayable log |
Positioned |
deliveries that report their own log position |
DescribeServer |
reporting a ServerSpec for AsyncAPI |
Seekable mints its Seeker handle before the stream borrows the subscriber, so a running
subscription can be repositioned from outside the dispatch loop. Positions are broker-owned
(KafkaPosition-style constructors on your own type); a position captured from a delivered
message via Positioned::position carries a pinned contract - seeking to it redelivers exactly
that message - while constructed positions keep the semantics your position type documents.
Document what one seek covers (a consumer instance, or a shared group cursor) and reset any ack
bookkeeping the reposition invalidates.
Per-delivery context and Ctx keys¶
A broker with native delivery metadata (a partition, an offset, a stream sequence) exposes it as a
typed per-delivery context: a #[non_exhaustive] struct the subscriber names, plus ContextField
key types so handlers can bind single fields as parameters with the
Ctx<K> extractor. Keys are unit structs; values are
owned. No type-map and no heap on the delivery path.
/// Per-delivery context of this broker.
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct MyContext {
pub partition: i32,
}
/// `Ctx<Partition>` in a handler binds the delivery's partition.
#[derive(Debug, Default, Clone, Copy)]
pub struct Partition;
impl ContextField for Partition {
type Context = MyContext;
type Value = i32;
fn read(self, src: &MyContext) -> i32 {
src.partition
}
}
A broker with no per-delivery fields uses () and skips all of this.
Middleware on the async edges¶
Integrations that need async I/O around encode and decode (a schema registry, a wire-format
envelope) do not belong in a Codec: the core codec is synchronous and handlers should stay on the
default one. Put them on the async edges instead - transcode incoming payloads on the
subscription's delivery path (before the codec sees them), and frame outgoing ones with a core
PublishLayer added app-wide via RustStream::publish_layer. The publish layer is async and
fallible, and Outgoing::payload_mut exists exactly for envelope wrapping.
Config and defaults¶
Your crate owns its Config. The core carries no broker-specific config, which is what keeps an
upstream change scoped to one broker crate. If a config field has no sane default, do not implement
Default for it; force the user to set it explicitly rather than shipping a default that might break
later.
Errors¶
Use thiserror for a single crate-level error enum, with variants by source. Mark public error
enums #[non_exhaustive]. Never use anyhow in a library crate.
Test support¶
Ship an in-process transport implementing TestableBroker on its connected form under a
testing feature (registered with register_testable_broker! for that connected type, since the
harness connects every broker before recovering its transport) so users can unit-test handlers
against your broker with the TestApp harness. The transport does core routing only: it dispatches published messages to matching
subscribers and treats ack/nack as effectively a no-op. Do not simulate broker-specific semantics
(durable cursors, redelivery timers, offsets, dead-letter routing) in it; those are verified end to
end against a real server.
The reference is the in-memory broker's own implementation (on ConnectedMemoryBroker):
// The harness drives the connected form: TestApp connects every registered broker before it
// recovers the in-process transport, and run_suite scenarios receive connected brokers.
#[cfg(feature = "testing")]
impl crate::testing::TestableBroker for ConnectedMemoryBroker {
fn install_coordinator(&self, coordinator: Coordinator) {
self.state.install_coordinator(coordinator);
}
fn inject(&self, message: OutgoingMessage<'_>) {
// Injecting into a shut-down bus is a harness bug (both run_suite and TestApp drive
// the bus strictly before shutdown), so fail loudly instead of losing the message.
self.state
.fanout(&MemoryOutbound {
name: message.name().to_owned(),
payload: Bytes::copy_from_slice(message.payload()),
headers: message.headers().clone(),
})
.expect("inject on a shut-down broker: drive the harness before shutdown");
}
fn published(&self, name: &str) -> Vec<RawMessage> {
self.state
.published
.lock()
.expect("memory broker mutex poisoned")
.get(name)
.cloned()
.unwrap_or_default()
}
}
#[cfg(feature = "testing")]
crate::register_testable_broker!(ConnectedMemoryBroker);
The transport calls Coordinator::enqueued on every enqueue into a subscriber and
Coordinator::consumed when a delivery is settled or dropped (so the harness can tell when the
reaction has settled), and routes delayed redeliveries through Coordinator::schedule_redelivery.
That one type then works with both TestApp and the conformance suite. See
Testing for the user-facing side, and Conformance to
prove the implementation with run_suite and the lifecycle ladder check.