rabbit_warren 0.1.0

An ergonomic, production-ready RabbitMQ client built on top of lapin, featuring automatic ack/nack strategies and distributed tracing.
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" } # Или путь к вашему workspace/git
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";
    // Подключение и установка 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. Отправка сообщений

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.

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)

Вы можете использовать методы расширения прямо в цепочке вызовов, чтобы явно указать судьбу сообщения при конкретной ошибке:

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(())
});

Или использовать кастомную логику маппинга ошибок:

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<()>> {
    // 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 (бесконечного цикла обработки битого сообщения):

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, 
        }
    }
}

Теперь сам код консюмера становится компактным, безопасным и сфокусированным на бизнес-логике:

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, чтобы централизовать логику повторных попыток:

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 или отбросить
        }
    }
}

// Использование:
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| {
    // Ваш кастомный Carrier для FieldTable
    // propagator.inject_context(&context, &mut carrier);
});
client.publish_json("exchange", "key", &msg, Some(headers)).await?;

Тесты

Для запуска интеграционных тестов у вас должен быть установлен Docker, так как тесты используют testcontainers для поднятия временного (ephemeral) экземпляра RabbitMQ:

cargo test --all-features