use std::time::Duration;
use futures::StreamExt;
use ruststream::{
Broker, Headers, IncomingMessage, OutgoingMessage, Publisher, RequestReply, Subscriber,
};
use ruststream_lapin::{LapinBroker, RabbitQueue};
#[tokio::main(flavor = "multi_thread")]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let broker = LapinBroker::new("amqp://localhost:5672").declare_topology(true);
Broker::connect(&broker).await?;
let mut responder = broker
.subscribe(RabbitQueue::new("rpc.ping").durable(false).exclusive(true))
.await?;
let reply_publisher = broker.publisher();
tokio::spawn(async move {
let mut requests = std::pin::pin!(responder.stream());
while let Some(Ok(request)) = requests.next().await {
let Some(reply_to) = request.headers().reply_to().map(str::to_owned) else {
continue;
};
let mut headers = Headers::new();
if let Some(correlation_id) = request.headers().correlation_id() {
headers.insert("correlation-id", correlation_id.as_bytes().to_vec());
}
let reply = OutgoingMessage::new(&reply_to, b"pong").with_headers(headers);
let _ = reply_publisher.publish(reply).await;
let _ = request.ack().await;
}
});
let requester = broker.requester();
let reply = requester
.request(
OutgoingMessage::new("rpc.ping", b"ping"),
Duration::from_secs(2),
)
.await?;
println!("reply: {}", String::from_utf8_lossy(reply.payload()));
broker.shutdown().await?;
Ok(())
}