use backbone_queue::{
QueueService,
redis::RedisQueueBuilder,
types::{QueueMessage, QueuePriority}
};
use std::time::Duration;
use tokio::time::sleep;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
env_logger::init();
println!("๐ Basic Redis Queue Example");
println!("============================");
println!("๐ก Connecting to Redis...");
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("example_queue")
.key_prefix("example")
.pool_size(5)
.build()
.await?;
if !queue.test_connection().await? {
eprintln!("โ Failed to connect to Redis");
return Ok(());
}
println!("โ
Connected to Redis successfully");
println!("๐งน Clearing existing messages...");
queue.purge().await?;
println!("โ
Queue cleared");
println!("\n๐ฌ Example 1: Basic Message Operations");
println!("-------------------------------------");
let message = QueueMessage::builder()
.payload("Hello, Redis Queue!").expect("payload serialization")
.priority(QueuePriority::Normal)
.build();
println!("๐ค Enqueuing message: {}", message.payload);
let message_id = queue.enqueue(message.clone()).await?;
println!("โ
Message enqueued with ID: {}", message_id);
let size = queue.size().await?;
println!("๐ Queue size: {}", size);
println!("๐ฅ Dequeuing message...");
if let Some(received_message) = queue.dequeue().await? {
println!("โ
Received message: {}", received_message.payload);
println!("๐ Message ID: {}", received_message.id);
println!("โ๏ธ Priority: {}", received_message.priority);
println!("๐
Enqueued at: {}", received_message.enqueued_at);
queue.ack(&received_message.id).await?;
println!("โ
Message acknowledged");
} else {
println!("โ No message received");
}
println!("\n๐ Example 2: Priority Message Ordering");
println!("--------------------------------------");
let priority_messages = vec![
QueueMessage::builder()
.payload("Low priority task").expect("payload serialization")
.priority(QueuePriority::Low)
.build(),
QueueMessage::builder()
.payload("Critical emergency task").expect("payload serialization")
.priority(QueuePriority::Critical)
.build(),
QueueMessage::builder()
.payload("Normal background task").expect("payload serialization")
.priority(QueuePriority::Normal)
.build(),
QueueMessage::builder()
.payload("High priority task").expect("payload serialization")
.priority(QueuePriority::High)
.build(),
];
println!("๐ค Enqueuing messages with different priorities...");
for message in &priority_messages {
let id = queue.enqueue(message.clone()).await?;
println!(" - {} ({})", message.payload, message.priority);
}
println!("\n๐ฅ Dequeue order (should be by priority):");
for _ in 0..priority_messages.len() {
if let Some(message) = queue.dequeue().await? {
println!(" 1๏ธโฃ {} ({})", message.payload, message.priority);
queue.ack(&message.id).await?;
}
}
println!("\n๐ฆ Example 3: Batch Operations");
println!("----------------------------");
let mut batch_messages = Vec::new();
for i in 1..=10 {
batch_messages.push(QueueMessage::builder()
.payload(format!("Batch message {}", i)).expect("payload serialization")
.priority(QueuePriority::Normal)
.build());
}
println!("๐ค Enqueuing {} messages in batch...", batch_messages.len());
let start_time = std::time::Instant::now();
let message_ids = queue.enqueue_batch(batch_messages.clone()).await?;
let enqueue_time = start_time.elapsed();
println!("โ
Batch enqueued in {:?}", enqueue_time);
println!("๐ Enqueue rate: {:.2} messages/sec", message_ids.len() as f64 / enqueue_time.as_secs_f64());
let stats = queue.get_stats().await?;
println!("๐ Queue Stats:");
println!(" - Visible messages: {}", stats.visible_messages);
println!(" - Invisible messages: {}", stats.invisible_messages);
println!(" - Total messages: {}", stats.total_messages);
println!("\n๐ฅ Dequeuing messages in batch...");
let start_time = std::time::Instant::now();
let batch_result = queue.dequeue_batch(5).await?;
let dequeue_time = start_time.elapsed();
println!("โ
Received {} messages in {:?}", batch_result.messages.len(), dequeue_time);
println!("๐ Dequeue rate: {:.2} messages/sec", batch_result.messages.len() as f64 / dequeue_time.as_secs_f64());
let ack_ids: Vec<String> = batch_result.messages
.iter()
.map(|m| m.id.clone())
.collect();
let start_time = std::time::Instant::now();
let ack_count = queue.ack_batch(ack_ids).await?;
let ack_time = start_time.elapsed();
println!("โ
Acknowledged {} messages in {:?}", ack_count, ack_time);
println!("\n๐ท๏ธ Example 4: Messages with Attributes");
println!("------------------------------------");
use std::collections::HashMap;
let mut attributes = HashMap::new();
attributes.insert("source".to_string(), "api".to_string());
attributes.insert("user_id".to_string(), "12345".to_string());
attributes.insert("request_id".to_string(), "req-abc-123".to_string());
let message_with_attrs = QueueMessage::builder()
.payload("Process user data").expect("payload serialization")
.priority(QueuePriority::High)
.attributes(attributes)
.build();
println!("๐ค Enqueuing message with custom attributes...");
let id = queue.enqueue(message_with_attrs.clone()).await?;
println!("โ
Message enqueued with ID: {}", id);
if let Some(received_message) = queue.dequeue().await? {
println!("๐ฅ Received message with attributes:");
println!(" - Payload: {}", received_message.payload);
println!(" - Attributes:");
for (key, value) in &received_message.attributes {
println!(" {}: {}", key, value);
}
queue.ack(&received_message.id).await?;
}
println!("\n๐ฅ Example 5: Health Monitoring");
println!("------------------------------");
let health = queue.health_check().await?;
println!("๐ฅ Queue Health Status:");
println!(" - Status: {:?}", health.status);
println!(" - Queue size: {}", health.queue_size);
println!(" - Error rate: {:.2}%", health.error_rate * 100.0);
println!(" - Last activity: {:?}", health.last_activity);
println!(" - Checked at: {:?}", health.checked_at);
println!("\n๐ Final Queue Statistics");
println!("=========================");
let final_stats = queue.get_stats().await?;
println!(" - Visible messages: {}", final_stats.visible_messages);
println!(" - Invisible messages: {}", final_stats.invisible_messages);
println!(" - Dead letter messages: {}", final_stats.dead_letter_messages);
println!(" - Total processed: {}", final_stats.total_processed);
println!(" - Total failed: {}", final_stats.total_failed);
println!("\n๐งน Cleaning up...");
let purged_count = queue.purge().await?;
println!("โ
Purged {} messages", purged_count);
println!("\n๐ Example completed successfully!");
Ok(())
}