An async messaging framework for Rust: broker-agnostic traits, a router runtime, codecs, AsyncAPI generation, Prometheus and OpenTelemetry observability, and a conformance harness for broker authors.
RustStream connects your service to a message broker through a small set of generic traits, then gives you a router, middleware, codecs, and tooling on top. The core depends on no broker: each broker is an independent crate held to one contract. The core is 100% safe Rust.
- Broker-agnostic core. Brokers are separate crates, checked by a conformance harness.
- Misuse does not compile. Double acks, out-of-order lifecycle calls and under-specified publishes are compile errors that name the fix.
- Pluggable codecs: JSON, MessagePack, and CBOR behind cargo features, or raw bytes with none.
- Zero-boilerplate binaries.
#[subscriber]and#[ruststream::app]macros, and a CLI that scaffolds, runs and documents a service. - Observability: AsyncAPI 3.1, Prometheus metrics, a health probe, and OpenTelemetry.
- Tests without a broker. The service's own app runs in process under a test harness.
- Capability traits for batches, transactions, request-reply, partitioning and repositioning; a broker implements only what it supports.
[dependencies]
ruststream = { version = "0.7", features = ["macros", "memory", "json"] }
serde = { version = "1", features = ["derive"] }The CLI ships with the crate behind the cli feature:
cargo install ruststream --features cliuse ruststream::memory::prelude::*;
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct Order {
id: u64,
}
#[subscriber("orders")]
async fn handle(order: &Order) -> HandlerOutcome {
println!("got order {}", order.id);
HandlerOutcome::ack()
}
#[ruststream::app]
fn app() -> RustStream {
RustStream::new(AppInfo::new("orders", "0.1.0"))
.with_broker(MemoryBroker::new(), |b| { b.include(handle); })
}#[ruststream::app] generates main, so there is no runtime boilerplate.
ruststream run # start the service (or: cargo run -- run)
ruststream asyncapi gen # print the AsyncAPI documentScaffold a fresh project with cargo generate --git https://github.com/powersemmi/ruststream templates/memory --name my-service (each broker crate ships its own template). See the
quick start.
TestApp runs the service's own app in process, with no external broker.
use ruststream::testing::TestApp;
let tb = TestApp::start(service()).await?;
// Inject an order; the harness drives the handler to completion before returning.
tb.broker::<MemoryBroker>()
.message(&Order { id: 42 })
.to("orders")
.publish()
.await?;
// The handler ran once, decoded the order, and acked.
tb.broker::<MemoryBroker>()
.subscriber("orders")
.assert_called_once()
.with(&Order { id: 42 })
.settled(HandlerOutcome::ack());Full compiling example: examples/testing.rs.
| Broker | Crate |
|---|---|
| NATS | ruststream-nats |
| Redis / Valkey | ruststream-fred |
| RabbitMQ (AMQP 0.9.1) | ruststream-lapin |
| Apache Kafka | ruststream-rdkafka |
| AMQP 1.0 | ruststream-amqp |
| Google Cloud Pub/Sub | ruststream-gcp-pubsub |
| Amazon SQS / SNS | ruststream-sqs-sns |
| Apache Pulsar | ruststream-pulsar |
| MQTT 5 | ruststream-rumqttc |
| ZeroMQ | ruststream-zeromq |
| Files and stdio | ruststream-sea-file |
| Amazon Kinesis | ruststream-kinesis |
What each broker supports is on the broker index. To write a broker, see the broker-authors guide.
- Site: https://powersemmi.github.io/ruststream/latest
- API reference and topic overviews: https://docs.rs/ruststream
The MSRV is 1.95, edition 2024. The policy is a rolling one, as in tower: the MSRV rises only to a Rust release at least six months old, and only in a minor release. A broker crate may require a newer toolchain when its client does.
See CONTRIBUTING.md.
Licensed under the Apache-2.0 license.
Inspired by FastStream.