diff --git a/.gitignore b/.gitignore index 8094110..f935865 100644 --- a/.gitignore +++ b/.gitignore @@ -1,8 +1,10 @@ Cargo.lock target/ +.idea/ .vscode/ .*.swp +.!* # Using the `Dockerfile.template` in the root directory now. Dockerfile diff --git a/step13_testcontainers/code/.junie/guidelines.md b/step13_testcontainers/code/.junie/guidelines.md new file mode 100644 index 0000000..c756fa4 --- /dev/null +++ b/step13_testcontainers/code/.junie/guidelines.md @@ -0,0 +1,23 @@ +# High-level + +Never remove existing tests, or the logic from them -- your job is to make existing tests build and pass if they fail to. +Write no comments, keep the code self-descriptive. +Respect Rust formatting rules in `../../rustfmt.toml`. No need to bring this `rustfmt.toml` into this project's directory from `../../`. +Keep Rust edition at 2024 in `Cargo.toml`. + +# When Contributing + +Only make small, self-contained changes. +Make sure they are readable and understood in isolation -- with no comments, from the code alone! +Do not add doctests for the sake of adding doctests, although if a standalone function calls for a doctest -- do add it! +Feel free to use `cargo doc` to examine the docs, but do not add `--open` to `cargo doc`! +If the project is using `ntest`, add `#[ntest::timeout(5000)]` to new tests. Do not modify or add/remove timeouts for existing tests. + +# Before Committing + +Every source file including `Cargo.toml` should end with a newline. +Run `cargo fmt` to ensure the style is consistent. +Run `cargo test` to ensure all tests pass. +Run `cargo check` to verify compilation. +Remove unused dependencies. +Ensure no build warnings. diff --git a/step13_testcontainers/code/Cargo.toml b/step13_testcontainers/code/Cargo.toml new file mode 100644 index 0000000..f721319 --- /dev/null +++ b/step13_testcontainers/code/Cargo.toml @@ -0,0 +1,22 @@ +[package] +name = "testcontainers_example" +edition = "2024" + +[[bin]] +name = "redis_publisher" +path = "src/publish.rs" + +[[bin]] +name = "redis_subscriber" +path = "src/subscribe.rs" + +[dependencies] +redis = { version = "0.29.2", features = ["tokio-comp"] } +tokio = { version = "1", features = ["full"] } +anyhow = "1.0" +testcontainers = "0.15" +testcontainers-modules = { version = "0.3", features = ["redis"] } +futures-util = "0.3" + +[dev-dependencies] +ntest = "0.9.0" diff --git a/step13_testcontainers/code/src/lib.rs b/step13_testcontainers/code/src/lib.rs new file mode 100644 index 0000000..fa85af3 --- /dev/null +++ b/step13_testcontainers/code/src/lib.rs @@ -0,0 +1,482 @@ +use anyhow::Result; +use std::sync::LazyLock; +use std::time::Duration; +use tokio::sync::Mutex; + +const PUBLISH_SCRIPT: &str = r#" + local channel = ARGV[1] + local message = ARGV[2] + + local existing_id = redis.call('HGET', 'msg_content_to_id', message) + if existing_id then + return existing_id + end + + local id = redis.call('INCR', 'msg_counter') + local key = 'msg:' .. id + redis.call('SET', key, message) + redis.call('HSET', 'msg_content_to_id', message, tostring(id)) + + local stored_count = redis.call('HLEN', 'msg_content_to_id') + if stored_count > 25 then + local oldest_id = id - 25 + local oldest_key = 'msg:' .. oldest_id + local oldest_message = redis.call('GET', oldest_key) + if oldest_message then + redis.call('DEL', oldest_key) + redis.call('HDEL', 'msg_content_to_id', oldest_message) + end + end + + redis.call('PUBLISH', channel, message) + return tostring(id) +"#; + +static SCRIPT_SHA: LazyLock>> = LazyLock::new(|| Mutex::new(None)); + +pub async fn create_redis_client(redis_url: &str) -> Result { + Ok(redis::Client::open(redis_url)?) +} + +pub async fn publish_with_persistence( + connection: &mut redis::aio::MultiplexedConnection, channel: &str, message: &str, +) -> Result { + loop { + let sha = match { SCRIPT_SHA.lock().await.clone() } { + Some(cached_sha) => cached_sha, + None => { + let sha = redis::cmd("SCRIPT").arg("LOAD").arg(PUBLISH_SCRIPT).query_async::(connection).await?; + *SCRIPT_SHA.lock().await = Some(sha.clone()); + sha + } + }; + + match redis::cmd("EVALSHA").arg(&sha).arg(0).arg(channel).arg(message).query_async::(connection).await { + Ok(result) => return Ok(result), + Err(e) => { + if e.to_string().contains("NoScriptError") { + *SCRIPT_SHA.lock().await = None; + tokio::time::sleep(Duration::from_millis(500)).await; + continue; + } + return Err(e.into()); + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use anyhow::anyhow; + use futures_util::StreamExt; + use redis::AsyncCommands; + use testcontainers::{Container, clients}; + use testcontainers_modules::redis::Redis; + use tokio::time::timeout; + + fn get_redis_url(container: &Container) -> String { + let port = container.get_host_port_ipv4(6379); + format!("redis://127.0.0.1:{}", port) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_set_get() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url)?; + let mut con = client.get_multiplexed_async_connection().await?; + + let _: () = con.set("test_key", "test_value").await?; + let val: String = con.get("test_key").await?; + + assert_eq!(val, "test_value"); + println!("Redis set/get test passed! Got value: {}", val); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_pubsub_basic() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("test_channel").await?; + + let publish_task = tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(100)).await; + let _ = publish_with_persistence(&mut publisher, "test_channel", "hello world").await.unwrap(); + }); + + let mut stream = pubsub.on_message(); + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + assert_eq!(message.get_channel_name(), "test_channel"); + assert_eq!(message.get_payload::()?, "hello world"); + + publish_task.await?; + println!("Redis pub-sub basic test passed!"); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_pubsub_multiple_messages() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("multi_channel").await?; + + let messages = vec!["message1", "message2", "message3"]; + let expected_messages = messages.clone(); + + let publish_task = tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(100)).await; + for msg in messages { + let _ = publish_with_persistence(&mut publisher, "multi_channel", msg).await.unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + } + }); + + let mut received_messages = Vec::new(); + let mut stream = pubsub.on_message(); + for _ in 0..3 { + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + received_messages.push(message.get_payload::()?); + } + + publish_task.await?; + assert_eq!(received_messages, expected_messages); + println!("Redis pub-sub multiple messages test passed!"); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_pubsub_multiple_channels() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("channel1").await?; + pubsub.subscribe("channel2").await?; + + let publish_task = tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(100)).await; + let _ = publish_with_persistence(&mut publisher, "channel1", "message_ch1").await.unwrap(); + let _ = publish_with_persistence(&mut publisher, "channel2", "message_ch2").await.unwrap(); + }); + + let mut received_messages = Vec::new(); + let mut stream = pubsub.on_message(); + for _ in 0..2 { + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + received_messages.push((message.get_channel_name().to_string(), message.get_payload::()?)); + } + + publish_task.await?; + received_messages.sort_by(|a, b| a.0.cmp(&b.0)); + + assert_eq!(received_messages[0], ("channel1".to_string(), "message_ch1".to_string())); + assert_eq!(received_messages[1], ("channel2".to_string(), "message_ch2".to_string())); + println!("Redis pub-sub multiple channels test passed!"); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_pubsub_pattern_subscription() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.psubscribe("test_*").await?; + + let publish_task = tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(100)).await; + let _ = publish_with_persistence(&mut publisher, "test_pattern1", "pattern_message1").await.unwrap(); + let _ = publish_with_persistence(&mut publisher, "test_pattern2", "pattern_message2").await.unwrap(); + let _ = publish_with_persistence(&mut publisher, "other_channel", "should_not_receive").await.unwrap(); + }); + + let mut received_messages = Vec::new(); + let mut stream = pubsub.on_message(); + for _ in 0..2 { + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + received_messages.push((message.get_channel_name().to_string(), message.get_payload::()?)); + } + + publish_task.await?; + received_messages.sort_by(|a, b| a.0.cmp(&b.0)); + + assert_eq!(received_messages.len(), 2); + assert_eq!(received_messages[0], ("test_pattern1".to_string(), "pattern_message1".to_string())); + assert_eq!(received_messages[1], ("test_pattern2".to_string(), "pattern_message2".to_string())); + println!("Redis pub-sub pattern subscription test passed!"); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_pubsub_unsubscribe() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("unsub_channel").await?; + + let _ = publish_with_persistence(&mut publisher, "unsub_channel", "before_unsub").await?; + { + let mut stream = pubsub.on_message(); + let message = timeout(Duration::from_secs(2), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + assert_eq!(message.get_payload::()?, "before_unsub"); + } + + pubsub.unsubscribe("unsub_channel").await?; + + let _ = publish_with_persistence(&mut publisher, "unsub_channel", "after_unsub").await?; + + { + let mut stream = pubsub.on_message(); + let result = timeout(Duration::from_millis(500), stream.next()).await; + assert!(result.is_err(), "Should not receive message after unsubscribe"); + } + + println!("Redis pub-sub unsubscribe test passed!"); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_message_persistence() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("persist_channel").await?; + + let message_id = publish_with_persistence(&mut publisher, "persist_channel", "test_message").await?; + + let mut stream = pubsub.on_message(); + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + assert_eq!(message.get_payload::()?, "test_message"); + + let stored_message: String = publisher.get(format!("msg:{}", message_id)).await?; + assert_eq!(stored_message, "test_message"); + + println!("Redis message persistence test passed! Message ID: {}", message_id); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_message_idempotency() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("idempotent_channel").await?; + + let first_id = publish_with_persistence(&mut publisher, "idempotent_channel", "duplicate_message").await?; + let second_id = publish_with_persistence(&mut publisher, "idempotent_channel", "duplicate_message").await?; + let third_id = publish_with_persistence(&mut publisher, "idempotent_channel", "duplicate_message").await?; + + assert_eq!(first_id, second_id); + assert_eq!(second_id, third_id); + + let mut stream = pubsub.on_message(); + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + assert_eq!(message.get_payload::()?, "duplicate_message"); + + let result = timeout(Duration::from_millis(500), stream.next()).await; + assert!(result.is_err(), "Should not receive duplicate message"); + + let stored_message: String = publisher.get(format!("msg:{}", first_id)).await?; + assert_eq!(stored_message, "duplicate_message"); + + println!("Redis message idempotency test passed! All attempts returned ID: {}", first_id); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_different_messages_unique_ids() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url.clone())?; + let publisher_client = redis::Client::open(redis_url)?; + + let mut pubsub = client.get_async_pubsub().await?; + let mut publisher = publisher_client.get_multiplexed_async_connection().await?; + + pubsub.subscribe("unique_channel").await?; + + let id1 = publish_with_persistence(&mut publisher, "unique_channel", "message1").await?; + let id2 = publish_with_persistence(&mut publisher, "unique_channel", "message2").await?; + let id3 = publish_with_persistence(&mut publisher, "unique_channel", "message3").await?; + + assert_ne!(id1, id2); + assert_ne!(id2, id3); + assert_ne!(id1, id3); + + let mut received_messages = Vec::new(); + let mut stream = pubsub.on_message(); + for _ in 0..3 { + let message = timeout(Duration::from_secs(5), stream.next()).await?.ok_or_else(|| anyhow!("timeout"))?; + received_messages.push(message.get_payload::()?); + } + + assert_eq!(received_messages, vec!["message1", "message2", "message3"]); + + println!("Redis different messages unique IDs test passed! IDs: {}, {}, {}", id1, id2, id3); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_idempotency_with_eviction() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url)?; + let mut publisher = client.get_multiplexed_async_connection().await?; + + let test_message = "test_message_for_eviction"; + + let first_id = publish_with_persistence(&mut publisher, "idempotency_channel", test_message).await?; + let duplicate_id = publish_with_persistence(&mut publisher, "idempotency_channel", test_message).await?; + + assert_eq!(first_id, duplicate_id, "Duplicate message should return same ID"); + println!("Idempotency confirmed: message '{}' got ID {} both times", test_message, first_id); + + println!("Publishing 25 different messages to trigger eviction..."); + for i in 1..=25 { + let filler_message = format!("filler_message_{}", i); + publish_with_persistence(&mut publisher, "idempotency_channel", &filler_message).await?; + } + + let hash_size: usize = publisher.hlen("msg_content_to_id").await?; + assert_eq!(hash_size, 25, "Hash should contain exactly 25 entries after eviction"); + + let evicted_check: Option = publisher.hget("msg_content_to_id", test_message).await?; + assert!(evicted_check.is_none(), "Original message should be evicted from idempotency hash"); + + let new_id = publish_with_persistence(&mut publisher, "idempotency_channel", test_message).await?; + assert_ne!(first_id, new_id, "Evicted message should get new ID when republished"); + + println!("Eviction confirmed: message '{}' got new ID {} after eviction (was {})", test_message, new_id, first_id); + + let final_duplicate_id = publish_with_persistence(&mut publisher, "idempotency_channel", test_message).await?; + assert_eq!(new_id, final_duplicate_id, "Republished message should maintain idempotency"); + + println!("Redis idempotency with eviction test passed!"); + Ok(()) + } + + #[tokio::test] + #[ntest::timeout(30000)] + async fn test_redis_message_eviction_25_limit() -> Result<()> { + let docker = clients::Cli::default(); + let redis_container = docker.run(Redis::default()); + let redis_url = get_redis_url(&redis_container); + + let client = redis::Client::open(redis_url)?; + let mut publisher = client.get_multiplexed_async_connection().await?; + + println!("Publishing 30 messages to test 25-element limit..."); + let mut message_ids = Vec::new(); + + for i in 1..=30 { + let message = format!("test_message_{}", i); + let id = publish_with_persistence(&mut publisher, "eviction_channel", &message).await?; + message_ids.push((id.parse::()?, message)); + } + + let counter: i32 = publisher.get("msg_counter").await?; + assert_eq!(counter, 30, "Message counter should be 30"); + + let hash_size: usize = publisher.hlen("msg_content_to_id").await?; + assert_eq!(hash_size, 25, "Idempotency hash should contain exactly 25 entries"); + + println!("Verifying that only the last 25 messages are stored..."); + for (id, message) in &message_ids { + let key = format!("msg:{}", id); + let stored_message: Option = publisher.get(&key).await?; + + if *id <= 5 { + assert!(stored_message.is_none(), "Message {} should have been evicted", id); + } else { + assert_eq!(stored_message.as_ref(), Some(message), "Message {} should be stored", id); + } + } + + println!("Verifying idempotency hash contains only the last 25 messages..."); + let hash_contents: Vec<(String, String)> = publisher.hgetall("msg_content_to_id").await?; + assert_eq!(hash_contents.len(), 25, "Hash should contain exactly 25 entries"); + + for (stored_message, stored_id) in hash_contents { + let id = stored_id.parse::()?; + assert!(id > 5, "Only messages with ID > 5 should be in hash, found ID {}", id); + assert_eq!(stored_message, format!("test_message_{}", id), "Hash entry should match expected message"); + } + + println!("Verifying random access works for retained messages..."); + for id in 6..=30 { + let key = format!("msg:{}", id); + let stored_message: String = publisher.get(&key).await?; + assert_eq!(stored_message, format!("test_message_{}", id), "Random access failed for message {}", id); + } + + println!("Redis message eviction 25-limit test passed!"); + Ok(()) + } +} diff --git a/step13_testcontainers/code/src/publish.rs b/step13_testcontainers/code/src/publish.rs new file mode 100644 index 0000000..0c06489 --- /dev/null +++ b/step13_testcontainers/code/src/publish.rs @@ -0,0 +1,37 @@ +use anyhow::{Result, anyhow}; +use std::env; +use testcontainers_example::{create_redis_client, publish_with_persistence}; + +#[tokio::main] +async fn main() -> Result<()> { + let args: Vec = env::args().collect(); + + if args.len() < 2 { + eprintln!("Usage: {} [channel]", args[0]); + eprintln!("Environment variables:"); + eprintln!(" REDIS_URL - Redis connection URL (default: redis://127.0.0.1:6379)"); + std::process::exit(1); + } + + let message = &args[1]; + let channel = args.get(2).map(|s| s.as_str()).unwrap_or("default_channel"); + + let redis_url = env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); + + println!("Connecting to Redis at: {}", redis_url); + println!("Publishing message: '{}' to channel: '{}'", message, channel); + + let client = create_redis_client(&redis_url).await?; + let mut connection = client.get_multiplexed_async_connection().await?; + + match publish_with_persistence(&mut connection, channel, message).await { + Ok(message_id) => { + println!("Message published successfully! Message ID: {}", message_id); + Ok(()) + } + Err(e) => { + eprintln!("Failed to publish message: {}", e); + Err(anyhow!("Publication failed: {}", e)) + } + } +} diff --git a/step13_testcontainers/code/src/subscribe.rs b/step13_testcontainers/code/src/subscribe.rs new file mode 100644 index 0000000..18b27e1 --- /dev/null +++ b/step13_testcontainers/code/src/subscribe.rs @@ -0,0 +1,31 @@ +use anyhow::Result; +use futures_util::StreamExt; +use std::env; +use testcontainers_example::create_redis_client; + +#[tokio::main] +async fn main() -> Result<()> { + let args: Vec = env::args().collect(); + + let channel = args.get(1).map(|s| s.as_str()).unwrap_or("default_channel"); + + let redis_url = env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); + + println!("Connecting to Redis at: {}", redis_url); + println!("Subscribing to channel: '{}'", channel); + println!("Waiting for messages... (Press Ctrl+C to exit)"); + + let client = create_redis_client(&redis_url).await?; + let mut pubsub = client.get_async_pubsub().await?; + + pubsub.subscribe(channel).await?; + + let mut stream = pubsub.on_message(); + while let Some(message) = stream.next().await { + let channel_name = message.get_channel_name(); + let payload = message.get_payload::()?; + println!("[{}] {}", channel_name, payload); + } + + Ok(()) +} diff --git a/step13_testcontainers/run.sh b/step13_testcontainers/run.sh new file mode 100755 index 0000000..4446994 --- /dev/null +++ b/step13_testcontainers/run.sh @@ -0,0 +1,5 @@ +#!/bin/bash + +set -e + +(cd code; cargo test)