use crate::error::BrokerError;
use crate::models::Queue;
use redis::aio::MultiplexedConnection;
use redis::AsyncCommands;
use std::collections::HashSet;
pub struct QueueParser;
impl QueueParser {
pub async fn parse_queues(
connection: &MultiplexedConnection,
) -> Result<Vec<Queue>, BrokerError> {
let mut conn = connection.clone();
let mut queues = Vec::new();
let mut discovered_queues = HashSet::new();
let binding_keys: Vec<String> = conn.keys("_kombu.binding.*").await.unwrap_or_default();
for binding_key in binding_keys {
if let Some(queue_name) = binding_key.strip_prefix("_kombu.binding.") {
discovered_queues.insert(queue_name.to_string());
}
}
let common_queues = vec!["celery", "default", "priority", "high", "low"];
for queue_name in common_queues {
discovered_queues.insert(queue_name.to_string());
}
for queue_name in discovered_queues {
let length: u64 = conn.llen(&queue_name).await.unwrap_or(0);
if length > 0 || ["celery", "default"].contains(&queue_name.as_str()) {
let consumers = if length > 0 { 1 } else { 0 };
queues.push(Queue {
name: queue_name,
length,
consumers,
});
}
}
queues.sort_by(|a, b| a.name.cmp(&b.name));
Ok(queues)
}
}