# RabbitMQ Client
Надежная, эргономичная и готовая к production библиотека-клиент RabbitMQ для Rust, построенная поверх [`lapin`](https://crates.io/crates/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`:
```toml
[dependencies]
rabbit_warren = { path = "../libs/rabbit_warren" } # Или путь к вашему workspace/git
tokio = { version = "1", features = ["full"] }
serde = { version = "1.0", features = ["derive"] }
```
*Примечание: Вам не нужно добавлять lapin в ваш Cargo.toml. Все необходимые типы (например, Delivery) реэкспортируются этим крейтом.*
## Примеры
### 1. Инициализация
```rust
use rabbit_warren::RmqClient;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let url = "amqp://guest:guest@127.0.0.1:5672";
// Подключение и установка prefetch count в 10
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. Отправка сообщений
```rust
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() };
// Автоматически сериализует в JSON, устанавливает delivery_mode=2 (persistent),
// и ожидает publisher confirmation.
client.publish_json("my_exchange", "routing.key", &event, None).await?;
Ok(())
}
```
### 3. Потребление сообщений
Библиотека предоставляет трейт ResultExt, позволяющий декларативно указать, что должно произойти при возникновении ошибки, не загромождая ваш обработчик логикой ack/nack.
```rust
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 {
// 1. Парсинг сообщения. Если он не удался, делаем requeue (временная проблема или мы хотим повторить).
let event: MyEvent = rabbit_warren::deserialize_delivery(&delivery)
.map_err(|e| ProcessingError::retryable(e))?;
// 2. Бизнес-логика.
// Если это перманентная бизнес-ошибка (например, невалидные данные), делаем discard, чтобы избежать poison pills.
if event.id == 0 {
return Err(ProcessingError::permanent("Invalid ID".to_string()));
}
// 3. Симуляция временной ошибки (например, таймаут БД).
// Библиотека автоматически сделает паузу (backoff) и вернет сообщение в очередь (requeue).
if event.id == 99 {
return Err(ProcessingError::retryable("Database temporarily unavailable".to_string()));
}
println!("Successfully processed event: {}", event.id);
// Возврат Ok(()) автоматически подтверждает (ack) сообщение.
Ok(())
});
}
```
### 4. Эргономичная обработка ошибок (ResultExt)
Вы можете использовать методы расширения прямо в цепочке вызовов, чтобы явно указать судьбу сообщения при конкретной ошибке:
```rust
use rabbit_warren::{RmqClient, ResultExt};
client.consume("my_queue", |delivery| async move {
// Ошибка десериализации -> requeue (временная проблема)
let msg: MyMessage = rabbit_warren::deserialize_delivery(&delivery).requeue_on_err()?;
// Ошибка валидации -> discard (бесполезно повторять)
validate(&msg).discard_on_err()?;
// Ошибка БД -> requeue (временная проблема)
save_to_db(&msg).await.requeue_on_err()?;
Ok(())
});
```
Или использовать кастомную логику маппинга ошибок:
```rust
use rabbit_warren::{RmqClient, ResultExt, ProcessingError};
.await
.map_err(|e| {
if e.is_transient() {
ProcessingError::retryable(e)
} else {
ProcessingError::permanent(e)
}
})?;
Ok(())
});
```
### 5. Рефакторинг консюмера
Эта библиотека позволяет заменить десятки строк ручного управления каналом, циклами `while let` и явными вызовами `ack`/`nack` на чистую бизнес-логику.
**Как это выглядело раньше (с использованием "сырого" `lapin`):**
```rust
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<()>> {
// 1. Многословное объявление топологии
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?;
// 2. Ручное управление жизненным циклом задачи и стримом
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);
// 3. Вложенная бизнес-логика
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);
// 4. Ручное управление ack/nack (высокий риск poison pill!)
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` (бесконечного цикла обработки битого сообщения):
```rust
use rabbit_warren::{AckAction, IntoAckAction};
// Реализуем трейт для нашего типа ошибки
impl IntoAckAction for AppError {
fn ack_action(&self) -> AckAction {
match self {
// Перманентные ошибки (битый JSON, "null") -> удаляем (Discard)
AppError::InvalidJson(_) => AckAction::Discard,
// Временные ошибки (сеть, БД, таймауты) -> возвращаем в очередь (Requeue)
AppError::DatabaseError(_) => AckAction::Requeue,
}
}
}
```
Теперь сам код консюмера становится компактным, безопасным и сфокусированным на бизнес-логике:
```rust
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<()>> {
// 1. Декларативное и лаконичное объявление топологии
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?;
// 2. Чистая бизнес-логика в замыкании
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?; // Может вернуть AppError::DatabaseError
Ok(())
}.await;
// 3. Магия ResultExt: библиотека сама вызовет ack_action() и выберет Discard или Requeue
process_result.into_processing_error()
}
};
// 4. Запуск в фоновой задаче (библиотека сама управляет циклом и ack/nack)
Ok(rmq_client.consume(&queue_name, handler))
}
```
## Возможность кастомной стратегии для Ack
Если в вашем приложении есть собственный enum ошибок, вы можете реализовать IntoAckAction, чтобы централизовать логику повторных попыток:
```rust
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, // Перманентный сбой, отправить в DLQ или отбросить
}
}
}
// Использование:
do_work().await.map_err(AppError::DbConnection)?;
Ok(())
});
```
*Примечание: По умолчанию любой тип ошибки, не реализующий IntoAckAction, будет обработан как AckAction::Requeue для предотвращения случайной потери данных.*
## Трассировка
Библиотека интегрируется с экосистемой `tracing`. Метод `consume` автоматически создает `info_span` для каждого сообщения.
Для полноценной распределенной трассировки (OpenTelemetry) вы можете легко извлекать и внедрять контекст в заголовки AMQP, используя стандартный крейт `opentelemetry` и передавая результат в аргумент `headers`:
```rust
use rabbit_warren::FieldTable;
use opentelemetry::global;
// Пример внедрения заголовков при публикации
let mut headers = FieldTable::default();
global::get_text_map_propagator(|propagator| {
// Ваш кастомный Carrier для FieldTable
// propagator.inject_context(&context, &mut carrier);
});
client.publish_json("exchange", "key", &msg, Some(headers)).await?;
```
## Тесты
Для запуска интеграционных тестов у вас должен быть установлен Docker, так как тесты используют [`testcontainers`](https://crates.io/crates/testcontainers) для поднятия временного (ephemeral) экземпляра RabbitMQ:
```bash
cargo test --all-features
```