use backbone_queue::{
QueueService,
redis::RedisQueueBuilder,
fifo::{FifoQueueService, FifoQueueServiceWrapper, FifoQueueConfig, MessageVolume},
types::{QueueMessage, QueuePriority}
};
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use serde_json::json;
fn create_fifo_message(
id: &str,
group_id: &str,
deduplication_id: &str,
payload: serde_json::Value,
priority: QueuePriority,
) -> QueueMessage {
QueueMessage {
id: id.to_string(),
payload,
priority,
receive_count: 0,
max_receive_count: 3,
enqueued_at: chrono::Utc::now(),
visible_at: chrono::Utc::now(),
expires_at: None,
visibility_timeout: 30,
status: backbone_queue::MessageStatus::Pending,
delay_seconds: None,
attributes: std::collections::HashMap::new(),
message_group_id: Some(group_id.to_string()),
message_deduplication_id: Some(deduplication_id.to_string()),
compressed: false,
original_size: None,
}
}
fn create_order_message(order_id: &str, action: &str) -> QueueMessage {
create_fifo_message(
&format!("order-{}", order_id),
format!("order-{}", order_id), format!("order-{}-{}", order_id, action), json!({
"order_id": order_id,
"action": action,
"timestamp": chrono::Utc::now(),
"data": {
"customer_id": format!("customer-{}", order_id),
"amount": (order_id.parse::<i64>().unwrap_or(0) * 100),
"items": vec![
{"product": "widget-{}".format!(order_id), "quantity": order_id.parse::<i64>().unwrap_or(1) % 10 + 1)},
{"product": "gadget-{}",format!(order_id), "quantity": order_id.parse::<i64>().unwrap_or(1) % 5 + 1)}
]
}
}),
QueuePriority::Normal,
)
}
fn create_activity_message(user_id: &str, activity: &str) -> QueueMessage {
create_fifo_message(
&format!("activity-{}", user_id),
format!("user-{}", user_id), format!("activity-{}-{}", user_id, chrono::Utc::now().timestamp()),
json!({
"user_id": user_id,
"activity": activity,
"timestamp": chrono::Utc::now(),
"metadata": {
"source": "mobile_app",
"version": "1.0.0"
}
}),
QueuePriority::Low,
)
}
async fn demo_basic_fifo_operations() -> Result<(), Box<dyn std::error::Error>> {
println!("๐ Basic FIFO Queue Operations");
println!("=============================");
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("fifo_test_queue")
.key_prefix("fifo")
.build()
.await?;
if !queue.test_connection().await? {
println!("โ Failed to connect to Redis, skipping demo");
return Ok(());
}
queue.purge().await?;
let config = backbone_queue::fifo::utils::get_recommended_config(MessageVolume::Medium);
let fifo_service = FifoQueueServiceWrapper::new(
Arc::new(queue),
config,
);
println!("\n๐ค Enqueuing FIFO messages...");
let messages = vec![
create_order_message("1001", "created"),
create_order_message("1002", "created"),
create_order_message("1003", "created"),
create_order_message("1001", "updated"),
create_order_message("1002", "updated"),
];
for (i, message) in messages.iter().enumerate() {
match fifo_service.enqueue_fifo(message.clone()).await {
Ok(id) => {
println!(" โ
Message {} enqueued: {} (Group: {})",
i + 1,
message.message_deduplication_id.as_ref().unwrap(),
message.message_group_id.as_ref().unwrap()
);
}
Err(e) => {
println!(" โ Failed to enqueue message {}: {}", i + 1, e);
}
}
let stats = fifo_service.get_fifo_stats().await?;
println!("\n๐ FIFO Queue Statistics:");
println!(" Total groups: {}", stats.total_groups);
println!(" Active groups: {}", stats.active_groups);
println!(" Deduplicated messages: {}", stats.deduplicated_messages);
println!(" Deduplication cache size: {}", stats.deduplication_cache.len());
Ok(())
}
async fn demo_message_deduplication() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐ซ Message Deduplication Demo");
println!("=============================");
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("dedup_test_queue")
.key_prefix("dedup")
.build()
.await?;
if !queue.test_connection().await? {
println!("โ Failed to connect to Redis, skipping demo");
return Ok(());
}
queue.purge().await?;
let config = FifoQueueConfig {
enabled: true,
deduplication_window_seconds: 300, max_message_groups: 1000,
enable_content_deduplication: true, content_deduplication_window_seconds: 60, max_deduplicated_messages: 10000,
};
let fifo_service = FifoQueueServiceWrapper::new(
Arc::new(queue),
config,
);
println!("\n๐ Testing deduplication IDs...");
let message1 = create_fifo_message(
"test-1",
"group-1",
"same-dedup-id",
json!({"action": "update", "data": "value1"}),
QueuePriority::Normal,
);
let message2 = create_fifo_message(
"test-2",
"group-1",
"same-dedup-id", json!({"action": "update", "data": "value2"}),
QueuePriority::Normal,
);
match fifo_service.enqueue_fifo(message1).await {
Ok(id) => println!(" โ
First message enqueued: {}", id),
Err(e) => println!(" โ Failed to enqueue first message: {}", e),
}
match fifo_service.enqueue_fifo(message2).await {
Ok(_) => println!(" โ ๏ธ Second message unexpectedly enqueued"),
Err(e) => println!(" โ
Second message correctly rejected: {}", e),
}
println!("\n๐ Testing content deduplication...");
let content = json!({
"order_id": "1234",
"amount": 99.99,
"items": ["item1", "item2"]
});
let message3 = create_fifo_message(
"test-3",
"group-2",
"content-dedup-1",
content.clone(),
QueuePriority::Normal,
);
let message4 = create_fifo_message(
"test-4",
"group-2",
"content-dedup-2", content, QueuePriority::Normal,
);
match fifo_service.enqueue_fifo(message3).await {
Ok(id) => println!(" โ
First content message enqueued: {}", id),
Err(e) => println!(" โ Failed to enqueue first content message: {}", e),
}
match fifo_service.enqueue_fifo(message4).await {
Ok(_) => println!(" โ ๏ธ Second content message unexpectedly enqueued"),
Err(e) => println!(" โ
Second content message correctly rejected: {}", e),
}
let stats = fifo_service.get_fifo_stats().await?;
println!("\n๐ Deduplication Statistics:");
println!(" Message deduplication cache: {} entries", stats.deduplication_cache.len());
println!(" Content deduplication cache: {} entries", stats.content_deduplication_cache.len());
Ok(())
}
async fn demo_group_statistics() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐ Message Group Statistics Demo");
println!("===============================");
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("stats_test_queue")
.key_prefix("stats")
.build()
.await?;
if !queue.test_connection().await? {
println!("โ Failed to connect to Redis, skipping demo");
return Ok(());
}
queue.purge().await?;
let fifo_service = FifoQueueServiceWrapper::new(
Arc::new(queue),
FifoQueueConfig::default(),
);
println!("\n๐ค Creating messages for different groups...");
let groups = vec![
("payments", 3),
("notifications", 5),
("analytics", 2),
("backup", 1),
];
for (group_name, count) in groups {
println!(" ๐ฆ Group '{}': {} messages", group_name, count);
for i in 1..=count {
let message = create_fifo_message(
&format!("{}-{}", group_name, i),
group_name,
&format!("{}-{}-{}", group_name, i, chrono::Utc::now().timestamp()),
json!({
"group": group_name,
"index": i,
"data": "test data for group processing"
}),
QueuePriority::Normal,
);
match fifo_service.enqueue_fifo(message).await {
Ok(_) => {}, Err(e) => println!(" โ Failed: {}", e),
}
}
}
println!("\nโ๏ธ Simulating message processing...");
let group_ids: Vec<String> = groups.iter().map(|(name, _)| name.to_string()).collect();
for group_id in group_ids {
let mut group_stats = fifo_service.get_group_stats(&group_id).await?.unwrap_or_default();
for _ in 0..group_stats.message_count / 2 {
{
let mut stats = fifo_service.stats.write().await;
stats.update_group(&group_id, true, 100 + (rand::random::<u64>() % 200));
}
}
for _ in 0..group_stats.message_count / 4 {
{
let mut stats = fifo_service.stats.write().await;
stats.update_group(&group_id, false, 300 + (rand::random::<u64>() % 200));
}
}
}
println!("\n๐ Message Group Statistics:");
let all_stats = fifo_service.get_all_group_stats().await?;
for (group_id, stats) in all_stats {
println!(" ๐ {}: {} messages, {:.1}% success rate, {:.1}ms avg processing time",
group_id,
stats.message_count,
stats.success_rate(),
stats.avg_processing_time_ms
);
}
Ok(())
}
async fn demo_cleanup_operations() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐งน Cleanup Operations Demo");
println!("==========================");
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("cleanup_test_queue")
.key_prefix("cleanup")
.build()
.await?;
if !queue.test_connection().await? {
println!("โ Failed to connect to Redis, skipping demo");
return Ok(());
}
queue.purge().await?;
let fifo_service = FifoQueueServiceWrapper::new(
Arc::new(queue),
FifoQueueConfig {
enabled: true,
deduplication_window_seconds: 2, max_message_groups: 100,
enable_content_deduplication: true,
content_deduplication_window_seconds: 2,
max_deduplicated_messages: 100,
},
);
println!("\n๐ Adding entries that will expire...");
{
let mut stats = fifo_service.stats.write().await;
let past_time = chrono::Utc::now() - chrono::Duration::seconds(10);
for i in 0..5 {
stats.deduplication_cache.insert(
format!("expired-dedup-{}", i),
past_time,
);
}
for i in 0..3 {
stats.content_deduplication_cache.insert(
format!("expired-content-{}", i),
past_time,
);
}
println!(" ๐๏ธ Added 5 expired message deduplication entries");
println!(" ๐๏ธ Added 3 expired content deduplication entries");
}
let stats_before = fifo_service.get_fifo_stats().await?;
println!("\n๐ Before cleanup:");
println!(" Deduplication cache: {} entries", stats_before.deduplication_cache.len());
println!(" Content deduplication cache: {} entries", stats_before.content_deduplication_cache.len());
println!("\n๐งน Performing cleanup...");
let removed_count = fifo_service.cleanup_deduplication().await?;
let stats_after = fifo_service.get_fifo_stats().await?;
println!("\n๐ After cleanup:");
println!(" Removed {} expired entries", removed_count);
println!(" Deduplication cache: {} entries", stats_after.deduplication_cache.len());
println!(" Content deduplication cache: {} entries", stats_after.content_deduplication_cache.len());
Ok(())
}
async fn demo_configuration_recommendations() -> Result<(), Box<dyn std::error::Error>> {
println!("\nโ๏ธ Configuration Recommendations");
println!("===============================");
let volumes = vec![
(MessageVolume::Low, "Small application (100-1K messages/day)"),
(MessageVolume::Medium, "Medium application (1K-10K messages/day)"),
(MessageVolume::High, "Large application (10K+ messages/day)"),
];
for (volume, description) in volumes {
let config = backbone_queue::fifo::utils::get_recommended_config(volume);
let errors = backbone_queue::fifo::utils::validate_config(&config);
println!("\n๐ {} Volume Recommendation:", volume as i32);
println!(" ๐ Description: {}", description);
println!(" โฐ Deduplication window: {} seconds", config.deduplication_window_seconds);
println!(" ๐ Max message groups: {}", config.max_message_groups);
println!(" ๐ Content deduplication: {}", config.enable_content_deduplication);
if config.enable_content_deduplication {
println!(" โฐ Content deduplication window: {} seconds", config.content_deduplication_window_seconds);
}
if !errors.is_empty() {
println!(" โ ๏ธ Validation issues:");
for error in errors {
println!(" โ {}", error);
}
} else {
println!(" โ
Configuration is valid");
}
}
Ok(())
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
env_logger::init();
println!("๐ FIFO Queue Demo");
println!("===================");
demo_basic_fifo_operations().await?;
demo_message_deduplication().await?;
demo_group_statistics().await?;
demo_cleanup_operations().await?;
demo_configuration_recommendations().await?;
println!("\n๐ FIFO queue demo completed!");
println!("\n๐ Key Takeaways:");
println!(" โข FIFO ensures exact message ordering within groups");
println!(" โข Message deduplication prevents duplicate processing");
println!(" โข Content deduplication provides additional safety");
println!(" โข Message groups allow parallel processing of independent streams");
println!(" โข Group statistics help monitor individual processing flows");
println!(" โข Automatic cleanup prevents memory leaks");
println!(" โข Configuration recommendations optimize for different volumes");
Ok(())
}