#![cfg(feature = "integration-tests")]
mod support;
use support::install_crypto_provider;
use camel_api::cache::{CacheEntry, ContentType};
use camel_config::CamelConfig;
use redis::AsyncCommands;
use std::path::Path;
use std::sync::Arc;
use std::sync::Mutex;
use std::time::Duration;
use testcontainers::ContainerAsync;
use testcontainers::GenericImage;
use testcontainers::core::{ContainerPort, WaitFor};
use testcontainers::runners::AsyncRunner;
const REDIS_IMAGE_TAG: &str = "7-alpine";
async fn own_redis() -> (ContainerAsync<GenericImage>, String) {
let container = GenericImage::new("redis", REDIS_IMAGE_TAG)
.with_exposed_port(ContainerPort::Tcp(6379))
.with_wait_for(WaitFor::message_on_stdout("Ready to accept connections"))
.start()
.await
.expect("Redis container failed to start");
let port = container
.get_host_port_ipv4(6379)
.await
.expect("Redis port not available");
(container, format!("redis://127.0.0.1:{port}"))
}
fn disk_offload_toml(url: &str, payload_dir: &str, stale_retention: &str, extra: &str) -> String {
format!(
r#"
[default.cache_repo]
backend = "redis"
url = "{url}"
payload = "disk"
payload_dir = "{payload_dir}"
stale_retention = "{stale_retention}"
{extra}
"#
)
}
fn load_camel_toml(toml_str: &str) -> CamelConfig {
let dir = tempfile::TempDir::new().expect("tempdir");
let path = dir.path().join("Camel.toml");
std::fs::write(&path, toml_str).expect("write Camel.toml");
CamelConfig::from_file(path.to_str().unwrap()).expect("Camel.toml loads")
}
fn shared_log_buffer() -> Arc<Mutex<Vec<u8>>> {
struct SharedWriter(Arc<Mutex<Vec<u8>>>);
impl std::io::Write for SharedWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
static BUFFER: Mutex<Option<Arc<Mutex<Vec<u8>>>>> = Mutex::new(None);
let mut slot = BUFFER.lock().unwrap();
if let Some(buf) = slot.as_ref() {
return Arc::clone(buf);
}
let buf: Arc<Mutex<Vec<u8>>> = Arc::new(Mutex::new(Vec::new()));
let writer = Arc::clone(&buf);
let _ = tracing_subscriber::fmt()
.with_env_filter(tracing_subscriber::EnvFilter::new("camel_config=warn"))
.with_ansi(false)
.with_writer(move || SharedWriter(Arc::clone(&writer)))
.try_init();
*slot = Some(Arc::clone(&buf));
Arc::clone(&buf)
}
async fn context_with_disk_offload(
url: &str,
dir: &Path,
stale_retention: &str,
extra: &str,
) -> camel_core::CamelContext {
install_crypto_provider();
shared_log_buffer();
let cfg = load_camel_toml(&disk_offload_toml(
url,
dir.to_str().unwrap(),
stale_retention,
extra,
));
CamelConfig::configure_context(&cfg)
.await
.expect("context builds with redis disk-offload cache_repo")
}
async fn raw_connection(url: &str) -> redis::aio::MultiplexedConnection {
let client = redis::Client::open(url.to_string()).expect("raw client opens");
client
.get_multiplexed_async_connection()
.await
.expect("raw connection established")
}
fn blob_names(dir: &Path) -> Vec<String> {
std::fs::read_dir(dir)
.expect("payload dir is readable")
.map(|entry| {
entry
.expect("dir entry readable")
.file_name()
.to_string_lossy()
.to_string()
})
.filter(|name| name.ends_with(".blob"))
.collect()
}
fn the_blob(dir: &Path) -> std::path::PathBuf {
let blobs = blob_names(dir);
assert_eq!(blobs.len(), 1, "expected exactly one blob, got {blobs:?}");
dir.join(&blobs[0])
}
fn cache_entry(bytes: Vec<u8>) -> CacheEntry {
CacheEntry {
bytes,
payload_path: None,
content_type: ContentType::Bytes,
expires_at: None,
}
}
#[tokio::test(flavor = "multi_thread")]
async fn redis_disk_offload_round_trip() {
let (_container, url) = own_redis().await;
let dir = tempfile::TempDir::new().expect("payload dir tempdir");
let ctx = context_with_disk_offload(&url, dir.path(), "30s", "").await;
let repo = ctx
.cache_repository("redis")
.expect("redis cache repository registered when payload = disk");
assert_eq!(repo.name(), "redis");
let payload = vec![0xD4u8; 50 * 1024];
repo.set(
"k",
cache_entry(payload.clone()),
Some(Duration::from_secs(60)),
)
.await
.expect("set succeeds");
let blobs = blob_names(dir.path());
assert_eq!(
blobs.len(),
1,
"exactly one offloaded blob in the payload dir"
);
let got = repo
.get("k")
.await
.expect("get succeeds")
.expect("entry is present");
assert_eq!(
got.bytes, payload,
"hydrated payload must equal the stored one"
);
let logs = shared_log_buffer().lock().unwrap().clone();
let logs = String::from_utf8_lossy(&logs);
assert!(
logs.contains("offloaded entries under"),
"startup WARN must mention offloaded entries, got: {logs}"
);
assert!(
logs.contains(dir.path().to_str().unwrap()),
"startup WARN must name the payload dir, got: {logs}"
);
let mut conn = raw_connection(&url).await;
let raw: String = conn
.get("camel:cache:redis:k")
.await
.expect("raw GET returns the index row");
let row: serde_json::Value = serde_json::from_str(&raw).expect("index row is JSON");
assert!(
row["bytes"].as_array().is_some_and(Vec::is_empty),
"index row must store an empty bytes array, got: {raw}"
);
let stored = row["payload_path"]
.as_str()
.expect("index row must carry payload_path");
assert_eq!(
stored, blobs[0],
"index row must name the blob present in the payload dir"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn redis_disk_offload_early_sweep_is_miss() {
let (_container, url) = own_redis().await;
let dir = tempfile::TempDir::new().expect("payload dir tempdir");
let ctx = context_with_disk_offload(&url, dir.path(), "30s", "").await;
let repo = ctx
.cache_repository("redis")
.expect("redis cache repository registered");
repo.set(
"k",
cache_entry(b"early-sweep-payload".to_vec()),
Some(Duration::from_secs(60)),
)
.await
.expect("set succeeds");
let blob = the_blob(dir.path());
std::fs::remove_file(&blob).expect("blob file deleted");
let got = repo.get("k").await.expect("get succeeds");
assert!(got.is_none(), "vanished blob must degrade get to a miss");
let stale = repo.peek_stale("k").await.expect("peek_stale succeeds");
assert!(
stale.is_none(),
"vanished blob must degrade peek_stale to a miss"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn redis_disk_offload_sweeper_reclaims_orphan() {
let (_container, url) = own_redis().await;
let dir = tempfile::TempDir::new().expect("payload dir tempdir");
let ctx =
context_with_disk_offload(&url, dir.path(), "1ms", "payload_sweep_interval = \"1s\"").await;
let repo = ctx
.cache_repository("redis")
.expect("redis cache repository registered");
repo.set(
"k",
cache_entry(b"sweeper-orphan".to_vec()),
Some(Duration::from_millis(100)),
)
.await
.expect("set succeeds");
let blob = the_blob(dir.path());
support::wait::wait_until(
"sweeper reclaims the dead payload blob",
Duration::from_secs(8),
Duration::from_millis(100),
|| async { Ok(!blob.exists()) },
)
.await
.expect("blob must be gone from the payload dir");
}
#[tokio::test(flavor = "multi_thread")]
async fn redis_disk_offload_no_ttl_capped() {
let (_container, url) = own_redis().await;
let dir = tempfile::TempDir::new().expect("payload dir tempdir");
let ctx = context_with_disk_offload(
&url,
dir.path(),
"1ms",
"payload_sweep_interval = \"1s\"\npayload_max_ttl = \"1s\"",
)
.await;
let repo = ctx
.cache_repository("redis")
.expect("redis cache repository registered");
repo.set("k", cache_entry(b"capped-payload".to_vec()), None)
.await
.expect("set without ttl succeeds");
let blob = the_blob(dir.path());
support::wait::wait_until(
"no-ttl entry expires at payload_max_ttl",
Duration::from_secs(5),
Duration::from_millis(100),
|| async {
repo.get("k")
.await
.map(|entry| entry.is_none())
.map_err(|e| e.to_string())
},
)
.await
.expect("no-ttl entry must expire at payload_max_ttl");
support::wait::wait_until(
"sweeper reclaims the capped payload blob",
Duration::from_secs(8),
Duration::from_millis(100),
|| async { Ok(!blob.exists()) },
)
.await
.expect("blob must be swept from the payload dir");
}