use backbone_queue::{
QueueService,
redis::RedisQueueBuilder,
compression::{MessageCompressor, CompressionConfig, CompressionAlgorithm, CompressedMessageBuilder},
types::{QueueMessage, QueuePriority}
};
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use serde_json::json;
fn create_large_payload(size_mb: usize) -> serde_json::Value {
let base_data = "Lorem ipsum dolor sit amet, consectetur adipiscing elit. Sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat. ";
let mut data = String::new();
let target_size = size_mb * 1024 * 1024; while data.len() < target_size {
data.push_str(base_data);
}
json!({
"id": "large-document-123",
"title": "Large Document",
"content": data,
"metadata": {
"size_mb": size_mb,
"type": "document",
"compression_test": true,
"created_at": chrono::Utc::now()
},
"attachments": vec![
{
"name": "file1.pdf",
"size": 2048576,
"type": "application/pdf"
},
{
"name": "image1.png",
"size": 1024000,
"type": "image/png"
}
]
})
}
fn create_small_payload() -> serde_json::Value {
json!({
"id": "small-message-456",
"text": "Hello, World!",
"type": "notification",
"timestamp": chrono::Utc::now()
})
}
async fn benchmark_compression() -> Result<(), Box<dyn std::error::Error>> {
println!("๐ง Compression Benchmark");
println!("=======================");
let mut results = HashMap::new();
let sizes = vec![1, 5, 10, 25];
for size_mb in sizes {
println!("\n๐ Testing {}MB payload:", size_mb);
let payload = create_large_payload(size_mb);
let raw_size = serde_json::to_vec(&payload).unwrap().len();
{
let config = CompressionConfig {
algorithm: CompressionAlgorithm::Gzip,
level: 6,
min_size: 0,
force_compression: true,
max_attempts: 3,
};
let compressor = MessageCompressor::new(config);
let mut message = QueueMessage {
id: format!("gzip-test-{}mb", size_mb),
payload: payload.clone(),
priority: QueuePriority::Normal,
created_at: chrono::Utc::now(),
expires_at: None,
receive_count: 0,
max_receive_count: 3,
visibility_timeout: 30,
delay_until: None,
metadata: HashMap::new(),
};
let start = Instant::now();
compressor.compress_message(&mut message).await?;
let compression_time = start.elapsed();
let compressed_size = message.metadata["compression_compressed_size"]
.as_u64().unwrap() as usize;
let compression_ratio = compressed_size as f64 / raw_size as f64;
println!(" ๐ฆ GZIP: {} -> {} bytes ({:.1}% compression) in {:?}",
raw_size, compressed_size, compression_ratio * 100.0, compression_time);
results.insert(format!("gzip_{}mb", size_mb),
(compressed_size, compression_ratio, compression_time));
}
{
let config = CompressionConfig {
algorithm: CompressionAlgorithm::Zlib,
level: 6,
min_size: 0,
force_compression: true,
max_attempts: 3,
};
let compressor = MessageCompressor::new(config);
let mut message = QueueMessage {
id: format!("zlib-test-{}mb", size_mb),
payload: payload.clone(),
priority: QueuePriority::Normal,
created_at: chrono::Utc::now(),
expires_at: None,
receive_count: 0,
max_receive_count: 3,
visibility_timeout: 30,
delay_until: None,
metadata: HashMap::new(),
};
let start = Instant::now();
compressor.compress_message(&mut message).await?;
let compression_time = start.elapsed();
let compressed_size = message.metadata["compression_compressed_size"]
.as_u64().unwrap() as usize;
let compression_ratio = compressed_size as f64 / raw_size as f64;
println!(" ๐ฆ ZLIB: {} -> {} bytes ({:.1}% compression) in {:?}",
raw_size, compressed_size, compression_ratio * 100.0, compression_time);
results.insert(format!("zlib_{}mb", size_mb),
(compressed_size, compression_ratio, compression_time));
}
println!(" ๐ง GZIP Compression Levels:");
for level in [1, 3, 6, 9] {
let config = CompressionConfig {
algorithm: CompressionAlgorithm::Gzip,
level,
min_size: 0,
force_compression: true,
max_attempts: 3,
};
let compressor = MessageCompressor::new(config);
let mut message = QueueMessage {
id: format!("gzip-lvl{}-{}mb", level, size_mb),
payload: payload.clone(),
priority: QueuePriority::Normal,
created_at: chrono::Utc::now(),
expires_at: None,
receive_count: 0,
max_receive_count: 3,
visibility_timeout: 30,
delay_until: None,
metadata: HashMap::new(),
};
let start = Instant::now();
compressor.compress_message(&mut message).await?;
let compression_time = start.elapsed();
let compressed_size = message.metadata["compression_compressed_size"]
.as_u64().unwrap() as usize;
let compression_ratio = compressed_size as f64 / raw_size as f64;
println!(" Level {}: {} bytes ({:.1}% ratio) in {:?}",
level, compressed_size, compression_ratio * 100.0, compression_time);
}
}
println!("\n๐ Best Compression Results:");
for size_mb in sizes {
let gzip_key = format!("gzip_{}mb", size_mb);
let zlib_key = format!("zlib_{}mb", size_mb);
if let (Some(gzip_result), Some(zlib_result)) = (results.get(&gzip_key), results.get(&zlib_key)) {
let gzip_ratio = gzip_result.1;
let zlib_ratio = zlib_result.1;
if gzip_ratio < zlib_ratio {
println!(" {}MB: GZIP wins ({:.1}% vs {:.1}%)",
size_mb, gzip_ratio * 100.0, zlib_ratio * 100.0);
} else {
println!(" {}MB: ZLIB wins ({:.1}% vs {:.1}%)",
size_mb, zlib_ratio * 100.0, gzip_ratio * 100.0);
}
}
}
Ok(())
}
async fn test_automatic_compression() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐ค Automatic Compression Test");
println!("===============================");
let config = CompressionConfig {
algorithm: CompressionAlgorithm::Gzip,
level: 6,
min_size: 1024, force_compression: false,
max_attempts: 3,
};
let compressor = MessageCompressor::new(config);
{
let mut small_message = QueueMessage {
id: "small-test".to_string(),
payload: create_small_payload(),
priority: QueuePriority::Normal,
created_at: chrono::Utc::now(),
expires_at: None,
receive_count: 0,
max_receive_count: 3,
visibility_timeout: 30,
delay_until: None,
metadata: HashMap::new(),
};
let original_payload = small_message.payload.clone();
compressor.compress_message(&mut small_message).await?;
println!("๐จ Small message: {} -> {} bytes",
serde_json::to_vec(&original_payload).unwrap().len(),
serde_json::to_vec(&small_message.payload).unwrap().len());
if small_message.metadata.contains_key("compression_algorithm") {
println!(" โ Unexpectedly compressed");
} else {
println!(" โ
Correctly not compressed (below threshold)");
}
}
{
let mut large_message = QueueMessage {
id: "large-test".to_string(),
payload: create_large_payload(2),
priority: QueuePriority::Normal,
created_at: chrono::Utc::now(),
expires_at: None,
receive_count: 0,
max_receive_count: 3,
visibility_timeout: 30,
delay_until: None,
metadata: HashMap::new(),
};
let original_size = serde_json::to_vec(&large_message.payload).unwrap().len();
compressor.compress_message(&mut large_message).await?;
let compressed_size = serde_json::to_vec(&large_message.payload).unwrap().len();
println!("๐ฆ Large message: {} -> {} bytes", original_size, compressed_size);
if large_message.metadata.contains_key("compression_algorithm") {
let ratio = compressed_size as f64 / original_size as f64;
println!(" โ
Correctly compressed ({:.1}% of original)", ratio * 100.0);
} else {
println!(" โ Unexpectedly not compressed");
}
}
Ok(())
}
async fn test_compression_roundtrip() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐ Compression Roundtrip Test");
println!("=============================");
let compressor = Arc::new(MessageCompressor::default());
let test_cases = vec![
("small_json", create_small_payload()),
("large_text", create_large_payload(1)),
("mixed_data", json!({
"text": "Hello, World!".repeat(1000),
"numbers": vec![1i32; 10000],
"nested": {
"data": "test".repeat(5000),
"array": vec![json!({"item": "value"}); 1000]
}
})),
];
for (name, payload) in test_cases {
println!("\n๐ Testing case: {}", name);
let builder = CompressedMessageBuilder::new(compressor.clone())
.id(format!("roundtrip-{}", name))
.payload(payload.clone())
.priority(QueuePriority::High)
.metadata("test_case".to_string(), json!(name));
let mut compressed_message = builder.build().await?;
println!(" โ
Message compressed");
assert!(compressed_message.metadata.contains_key("compression_algorithm"));
assert!(compressed_message.metadata.contains_key("compression_original_size"));
assert!(compressed_message.metadata.contains_key("compression_compressed_size"));
let original_size = compressed_message.metadata["compression_original_size"]
.as_u64().unwrap() as usize;
let compressed_size = compressed_message.metadata["compression_compressed_size"]
.as_u64().unwrap() as usize;
println!(" ๐ Size: {} -> {} bytes ({:.1}% ratio)",
original_size, compressed_size, (compressed_size as f64 / original_size as f64) * 100.0);
compressor.decompress_message(&mut compressed_message).await?;
println!(" โ
Message decompressed");
let decompressed_payload = compressed_message.payload;
if decompressed_payload == payload {
println!(" โ
Data integrity verified");
} else {
println!(" โ Data integrity check failed");
return Err("Data integrity check failed".into());
}
assert!(!compressed_message.metadata.contains_key("compression_algorithm"));
}
Ok(())
}
async fn test_compression_stats() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐ Compression Statistics Test");
println!("=============================");
let compressor = MessageCompressor::default();
let mut messages = Vec::new();
for i in 0..10 {
let mut message = QueueMessage {
id: format!("stats-test-{}", i),
payload: create_large_payload(1), priority: QueuePriority::Normal,
created_at: chrono::Utc::now(),
expires_at: None,
receive_count: 0,
max_receive_count: 3,
visibility_timeout: 30,
delay_until: None,
metadata: HashMap::new(),
};
compressor.compress_message(&mut message).await?;
messages.push(message);
}
let stats = compressor.get_stats().await;
println!("๐ Compression Statistics:");
println!(" Total messages: {}", stats.total_messages);
println!(" Compressed messages: {}", stats.compressed_messages);
println!(" Original bytes: {}", stats.original_bytes);
println!(" Compressed bytes: {}", stats.compressed_bytes);
println!(" Compression ratio: {:.3}", stats.compression_ratio);
println!(" Avg compression time: {:.2}ms", stats.avg_compression_time_ms());
let space_saved = stats.original_bytes - stats.compressed_bytes;
println!(" Space saved: {} bytes ({:.1}%)",
space_saved, (space_saved as f64 / stats.original_bytes as f64) * 100.0);
for message in messages.iter().take(5) {
let mut message_clone = message.clone();
compressor.decompress_message(&mut message_clone).await?;
}
let final_stats = compressor.get_stats().await;
println!(" Decompressed messages: {}", final_stats.decompressed_messages);
println!(" Avg decompression time: {:.2}ms", final_stats.avg_decompression_time_ms());
Ok(())
}
async fn test_redis_integration() -> Result<(), Box<dyn std::error::Error>> {
println!("\n๐ Redis Integration Test");
println!("=========================");
let queue = RedisQueueBuilder::new()
.url("redis://localhost:6379")
.queue_name("compression_test_queue")
.key_prefix("compression")
.build()
.await?;
if !queue.test_connection().await? {
println!("โ Failed to connect to Redis, skipping integration test");
return Ok(());
}
queue.purge().await?;
let compressor = Arc::new(MessageCompressor::default());
let large_payload = create_large_payload(2);
let builder = CompressedMessageBuilder::new(compressor.clone())
.id("redis-compression-test".to_string())
.payload(large_payload)
.priority(QueuePriority::High);
let compressed_message = builder.build().await?;
println!("โ
Created compressed message");
let message_id = queue.enqueue(compressed_message).await?;
println!("โ
Enqueued compressed message: {}", message_id);
if let Some(mut message) = queue.dequeue().await? {
println!("โ
Dequeued message");
if message.metadata.contains_key("compression_algorithm") {
println!("๐ฆ Message is compressed, decompressing...");
compressor.decompress_message(&mut message).await?;
println!("โ
Message decompressed successfully");
if let Some(content) = message.payload.get("content") {
if content.is_string() && content.as_str().unwrap().len() > 1000000 {
println!("โ
Large content restored ({} chars)", content.as_str().unwrap().len());
}
}
}
queue.ack(&message.id).await?;
println!("โ
Message acknowledged");
}
queue.purge().await?;
println!("๐งน Queue cleaned up");
Ok(())
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
env_logger::init();
println!("๐ Message Compression Demo");
println!("==========================");
benchmark_compression().await?;
test_automatic_compression().await?;
test_compression_roundtrip().await?;
test_compression_stats().await?;
test_redis_integration().await?;
println!("\n๐ Compression demo completed!");
println!("\n๐ Key Takeaways:");
println!(" โข Compression is most effective for repetitive/large data");
println!(" โข GZIP provides good balance of speed and ratio");
println!(" โข ZLIB offers slightly better compression for some data");
println!(" โข Automatic thresholds prevent unnecessary compression");
println!(" โข Statistics help monitor compression efficiency");
println!(" โข Integration with queue operations is seamless");
Ok(())
}