ruststream-amqp 0.6.0

AMQP 1.0 broker implementation for the RustStream messaging framework (ActiveMQ Artemis, Azure Service Bus, RabbitMQ 4.x, and other AMQP 1.0 brokers).
Documentation

ruststream-amqp implements the RustStream broker contract over fe2o3-amqp. Handlers, routers, codecs, and middleware come from the framework; this crate supplies the transport - and nothing broker-specific leaks back into the framework.

AMQP 1.0 is an ISO-standard protocol spoken by ActiveMQ Artemis and Classic, RabbitMQ 4.x (a separate protocol stack from the 0.9.1 that ruststream-lapin speaks), Azure Service Bus and Event Hubs, Amazon MQ, Solace, Apache Qpid, and IBM MQ - one crate serves the whole family.

Features

  • Lazy startup contract. AmqpBroker::new(url) is synchronous and does no I/O; the runtime connects once at startup, so the broker composes with #[ruststream::app]. SASL (ANONYMOUS, PLAIN, EXTERNAL) and the container id are builder options.
  • Acknowledgement as dispositions. ack maps to accept, nack(requeue = true) to release, nack(requeue = false) to reject - the broker's own dead-letter policy applies. At-most-once subscriptions report AckError::Unsupported instead of pretending.
  • Explicit addressing. The protocol standardises the wire, not the meaning of an address: AmqpAddress::queue (anycast), AmqpAddress::topic (multicast), AmqpAddress::raw (verbatim, for deployments with their own convention), plus credit (prefetch as protocol-level flow control) and the settle guarantee.
  • Native request/reply. AmqpPublisher implements the RequestReply capability over reply-to, correlation-id, and a dynamic receiver link.
  • Transactions (feature transaction). A distinct AmqpTransactionalPublish policy pairs into a TransactionalPublisher built on the protocol's transactional posting; the plain publisher carries no transactional surface.
  • Headers without an envelope. Well-known headers ride the properties section (content-type, correlation-id, reply-to, message-id, the partition key as group-id); everything else rides application-properties, so non-Rust peers see plain AMQP messages.
  • In-process test broker (feature testing). AmqpTestBroker reproduces core routing with no server, implements ruststream::testing::TestableBroker, and passes the framework's conformance suite in process.

Status

Implemented and verified against ActiveMQ Artemis (the framework's conformance lifecycle, request/reply, and transactions suites run in CI against a live broker). Built on ruststream 0.6 from crates.io; the crate itself is not published yet. Design and scope are tracked in powersemmi/ruststream#187.

Write a service

use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream};
use ruststream::subscriber;
use ruststream_amqp::{AmqpAddress, AmqpBroker};
use serde::Deserialize;

#[derive(Debug, Deserialize)]
struct Order {
    id: u64,
}

#[subscriber(AmqpAddress::queue("orders"))]
async fn handle(order: &Order) -> HandlerResult {
    println!("got order {}", order.id);
    HandlerResult::Ack
}

#[ruststream::app]
fn app() -> impl App {
    RustStream::new(AppInfo::new("orders", "0.1.0"))
        .with_broker(AmqpBroker::new("amqp://localhost:5672"), |b| b.include(handle))
}

The descriptor carries the AMQP-specific options inline in the decorator:

#[subscriber(AmqpAddress::queue("orders").credit(64))]
async fn handle(order: &Order) -> HandlerResult { /* ... */ }

Request/reply

The requester side is a first publish, so it belongs in the scope's after_startup hook: the publisher arrives live, already paired with the connected broker.

use std::io;
use std::time::Duration;

use ruststream::runtime::{App, AppInfo, RustStream};
use ruststream::{IncomingMessage, OutgoingMessage, RequestReply};
use ruststream_amqp::{AmqpBroker, AmqpPublish};

#[ruststream::app]
fn app() -> impl App {
    RustStream::new(AppInfo::new("greeter-client", "0.1.0"))
        .with_broker(AmqpBroker::new("amqp://localhost:5672"), |b| {
            b.after_startup(AmqpPublish, async move |publisher| -> io::Result<()> {
                let reply = publisher
                    .request(
                        OutgoingMessage::new("greeter", b"hello".as_slice()),
                        Duration::from_secs(5),
                    )
                    .await
                    .map_err(io::Error::other)?;
                println!("{}", String::from_utf8_lossy(reply.payload()));
                Ok(())
            });
        })
}

The responder answers on the requester's dynamic reply-to address, so it publishes through an injected publisher rather than the fixed-destination publish(..) reply form; see examples/amqp_request_reply.rs for both sides in one app.

Test it

The testing feature runs handlers against an in-process AMQP stand-in - no server, same routing. Broker-specific behaviour (dispositions, credit, dead-lettering) is covered by the env-gated live suite instead: just test-brokers spins up ActiveMQ Artemis and runs the integration tests plus the framework conformance suites against it.

Layout

ruststream-amqp/
├── crates/
│   └── ruststream-amqp/        the published crate
│       └── examples/           runnable amqp_* examples
├── docker-compose.test.yml     ActiveMQ Artemis for the live suite
└── Cargo.toml                  workspace

Contributing

just check          # fmt, clippy, feature checks
just test           # handler-stub tests, no server
just test-brokers   # live integration + conformance against ActiveMQ Artemis

License

Licensed under the Apache-2.0 license.