rAMQP
A from-scratch, clean-room AMQP 1.0 stack for Rust on tokio — a published
client, a shared protocol engine, and a performance-first, highly-available
broker — with no external AMQP dependencies anywhere.
| Crate | What it is | Status |
|---|---|---|
ramqp |
The async client — connects to RabbitMQ 4.x, ActiveMQ Artemis, and other AMQP 1.0 brokers | Published (0.8.1 — see Upgrading to 0.8) |
ramqp-core |
The role-agnostic engine: clean-room codec + type system, framing, session/link state machines, SASL (both directions) | Published (0.2.4) |
ramqp-broker |
The broker: store-and-forward AMQP 1.0 server with transient, durable, and Raft-replicated quorum queues | Working, pre-1.0 — first crates.io release (0.9.0) ships with the next tag; see The broker |
Everything is #![forbid(unsafe_code)], async-first, and MIT.
The client (ramqp)
use ;
async
- Single-pass, zero-copy framing — bodies are exposed as
bytes::Bytesslices; the transfer/body split is computed once from the negotiatedmax-frame-size, never by trial re-serialization. - Lock-free actor runtime — one owning driver task per connection holds all protocol state; user handles are cheap clones that exchange messages over bounded channels. No locks on the per-message path.
- Lazy receive by default —
recv()yields aDeliveryexposing raw bytes; typed decoding (.message(),.decode::<T>()) is an explicit, cheap opt-in. - Flat, classified errors — one error type per operation (
ConnectError,SendError,RecvError,SessionError,LinkError) withsource()chains,is_retryable()/is_fatal(), and typed access to any peer-sent error. - Transparent reconnect — opt in with
ConnectionBuilder::reconnecting(true)and yourProducer/Consumerhandles survive a broker drop: reconnection with backoff, re-attach, and in-flight replay happen behind the scenes. - Pluggable observability — a
Metricstrait and a connection-event stream, usable without enablingtracing.
| Feature | Effect |
|---|---|
rustls |
amqps:// via rustls + webpki-roots |
native-tls |
amqps:// via native-tls |
ws |
ws:// (AMQP over WebSocket); wss:// also needs rustls |
scram |
SASL SCRAM-SHA-1/256/512 |
transaction |
Transaction controller (clean-room, spec part 4) |
ANONYMOUS / PLAIN / EXTERNAL SASL are always available.
Verified against live RabbitMQ 4.x and ActiveMQ Artemis — the interop suite runs in CI against both: SASL, open/begin/attach, credit flow, transfer, all four settlement outcomes, custom-CA TLS, AMQP-over-WebSocket, and transparent mid-stream reconnect.
How the client compares
Benchmarked head-to-head against the established fe2o3-amqp client over a
live broker, with matched credit windows and warmup (see
bench-compare/). On RabbitMQ 4.x, 5000 messages, receive
throughput (median msg/s):
| body | ramqp (per-msg ack) | fe2o3-amqp (per-msg ack) | ramqp (batched ack) |
|---|---|---|---|
| 64 B | ~180,000 | ~177,000 | ~230,000 |
| 1 KB | ~119,000 | ~129,000 | ~150,000 |
| 8 KB | ~42,000 | ~42,000 | ~53,000 |
At parity on the standard per-message path; ~1.2–1.4× faster with batched
ranged settlement (accept_through). Numbers are directional (single-node
RabbitMQ); reproduce them with the harness in bench-compare/.
Upgrading to 0.8
Your Cargo.toml and code need no changes. ramqp = "0.8" behaves exactly
like 0.7: the 0.8.0 release is an internal restructure — the protocol engine
moved into the new ramqp-core crate, and ramqp re-exports every moved
module, so all ramqp::... paths, feature names, and types are unchanged
(compile-time-locked by tests/public_api.rs). Cargo pulls in ramqp-core
automatically as an ordinary transitive dependency; the scram/transaction
features transparently enable their ramqp-core counterparts.
What's new is optional, for your Cargo.toml only if you want it:
= "0.8" # the client, exactly as before
= "0.2" # just the engine (codec/types/state machines), no client
= "0.9" # embed the broker (working, pre-1.0; config API still settling)
A client-only build never compiles broker code, and vice versa — isolation is by crate boundary, not feature flags.
The broker (ramqp-broker)
A performance-first, highly-available AMQP 1.0 broker on the same clean-room
engine. Working, pre-1.0 — the wire behavior is exercised by conformance,
cross-client interop, and partition/chaos suites; the Rust config API is
still settling (#[non_exhaustive] types). The design, targets, and phased
plan live in broker.md; its §11 checkboxes are the live status.
Working today:
- TCP acceptor, server-side handshake, SASL ANONYMOUS/PLAIN with a pluggable authenticator — any AMQP 1.0 client connects out of the box.
- Transient queues (
/queues/<name>, auto-declared): in-memory queue actors with competing consumers, credit-based dispatch, accept/release/modified/reject settlement, redelivery on consumer failure, and bounded depth with overflow rejection. - Quorum queues (
/quorum/<name>): backed by a per-queue Raft group (openraft) — a publish is acknowledged only after the enqueue commits to the replicated log; snapshots + log compaction keep memory tracking queue depth, not history. Single-replica today; multi-node placement and failover routing are the current work. - Cluster foundation: a metadata Raft group (replicated queue catalog) over a real TCP inter-node transport with static-seed bootstrap. Three-node clusters form, replicate, and survive leader failure with re-election — with zero committed-message loss (tested).
- The
ramqp-brokerddaemon.
# then point any AMQP 1.0 client at it:
RAMQP_URL=amqp://localhost:5672 RAMQP_ADDRESS=/queues/demo \
First broker numbers
Same machine, same harness, same client stack on both legs; 256 B payloads,
untuned defaults (methodology and honest caveats in
bench-compare/README.md):
| 256 B closed-loop | ramqp-broker | RabbitMQ 4.3.1 | Artemis |
|---|---|---|---|
| p50 / p99 / p99.9 latency | 89 / 213 / 428 µs | 251 / 519 / 777 µs | 227 / 576 / 833 µs |
| blast throughput | 326k msg/s | 48k msg/s | 79k msg/s |
| broker memory | ~40 MiB (incl. client) | 133 MiB | 715 MiB |
Quorum queues (every message Raft-committed; single replica): ~202k msg/s, depth-flat to a 50k backlog, p50 ~116 µs.
Performance is the product for this broker — the targets, hot-path rules, and
benchmark-as-merge-gate policy are broker.md §3.
Architecture
ramqp (client)
PUBLIC API Connection · Session · Producer · Consumer (ramqp/src/api)
RESILIENCE supervisor · reconnect · replay · pool (ramqp/src/resilience)
DIAL + DRIVER connect TCP/TLS/WS · client driver task · SASL (ramqp/src/{transport,connection,sasl})
ramqp-core (shared engine)
LINK sender/receiver · settlement · credit · delivery (ramqp-core/src/link)
SESSION begin/end · windows · registry · both polarities
(client attach + server accept) (ramqp-core/src/session)
CONNECTION open negotiation · mux · heartbeat (ramqp-core/src/connection)
TRANSPORT single-pass frame codec · header (both orders) (ramqp-core/src/transport)
SASL SCRAM math · server-side machinery (ramqp-core/src/sasl)
CONTRACTS errors · ids · config · metrics/events · proto (ramqp-core/src/{error,ids,config,observe,proto})
CODEC + TYPES clean-room AMQP 1.0 type system + wire codec (ramqp-core/src/{codec,types})
ramqp-broker (the broker)
FRONTEND acceptor · server handshake · connection driver (ramqp-broker/src/{broker,connection,auth})
QUEUES transient actors · quorum (Raft-backed) actors (ramqp-broker/src/{queue,quorum,registry})
CLUSTER metadata group · per-queue groups · TCP Raft
transport · static-seed bootstrap (ramqp-broker/src/cluster)
DAEMON ramqp-brokerd (ramqp-broker/src/bin)
Tests & benchmarks
# Live-broker interop — #[ignore]d so a plain `cargo test` stays green;
# CI runs this for real against RabbitMQ 4.x and Artemis:
RAMQP_BROKER_URL=amqp://guest:guest@localhost:5672 \
RAMQP_BROKER_ADDRESS=/queues/ramqp_it \
# Broker latency/throughput/RSS harness (ours in-process, or any AMQP URL):
Releasing
Publish order matters now that the repo is a workspace — see
RELEASING.md. Short version: ramqp-core must be on
crates.io before ramqp 0.8.0 (the client's manifest requires it by
version); ramqp-broker stays unpublished until its API stabilizes;
bench-compare is never published.
License
MIT.