Brokers¶
A broker connects RustStream to a message transport. The framework ships an in-memory broker for development and tests; production brokers are separate crates you add as a dependency.
Handlers, routers, codecs, and middleware are broker-agnostic, so moving between brokers is a
one-line change at with_broker.
Each broker crate is documented where the Docs column points: on its own site, or in its repository until a site exists.
| Broker | Crate | Transport | Docs |
|---|---|---|---|
| Memory | ruststream (feature memory) |
in-process, for development and tests | this site |
| NATS | ruststream-nats |
Core NATS and JetStream | powersemmi.github.io/ruststream-nats |
| Redis | ruststream-fred |
Redis Streams (standalone, cluster, sentinel) | powersemmi.github.io/ruststream-fred |
| RabbitMQ | ruststream-lapin |
AMQP 0.9.1 (queues, exchanges, publisher confirms, direct reply-to) | powersemmi.github.io/ruststream-lapin |
| Kafka | ruststream-rdkafka |
Apache Kafka (consumer groups, tracked commits, transactions, exactly-once pipelines) | powersemmi.github.io/ruststream-rdkafka |
| AMQP 1.0 | ruststream-amqp |
ActiveMQ Artemis, RabbitMQ 4.x, Azure Service Bus, and the wider AMQP 1.0 family (request/reply, transactions) | repository |
| Google Cloud Pub/Sub | ruststream-gcp-pubsub |
Pub/Sub over the official client (ordering keys, exactly-once acknowledgement, dead-letter policies) | repository |
| AWS SQS / SNS | ruststream-sqs-sns |
SQS queues with SNS fan-out (FIFO groups, visibility management, native deferred retry) | repository |
| Apache Pulsar | ruststream-pulsar |
Pulsar topics and patterns (subscription modes, dead-letter policies, repositioning) | repository |
| MQTT 5 | ruststream-rumqttc |
MQTT v5 (QoS levels, shared groups, retained messages) | repository |
| ZeroMQ | ruststream-zeromq |
Brokerless PUSH/PULL, PUB/SUB, and DEALER/ROUTER request/reply over TCP and IPC | repository |
| Stream files / stdio | ruststream-sea-file |
Persistent replayable stream files and shell pipelines; zero infrastructure, full repositioning | repository |
| AWS Kinesis | ruststream-kinesis |
Kinesis data streams (shard leasing, checkpointing, repositioning) | repository |
To implement a broker for another transport, see Broker authors.
Switching brokers¶
Every broker constructs synchronously and connects lazily (the runtime calls Broker::connect once
at startup), so the same handlers and routers run on any of them; only the broker construction
differs by one line inside with_broker.
use ruststream::runtime::{AppInfo, RustStream};
use ruststream_rdkafka::KafkaBroker;
#[ruststream::app]
fn app() -> RustStream {
RustStream::new(AppInfo::new("orders", "0.1.0"))
.with_broker(
KafkaBroker::new(["localhost:9092"]).default_group("orders"),
|b| {
b.include_router(routes::orders())
},
)
}
Each broker crate documents its own Config and connection options. Subscriptions that need
broker-specific options (consumer groups, durable names) use that broker's descriptor in the
#[subscriber(..)] decorator; see
broker-specific descriptors.