use backbone_queue::{
rabbitmq_simple::{RabbitMQQueueSimple, RabbitMQConfig, ExchangeType},
traits::QueueService,
types::{QueueMessage, QueuePriority},
utils::rabbitmq_simple::*,
};
use std::collections::HashMap;
use std::time::Duration;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
env_logger::init();
println!("đ° RabbitMQ Queue Examples\n");
println!("========================\n");
println!("đŦ Example 1: Basic Direct Exchange");
basic_direct_exchange_example().await?;
println!();
println!("đĸ Example 2: Fanout Exchange Broadcasting");
fanout_exchange_example().await?;
println!();
println!("đ¯ Example 3: Topic Exchange Pattern Matching");
topic_exchange_example().await?;
println!();
println!("⥠Example 4: Priority Message Handling");
priority_message_example().await?;
println!();
println!("đĻ Example 5: Batch Operations");
batch_operations_example().await?;
println!();
println!("đˇī¸ Example 6: Messages with Custom Headers");
message_with_headers_example().await?;
println!();
println!("đĄī¸ Example 7: Configuration Validation");
validation_example().await?;
println!();
println!("đ Example 8: Queue Health Monitoring");
health_monitoring_example().await?;
println!();
println!("â
All examples completed successfully!");
Ok(())
}
async fn basic_direct_exchange_example() -> Result<(), Box<dyn std::error::Error>> {
let config = dev_config("user_notifications", "notifications_direct");
let queue = RabbitMQQueueSimple::new(config).await?;
let message = QueueMessage::builder()
.payload(serde_json::json!({
"type": "email_notification",
"user_id": 12345,
"email": "user@example.com",
"subject": "Welcome to our service!",
"body": "Thank you for signing up. Here's how to get started..."
}))
.expect("Failed to serialize payload")
.priority(QueuePriority::Normal)
.routing_key("user.email")
.build();
let message_id = queue.enqueue(message).await?;
println!(" â Enqueued notification message: {}", message_id);
let dequeued = queue.dequeue().await?;
if let Some(msg) = dequeued {
println!(" â Dequeued message: {}", msg.id);
queue.ack(&msg.id).await?;
println!(" â Acknowledged message");
} else {
println!(" âšī¸ No messages available (simulated)");
}
Ok(())
}
async fn fanout_exchange_example() -> Result<(), Box<dyn std::error::Error>> {
let config = RabbitMQConfig {
connection_url: "amqp://guest:guest@localhost:5672/%2f".to_string(),
queue_name: "system_events".to_string(),
exchange_name: "system_events_fanout".to_string(),
exchange_type: ExchangeType::Fanout,
routing_key: None, };
let queue = RabbitMQQueueSimple::new(config).await?;
let system_event = QueueMessage::builder()
.payload(serde_json::json!({
"event_type": "system_maintenance",
"timestamp": chrono::Utc::now().to_rfc3339(),
"message": "System maintenance scheduled for 2:00 AM UTC",
"affected_services": ["api", "database", "cache"],
"duration_minutes": 30
}))
.expect("Failed to serialize payload")
.priority(QueuePriority::High)
.build();
let event_id = queue.enqueue(system_event).await?;
println!(" â Broadcast system event: {}", event_id);
let services = ["logging_service", "monitoring_service", "alerting_service"];
for service in services {
println!(" đĄ Service '{}' would receive the broadcast event", service);
}
Ok(())
}
async fn topic_exchange_example() -> Result<(), Box<dyn std::error::Error>> {
let config = RabbitMQConfig {
connection_url: "amqp://guest:guest@localhost:5672/%2f".to_string(),
queue_name: "application_logs".to_string(),
exchange_name: "logs_topic".to_string(),
exchange_type: ExchangeType::Topic,
routing_key: Some("logs.*".to_string()), };
let queue = RabbitMQQueueSimple::new(config).await?;
let log_messages = vec![
(
"logs.auth.success",
serde_json::json!({
"level": "INFO",
"service": "auth",
"event": "login_success",
"user_id": 12345,
"ip": "192.168.1.100"
}),
),
(
"logs.api.error",
serde_json::json!({
"level": "ERROR",
"service": "api",
"event": "validation_failed",
"error_code": "INVALID_INPUT",
"endpoint": "/api/v1/users"
}),
),
(
"logs.db.warning",
serde_json::json!({
"level": "WARNING",
"service": "database",
"event": "slow_query",
"query_time_ms": 2500,
"table": "users"
}),
),
(
"logs.cache.info",
serde_json::json!({
"level": "INFO",
"service": "cache",
"event": "cache_hit",
"key": "user_profile_12345",
"hit_rate": 0.85
}),
),
];
for (routing_key, payload) in log_messages {
let log_message = QueueMessage::builder()
.payload(payload)
.expect("Failed to serialize log payload")
.routing_key(routing_key)
.build();
let message_id = queue.enqueue(log_message).await?;
println!(" â Log message: {} -> {}", routing_key, message_id);
}
println!(" đ Total log messages enqueued: {}", log_messages.len());
Ok(())
}
async fn priority_message_example() -> Result<(), Box<dyn std::error::Error>> {
let config = dev_config("priority_queue", "priority_exchange");
let queue = RabbitMQQueueSimple::new(config).await?;
let priority_messages = vec![
(
QueuePriority::Critical,
serde_json::json!({
"alert": "SYSTEM_DOWN",
"severity": "critical",
"action_required": "immediate"
}),
),
(
QueuePriority::High,
serde_json::json!({
"alert": "HIGH_CPU_USAGE",
"severity": "high",
"cpu_percent": 95.5,
"action_required": "monitor"
}),
),
(
QueuePriority::Normal,
serde_json::json!({
"alert": "USER_LOGIN",
"severity": "info",
"user_id": 67890
}),
),
(
QueuePriority::Low,
serde_json::json!({
"alert": "CLEANUP_TASK",
"severity": "low",
"task": "log_rotation"
}),
),
];
for (priority, payload) in priority_messages {
let message = QueueMessage::builder()
.payload(payload)
.expect("Failed to serialize priority payload")
.priority(priority)
.routing_key(format!("alerts.{}", priority.name().to_lowercase()))
.build();
let message_id = queue.enqueue(message).await?;
println!(" ⥠{} priority message: {}", priority.name(), message_id);
}
Ok(())
}
async fn batch_operations_example() -> Result<(), Box<dyn std::error::Error>> {
let config = dev_config("batch_processing", "batch_exchange");
let queue = RabbitMQQueueSimple::new(config).await?;
let mut batch_messages = Vec::new();
for i in 1..=100 {
let profile_update = QueueMessage::builder()
.payload(serde_json::json!({
"user_id": i,
"update_type": "profile_refresh",
"fields_updated": ["last_seen", "status"],
"timestamp": chrono::Utc::now().to_rfc3339()
}))
.expect("Failed to serialize profile update")
.priority(QueuePriority::Normal)
.routing_key("user.profile.update")
.build();
batch_messages.push(profile_update);
}
println!(" đĻ Preparing batch of {} profile updates", batch_messages.len());
let start_time = std::time::Instant::now();
let message_ids = queue.enqueue_batch(batch_messages).await?;
let processing_time = start_time.elapsed();
println!(" â Batch processed in {:?}", processing_time);
println!(" â Messages per second: {:.2}", message_ids.len() as f64 / processing_time.as_secs_f64());
println!(" â Generated {} message IDs", message_ids.len());
let batch_result = queue.dequeue_batch(10).await?;
println!(" â Retrieved batch: {} messages out of {} requested",
batch_result.messages.len(), batch_result.requested);
Ok(())
}
async fn message_with_headers_example() -> Result<(), Box<dyn std::error::Error>> {
let config = dev_config("routing_by_headers", "headers_exchange");
let queue = RabbitMQQueueSimple::new(config).await?;
let mut headers = HashMap::new();
headers.insert("source_service".to_string(), serde_json::Value::String("payment_service".to_string()));
headers.insert("correlation_id".to_string(), serde_json::Value::String("req_123456".to_string()));
headers.insert("retry_count".to_string(), serde_json::Value::Number(serde_json::Number::from(0)));
headers.insert("requires_ack".to_string(), serde_json::Value::Bool(true));
headers.insert("timeout_ms".to_string(), serde_json::Value::Number(serde_json::Number::from(5000)));
let mut message = QueueMessage::builder()
.payload(serde_json::json!({
"transaction_id": "txn_789012",
"amount": 99.99,
"currency": "USD",
"merchant_id": "merchant_456",
"payment_method": "credit_card"
}))
.expect("Failed to serialize payment payload")
.priority(QueuePriority::High)
.routing_key("payment.process");
let final_message = message.build();
let message_id = queue.enqueue(final_message).await?;
println!(" â Payment message with headers: {}", message_id);
println!(" đ Headers:");
for (key, value) in &headers {
println!(" {}: {}", key, value);
}
Ok(())
}
async fn validation_example() -> Result<(), Box<dyn std::error::Error>> {
println!(" đĄī¸ Testing configuration validation...");
let valid_config = dev_config("valid_queue", "valid_exchange");
match RabbitMQQueueSimple::new(valid_config).await {
Ok(queue) => println!(" â Valid configuration accepted"),
Err(e) => println!(" â Unexpected error with valid config: {}", e),
}
let invalid_config = RabbitMQConfig {
connection_url: "invalid://not-a-real-protocol".to_string(),
queue_name: "test".to_string(),
exchange_name: "test".to_string(),
exchange_type: ExchangeType::Direct,
routing_key: None,
};
match RabbitMQQueueSimple::new(invalid_config).await {
Ok(_) => println!(" â Invalid configuration was incorrectly accepted"),
Err(e) => println!(" â Invalid configuration correctly rejected: {}", e),
}
let empty_config = RabbitMQConfig {
connection_url: "".to_string(),
queue_name: "test".to_string(),
exchange_name: "test".to_string(),
exchange_type: ExchangeType::Direct,
routing_key: None,
};
match RabbitMQQueueSimple::new(empty_config).await {
Ok(_) => println!(" â Empty URL was incorrectly accepted"),
Err(e) => println!(" â Empty URL correctly rejected: {}", e),
}
let prod_config = prod_config(
"amqps://user:pass@rabbitmq.example.com:5671/%2f",
"production_queue",
"production_exchange",
ExchangeType::Topic,
);
match RabbitMQQueueSimple::new(prod_config).await {
Ok(queue) => println!(" â Production TLS configuration accepted"),
Err(e) => println!(" â ī¸ Production config validation (connection test would fail): {}", e),
}
Ok(())
}
async fn health_monitoring_example() -> Result<(), Box<dyn std::error::Error>> {
let config = dev_config("health_check_queue", "health_exchange");
let queue = RabbitMQQueueSimple::new(config).await?;
println!(" đ Performing health check...");
let health = queue.health_check().await?;
println!(" â Health Status: {:?}", health.status);
println!(" â Queue Size: {}", health.queue_size);
println!(" â Error Rate: {:.2}%", health.error_rate * 100.0);
println!(" â Last Check: {}", health.checked_at.format("%Y-%m-%d %H:%M:%S UTC"));
println!("\n đ Queue Statistics:");
let stats = queue.get_stats().await?;
println!(" â Total Messages: {}", stats.total_messages);
println!(" â Visible Messages: {}", stats.visible_messages);
println!(" â Invisible Messages: {}", stats.invisible_messages);
println!(" â Total Processed: {}", stats.total_processed);
println!(" â Total Failed: {}", stats.total_failed);
println!(" â Success Rate: {:.2}%", stats.success_rate() * 100.0);
println!("\n đ Testing connection...");
let connection_ok = queue.test_connection().await?;
println!(" â Connection Test: {}", if connection_ok { "â
PASS" } else { "â FAIL" });
println!("\n âī¸ Validating configuration...");
let config_ok = queue.validate_config().await?;
println!(" â Configuration Validation: {}", if config_ok { "â
PASS" } else { "â FAIL" });
println!("\n đ Queue Size Information:");
let size = queue.size().await?;
let is_empty = queue.is_empty().await?;
println!(" â Current Size: {} messages", size);
println!(" â Is Empty: {}", if is_empty { "â
YES" } else { "â NO" });
Ok(())
}
async fn simulate_consumer(queue_name: &str) -> Result<(), Box<dyn std::error::Error>> {
println!(" đ¤ Simulating consumer for queue: {}", queue_name);
let config = dev_config(queue_name, "consumer_exchange");
let queue = RabbitMQQueueSimple::new(config).await?;
let mut processed_count = 0;
let max_messages = 5;
while processed_count < max_messages {
match queue.dequeue().await {
Ok(Some(message)) => {
println!(" â Processed message: {} - {}", message.id, message.payload);
queue.ack(&message.id).await?;
processed_count += 1;
}
Ok(None) => {
println!(" âšī¸ No more messages in queue");
break;
}
Err(e) => {
println!(" â Error dequeuing message: {}", e);
break;
}
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
println!(" đ Consumer processed {} messages", processed_count);
Ok(())
}
async fn microservices_communication_example() -> Result<(), Box<dyn std::error::Error>> {
println!("đī¸ Microservices Communication Pattern");
println!("=====================================");
let user_service = RabbitMQQueueSimple::new(
RabbitMQConfig {
connection_url: "amqp://guest:guest@localhost:5672/%2f".to_string(),
queue_name: "user_events".to_string(),
exchange_name: "domain_events".to_string(),
exchange_type: ExchangeType::Topic,
routing_key: Some("user.created".to_string()),
}
).await?;
let user_event = QueueMessage::builder()
.payload(serde_json::json!({
"event_type": "user.created",
"user_id": 98765,
"email": "newuser@example.com",
"timestamp": chrono::Utc::now().to_rfc3339(),
"metadata": {
"source": "user_service",
"version": "1.0"
}
}))
.expect("Failed to serialize user event")
.routing_key("user.created")
.build();
let event_id = user_service.enqueue(user_event).await?;
println!(" đ¤ User Service published event: {}", event_id);
let notification_service = RabbitMQQueueSimple::new(
RabbitMQConfig {
connection_url: "amqp://guest:guest@localhost:5672/%2f".to_string(),
queue_name: "notifications_queue".to_string(),
exchange_name: "domain_events".to_string(),
exchange_type: ExchangeType::Topic,
routing_key: Some("user.*".to_string()), }
).await?;
let analytics_service = RabbitMQQueueSimple::new(
RabbitMQConfig {
connection_url: "amqp://guest:guest@localhost:5672/%2f".to_string(),
queue_name: "analytics_queue".to_string(),
exchange_name: "domain_events".to_string(),
exchange_type: ExchangeType::Fanout,
routing_key: None, }
).await?;
println!(" đ§ Notification Service: Ready to process user events");
println!(" đ Analytics Service: Ready to process all events");
println!(" đ Event-driven architecture established!");
Ok(())
}