docs.rs failed to build rabbit_warren-0.1.0
Please check the
build logs for more information.
See
Builds for ideas on how to fix a failed build,
or
Metadata for how to configure docs.rs builds.
If you believe this is docs.rs' fault,
open an issue.
RabbitMQ Client
Надежная, эргономичная и готовая к production библиотека-клиент RabbitMQ для Rust, построенная поверх lapin.
Она абстрагирует типичный boilerplate-код (Publisher Confirms, флаг mandatory, автоматические стратегии ack/nack и распределенную трассировку), позволяя разработчикам сосредоточиться на бизнес-логике.
- Safe Publishing: Автоматически включает Publisher Confirms и использует флаг
mandatory для обнаружения немаршрутизируемых (unroutable) сообщений.
- Smart Consuming: Автоматический
ack при успехе. При ошибке вы сами выбираете стратегию: Requeue (со встроенным backoff для предотвращения tight loops) или Discard.
- Distributed Tracing: Встроенная пропаганда контекста OpenTelemetry (инъекция/извлечение trace-заголовков из AMQP-свойств).
- Zero Boilerplate: Нет необходимости вручную управлять каналами, ackers или опциями nack в коде вашего приложения.
- JSON Helpers: Встроенные утилиты
publish_json и deserialize_delivery.
Подключение
Добавьте это в ваш Cargo.toml:
[dependencies]
rabbit_warren = { path = "../libs/rabbit_warren" }
tokio = { version = "1", features = ["full"] }
serde = { version = "1.0", features = ["derive"] }
Примечание: Вам не нужно добавлять lapin в ваш Cargo.toml. Все необходимые типы (например, Delivery) реэкспортируются этим крейтом.
Примеры
1. Инициализация
use rabbit_warren::RmqClient;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let url = "amqp://guest:guest@127.0.0.1:5672";
let client = RmqClient::connect(url, 10).await?;
client.declare_exchange("my_exchange", rabbit_warren::ExchangeKind::Topic, true).await?;
client.declare_queue("my_queue", true).await?;
client.bind_queue("my_queue", "my_exchange", "routing.key").await?;
Ok(())
}
2. Отправка сообщений
use serde::Serialize;
#[derive(Serialize)]
struct MyEvent {
id: u32,
payload: String,
}
async fn publish_event(client: &RmqClient) -> anyhow::Result<()> {
let event = MyEvent { id: 1, payload: "hello".into() };
client.publish_json("my_exchange", "routing.key", &event, None).await?;
Ok(())
}
3. Потребление сообщений
Библиотека предоставляет трейт ResultExt, позволяющий декларативно указать, что должно произойти при возникновении ошибки, не загромождая ваш обработчик логикой ack/nack.
use rabbit_warren::{ProcessingError, ResultExt};
use serde::Deserialize;
#[derive(Deserialize)]
struct MyEvent {
id: u32,
payload: String,
}
async fn start_consuming(client: &RmqClient) {
client.consume("my_queue", |delivery| async move {
let event: MyEvent = rabbit_warren::deserialize_delivery(&delivery)
.map_err(|e| ProcessingError::retryable(e))?;
if event.id == 0 {
return Err(ProcessingError::permanent("Invalid ID".to_string()));
}
if event.id == 99 {
return Err(ProcessingError::retryable("Database temporarily unavailable".to_string()));
}
println!("Successfully processed event: {}", event.id);
Ok(())
});
}
4. Эргономичная обработка ошибок (ResultExt)
Вы можете использовать методы расширения прямо в цепочке вызовов, чтобы явно указать судьбу сообщения при конкретной ошибке:
use rabbit_warren::{RmqClient, ResultExt};
client.consume("my_queue", |delivery| async move {
let msg: MyMessage = rabbit_warren::deserialize_delivery(&delivery).requeue_on_err()?;
validate(&msg).discard_on_err()?;
save_to_db(&msg).await.requeue_on_err()?;
Ok(())
});
Или использовать кастомную логику маппинга ошибок:
use rabbit_warren::{RmqClient, ResultExt, ProcessingError};
client.consume("queue", |delivery| async move {
process(delivery)
.await
.map_err(|e| {
if e.is_transient() {
ProcessingError::retryable(e)
} else {
ProcessingError::permanent(e)
}
})?;
Ok(())
});
5. Рефакторинг консюмера
Эта библиотека позволяет заменить десятки строк ручного управления каналом, циклами while let и явными вызовами ack/nack на чистую бизнес-логику.
Как это выглядело раньше (с использованием "сырого" lapin):
use std::sync::Arc;
use futures_util::stream::StreamExt;
use lapin::{
options::{BasicAckOptions, BasicConsumeOptions, BasicNackOptions, ExchangeDeclareOptions, QueueBindOptions, QueueDeclareOptions},
types::FieldTable, Channel, ExchangeKind,
};
async fn create_order_listener(
channel: Arc<Channel>,
processor: Arc<dyn OrderProcessor>
) -> Result<tokio::task::JoinHandle<()>> {
let queue_name = std::env::var("ORDERS_QUEUE_NAME").unwrap_or_default();
let exchange = std::env::var("ORDERS_EXCHANGE_NAME").unwrap_or_default();
let routing_key = std::env::var("ORDERS_ROUTING_KEY").unwrap_or_default();
channel
.exchange_declare(
&exchange,
ExchangeKind::Topic,
ExchangeDeclareOptions { durable: true, ..Default::default() },
FieldTable::default(),
)
.await?;
channel
.queue_declare(
&queue_name,
QueueDeclareOptions { durable: true, ..Default::default() },
FieldTable::default()
)
.await?;
channel
.queue_bind(
&queue_name,
&exchange,
&routing_key,
QueueBindOptions::default(),
FieldTable::default(),
)
.await?;
let mut consumer = channel
.basic_consume(
&queue_name,
"order-listener-consumer",
BasicConsumeOptions::default(),
FieldTable::default(),
)
.await?;
let listener_task = tokio::task::spawn(async move {
let processor_clone = processor.clone();
while let Some(delivery_result) = consumer.next().await {
let delivery = match delivery_result {
Ok(delivery) => delivery,
Err(error) => {
error!("Failed to receive delivery from stream: {}", error);
continue;
}
};
debug!("Message consumed from {}", queue_name);
let span = common_opentelemetry::correlate_trace_from_delivery(&delivery);
let process_message = async {
let payload = std::str::from_utf8(&delivery.data)?;
if payload.trim() == "null" {
tracing::error!("Received null payload");
return Err(anyhow::anyhow!("Null payload received"));
}
let order: Order = serde_json::from_str(payload)?;
processor_clone.process(order).await?;
Ok(())
}
.instrument(span);
match process_message.await {
Ok(_) => {
delivery.ack(BasicAckOptions::default()).await.call_safe();
}
Err(error) => {
error!("Processing failed: {}", error);
let nack_options = BasicNackOptions {
requeue: true,
..Default::default()
};
delivery.nack(nack_options).await.call_safe();
}
}
}
});
Ok(listener_task)
}
Как это выглядит с rabbit_warren
Сначала мы централизованно определяем стратегию обработки ошибок для нашего приложения, чтобы избежать ситуации poison pill (бесконечного цикла обработки битого сообщения):
use rabbit_warren::{AckAction, IntoAckAction};
impl IntoAckAction for AppError {
fn ack_action(&self) -> AckAction {
match self {
AppError::InvalidJson(_) => AckAction::Discard,
AppError::DatabaseError(_) => AckAction::Requeue,
}
}
}
Теперь сам код консюмера становится компактным, безопасным и сфокусированным на бизнес-логике:
use rabbit_warren::{RmqClient, ResultExt, Delivery};
use std::sync::Arc;
async fn create_order_listener(
rmq_client: RmqClient,
processor: Arc<dyn OrderProcessor>
) -> Result<tokio::task::JoinHandle<()>> {
let queue_name = std::env::var("ORDERS_QUEUE_NAME").unwrap_or_default();
let exchange = std::env::var("ORDERS_EXCHANGE_NAME").unwrap_or_default();
let routing_key = std::env::var("ORDERS_ROUTING_KEY").unwrap_or_default();
rmq_client.declare_exchange(&exchange, ExchangeKind::Topic, true).await?;
rmq_client.declare_queue(&queue_name, true).await?;
rmq_client.bind_queue(&queue_name, &exchange, &routing_key).await?;
let handler = move |delivery: Delivery| {
let processor = processor.clone();
async move {
let process_result: std::result::Result<(), AppError> = async {
let payload = rabbit_warren::payload_as_utf8(&delivery)?;
if payload.trim() == "null" {
tracing::error!("Received null payload");
return Err(AppError::InvalidJson("null payload".into()));
}
let order: Order = serde_json::from_str(payload)?;
processor.process(order).await?;
Ok(())
}.await;
process_result.into_processing_error()
}
};
Ok(rmq_client.consume(&queue_name, handler))
}
Возможность кастомной стратегии для Ack
Если в вашем приложении есть собственный enum ошибок, вы можете реализовать IntoAckAction, чтобы централизовать логику повторных попыток:
use rabbit_warren::{AckAction, IntoAckAction};
#[derive(Debug, thiserror::Error)]
enum AppError {
#[error("Database connection lost")]
DbConnection,
#[error("Validation failed: {0}")]
Validation(String),
}
impl IntoAckAction for AppError {
fn ack_action(&self) -> AckAction {
match self {
AppError::DbConnection => AckAction::Requeue, AppError::Validation(_) => AckAction::Discard, }
}
}
client.consume("queue", |delivery| async move {
do_work().await.map_err(AppError::DbConnection)?;
Ok(())
});
Примечание: По умолчанию любой тип ошибки, не реализующий IntoAckAction, будет обработан как AckAction::Requeue для предотвращения случайной потери данных.
Трассировка
Библиотека интегрируется с экосистемой tracing. Метод consume автоматически создает info_span для каждого сообщения.
Для полноценной распределенной трассировки (OpenTelemetry) вы можете легко извлекать и внедрять контекст в заголовки AMQP, используя стандартный крейт opentelemetry и передавая результат в аргумент headers:
use rabbit_warren::FieldTable;
use opentelemetry::global;
let mut headers = FieldTable::default();
global::get_text_map_propagator(|propagator| {
});
client.publish_json("exchange", "key", &msg, Some(headers)).await?;
Тесты
Для запуска интеграционных тестов у вас должен быть установлен Docker, так как тесты используют testcontainers для поднятия временного (ephemeral) экземпляра RabbitMQ:
cargo test --all-features