pub mod memory;
#[cfg(feature = "rabbitmq")]
pub mod rabbitmq;
#[cfg(feature = "redis")]
pub mod redis;
use crate::config::QueueKind;
use crate::{RagConfig, RagError, Result};
use async_trait::async_trait;
use std::sync::Arc;
pub const TOPIC: &str = "docling_rag_ingest";
#[async_trait]
pub trait QueueReceiver: Send {
async fn recv(&mut self) -> Option<Vec<u8>>;
}
#[async_trait]
pub trait MessageQueue: Send + Sync {
async fn publish(&self, payload: &[u8]) -> Result<()>;
async fn subscribe(&self) -> Result<Box<dyn QueueReceiver>>;
}
pub async fn from_config(cfg: &RagConfig) -> Result<Arc<dyn MessageQueue>> {
match cfg.queue {
QueueKind::Memory => Ok(Arc::new(memory::MemoryQueue::new())),
QueueKind::RabbitMq => {
#[cfg(feature = "rabbitmq")]
{
let url = cfg.rabbitmq_url.clone().ok_or_else(|| {
RagError::config("RABBITMQ_URL is required for the rabbitmq queue")
})?;
Ok(Arc::new(rabbitmq::RabbitMqQueue::connect(&url).await?))
}
#[cfg(not(feature = "rabbitmq"))]
{
Err(RagError::FeatureDisabled(
"rabbitmq".into(),
"rabbitmq".into(),
))
}
}
QueueKind::Redis => {
#[cfg(feature = "redis")]
{
let url = cfg
.redis_url
.clone()
.ok_or_else(|| RagError::config("REDIS_URL is required for the redis queue"))?;
Ok(Arc::new(redis::RedisQueue::connect(&url).await?))
}
#[cfg(not(feature = "redis"))]
{
Err(RagError::FeatureDisabled("redis".into(), "redis".into()))
}
}
}
}