use std::process::Command;
use std::sync::Arc;
use std::time::Duration;
use oxcache::backend::interface::{BackendKind, CacheConnector, CacheReader, CacheWriter};
use oxcache::backend::{AerospikeBackend, AerospikeConfig};
use tokio::sync::OnceCell;
const CONTAINER_NAME: &str = "oxcache-as-integration";
const HOST_PORT: u16 = 3001;
static SHARED_CONFIG: OnceCell<Option<AerospikeConfig>> = OnceCell::const_new();
async fn shared_config() -> Option<&'static AerospikeConfig> {
SHARED_CONFIG.get_or_init(start_container).await.as_ref()
}
async fn start_container() -> Option<AerospikeConfig> {
let _ = Command::new("docker").args(["rm", "-f", CONTAINER_NAME]).output();
let output = Command::new("docker")
.args([
"run",
"-d",
"--name",
CONTAINER_NAME,
"-p",
&format!("{HOST_PORT}:3000"),
"aerospike/aerospike-server:8.0",
])
.output()
.ok()?;
if !output.status.success() {
eprintln!("skip: docker run failed: {}", String::from_utf8_lossy(&output.stderr));
return None;
}
let start = std::time::Instant::now();
let timeout = Duration::from_secs(60);
let mut ready = false;
while start.elapsed() < timeout {
let logs = Command::new("docker").args(["logs", CONTAINER_NAME]).output().ok()?;
let stderr = String::from_utf8_lossy(&logs.stderr);
if stderr.contains("migrations: complete") || stderr.contains("service ready") {
ready = true;
break;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
if !ready {
eprintln!("skip: Aerospike container start timeout");
return None;
}
let _ = Command::new("docker")
.args([
"exec", CONTAINER_NAME, "sh", "-c",
&format!(
"sed -i 's|# access-address <IPADDR>|access-address 127.0.0.1\\n\\t\\taccess-port {HOST_PORT}|' /etc/aerospike/aerospike.conf"
),
])
.output();
let _ = Command::new("docker").args(["stop", CONTAINER_NAME]).output();
let _ = Command::new("docker").args(["start", CONTAINER_NAME]).output();
tokio::time::sleep(Duration::from_secs(5)).await;
let start = std::time::Instant::now();
let mut ready = false;
while start.elapsed() < timeout {
let logs = Command::new("docker")
.args(["logs", "--since", "10s", CONTAINER_NAME])
.output()
.ok()?;
let stderr = String::from_utf8_lossy(&logs.stderr);
if stderr.contains("migrations: complete") || stderr.contains("service ready") {
ready = true;
break;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
if !ready {
eprintln!("skip: Aerospike container restart timeout");
return None;
}
Some(AerospikeConfig {
seed_nodes: vec![format!("127.0.0.1:{HOST_PORT}")],
namespace: "test".to_string(),
set_name: "oxcache_test".to_string(),
default_ttl: 0,
ip_map: None,
})
}
async fn make_backend() -> Option<AerospikeBackend> {
let config = shared_config().await?.clone();
let mut last_err = None;
for _ in 0..5 {
match AerospikeBackend::new(config.clone()).await {
Ok(backend) => return Some(backend),
Err(e) => {
last_err = Some(e);
tokio::time::sleep(Duration::from_secs(2)).await;
}
}
}
eprintln!("skip: Aerospike connect failed: {}", last_err.unwrap());
None
}
#[tokio::test]
async fn test_aerospike_backend_kind() {
let Some(backend) = make_backend().await else { return };
assert_eq!(backend.backend_kind(), BackendKind::Aerospike);
}
#[tokio::test]
async fn test_aerospike_set_get_delete() {
let Some(backend) = make_backend().await else { return };
backend
.set(Arc::from("as:key1"), Arc::new(b"value1".to_vec()), None)
.await
.expect("set failed");
let val = backend.get("as:key1").await.unwrap();
assert_eq!(val, Some(b"value1".to_vec()));
assert!(backend.exists("as:key1").await.unwrap());
backend.delete("as:key1").await.expect("delete failed");
let val = backend.get("as:key1").await.unwrap();
assert_eq!(val, None);
assert!(!backend.exists("as:key1").await.unwrap());
}
#[tokio::test]
async fn test_aerospike_set_with_ttl() {
let Some(backend) = make_backend().await else { return };
backend
.set(
Arc::from("as:ttl_key"),
Arc::new(b"ttl_value".to_vec()),
Some(Duration::from_secs(120)),
)
.await
.expect("set with TTL failed");
let val = backend.get("as:ttl_key").await.unwrap();
assert_eq!(val, Some(b"ttl_value".to_vec()));
let ttl = backend.ttl("as:ttl_key").await.unwrap();
assert!(ttl.is_some());
let ttl = ttl.unwrap();
assert!(ttl > Duration::from_secs(100));
}
#[tokio::test]
async fn test_aerospike_expire() {
let Some(backend) = make_backend().await else { return };
backend
.set(Arc::from("as:exp_key"), Arc::new(b"exp_value".to_vec()), None)
.await
.unwrap();
let result = backend.expire("as:exp_key", Duration::from_secs(60)).await.unwrap();
assert!(result);
let ttl = backend.ttl("as:exp_key").await.unwrap();
assert!(ttl.is_some());
let result = backend.expire("as:nonexistent", Duration::from_secs(60)).await.unwrap();
assert!(!result);
}
#[tokio::test]
async fn test_aerospike_set_many_delete_many() {
let Some(backend) = make_backend().await else { return };
let items = vec![
(Arc::from("as:batch1"), Arc::new(b"b1".to_vec()), None),
(Arc::from("as:batch2"), Arc::new(b"b2".to_vec()), None),
(Arc::from("as:batch3"), Arc::new(b"b3".to_vec()), None),
];
backend.set_many(&items).await.expect("set_many failed");
assert!(backend.get("as:batch1").await.unwrap().is_some());
assert!(backend.get("as:batch2").await.unwrap().is_some());
assert!(backend.get("as:batch3").await.unwrap().is_some());
let keys = vec![
"as:batch1".to_string(),
"as:batch2".to_string(),
"as:batch3".to_string(),
];
backend.delete_many(&keys).await.expect("delete_many failed");
assert!(backend.get("as:batch1").await.unwrap().is_none());
assert!(backend.get("as:batch2").await.unwrap().is_none());
assert!(backend.get("as:batch3").await.unwrap().is_none());
}
#[tokio::test]
async fn test_aerospike_health_check_and_stats() {
let Some(backend) = make_backend().await else { return };
backend.health_check().await.expect("health_check failed");
let stats = backend.stats().await.unwrap();
assert_eq!(stats.get("backend_kind").unwrap(), "aerospike");
assert_eq!(stats.get("namespace").unwrap(), "test");
assert_eq!(stats.get("connected").unwrap(), "true");
assert!(backend.len().await.is_err());
assert!(backend.capacity().await.is_err());
assert!(backend.keys("*").await.is_err());
assert!(backend.clear().await.is_err());
backend.shutdown().await;
}
#[tokio::test]
async fn test_aerospike_chain_cache_basic() {
use oxcache::backend::MokaMemoryBackend;
use oxcache::cache::chain::{ChainCacheBuilder, ChainLink};
let Some(aerospike) = make_backend().await else { return };
let moka = MokaMemoryBackend::new();
let chain = ChainCacheBuilder::default()
.link(ChainLink::new(moka, 100, false, "moka"))
.link(ChainLink::new(aerospike, 30, true, "aerospike"))
.build();
chain
.set("chain:as_key1", b"chain_value".to_vec(), None)
.await
.expect("chain set failed");
let val = chain.get("chain:as_key1").await.unwrap();
assert_eq!(val, Some(b"chain_value".to_vec()));
chain.health_check().await.expect("chain health_check failed");
}