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, so each broker is an independent crate held to one contract; broker-specific configuration never leaks into the framework.
The core is 100% safe Rust: every crate carries #![forbid(unsafe_code)] and CI rejects any unsafe
block, so the guarantee cannot regress.
Features
- Broker-agnostic core. Just traits and types, zero broker dependencies. Brokers are separate crates, and the contract is checked by a conformance harness.
- Fully async on tokio. No blocking APIs in the public surface.
- Subscribers are
Streams, not callbacks. Back-pressure comes for free. - Misuse does not compile. Ack consumes
self(no double-ack); the broker lifecycle is a ladder of consuming transitions (connect(self)yields the connected form,shutdown(self)a terminal witness), so out-of-order lifecycle calls are compile errors; transactions settle by consuming their scope. - Publishers pair at startup. Reply wiring and the
Out(..)handler parameter attach a publish policy where the handler is included; the runtime pairs it against the connected broker, so a handler never sees a "not connected" publisher. - Pluggable codecs: JSON, MessagePack, and CBOR behind cargo features - or none at all:
rawsubscribers andpublish_rawreplies move payload bytes untouched. - Zero-boilerplate binaries.
#[ruststream::app]generatesmain; theruststreamCLI scaffolds projects, runs them, and generates the AsyncAPI document. - AsyncAPI 3.0 and Prometheus metrics, served from your own HTTP stack.
- OpenTelemetry behind the
otelfeature: OTLP export for traces and metrics, per-handler dispatch metrics following the messaging semantic conventions, and W3C trace-context propagation across the consume-transform-produce chain. - A cloneable health probe off the running app, so a sibling healthz route keeps reporting the terminal state after the messaging side stops.
- Colored console logging behind the
loggingfeature; the generated CLI installs it onrun, with verbosity driven byRUST_LOG. - Capability traits for optional features (batch subscribe, borrowed and owned transactions, request-reply, partitioning, repositioning a live subscription in a replayable log); a broker implements only what it supports.
Install
[]
= { = "0.6", = ["macros", "memory", "json"] }
= { = "1", = ["derive"] }
= "1"
The CLI ships with the crate behind the cli feature:
Write a service
use MemoryBroker;
use ;
use subscriber;
use JsonSchema;
use Deserialize;
async
#[ruststream::app] generates main, so there is no runtime boilerplate.
Injecting dependencies
Declare app state, derive FromRef, and take a dependency as a State<T> handler argument instead of
reaching through ctx.state(). The state is built once in on_startup; #[derive(FromRef)] makes
each field injectable, so no extractor is written by hand.
use State;
use FromRef;
async
Full compiling example: examples/from_context.rs.
Run it
Scaffold 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.
Testing the service
Unit-test a built service against the in-memory broker, with no external service. MemoryBroker is a
real broker here, not a test double: the TestApp harness drives it through the same dispatch path
the production runtime uses, so you assert on handler behaviour, middleware, and decoding exactly as
in production.
use TestApp;
let tb = start.await?;
// Inject an order; the harness drives the handler to completion before returning.
tb.
.publish
.await?;
// The handler ran once, decoded the order, and acked.
tb.
.subscriber
.assert_called_once
.with
.settled;
// It published the matching receipt downstream.
tb.
.
.assert_called_once
.with;
Full compiling example: examples/testing.rs. See the
testing guide.
Project documentation
Build the AsyncAPI spec and the interactive viewer HTML programmatically from a built service, then
serve them from your own HTTP stack. The CLI ruststream asyncapi gen (see Run it above) prints the
same document to stdout; this is the in-process path.
use ;
let spec = build_spec.to_json?;
let viewer = render_viewer_html;
// serve `viewer` at `/` and `spec` at `/asyncapi.json` from your own HTTP stack
Full compiling example: examples/asyncapi_http.rs.
- Guide and tutorials: https://powersemmi.github.io/ruststream/latest
- API reference: https://docs.rs/ruststream
- Writing a broker: https://powersemmi.github.io/ruststream/latest/broker-authors/
Ecosystem
ruststream-nats: the NATS broker (Core NATS and JetStream).ruststream-fred: the Redis broker (Redis Streams with consumer groups; standalone, cluster, and sentinel topologies) via thefredclient.ruststream-lapin: the RabbitMQ broker (AMQP 0.9.1: topology descriptors, native dead-letter and delayed retry, keyed worker lanes, publisher confirms and server-side transactions) via thelapinclient.ruststream-rdkafka: the Apache Kafka broker (consumer groups, tracked and transactional commits, retry and dead-letter topics, partition-scoped transactions and exactly-once pipelines, a service template) via therdkafkaclient.
Concrete brokers live in their own crates and pull ruststream from crates.io.
Minimum supported Rust version
The MSRV is 1.85 (edition 2024, native async fn in trait). CI builds the crate on every
stable toolchain from 1.85 up to current stable, so any floor in that range works.
The policy:
- The published
rust-versionstays at the floor. Raising it is a breaking change (a minor version bump pre-1.0) and is reviewed against the broker crates' client requirements at each minor release. - Broker crates (
ruststream-nats, ...) may require a newer toolchain than the core when their underlying clients do; cargo allows a dependent crate to have a stricter floor than its dependency. Check the broker crate's ownrust-versionfor its floor.
Contributing
License
Licensed under the Apache-2.0 license.
Inspired by FastStream.