use std::time::Duration;
use camel_api::{Exchange, Message, Value};
use camel_test::CamelTestContext;
use tower::ServiceExt;
fn test_rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
std::sync::Arc::new(camel_component_api::NoOpComponentContext)
}
async fn send_to_direct(
h: &CamelTestContext,
endpoint_uri: &str,
exchange: Exchange,
timeout: Duration,
) {
let deadline = tokio::time::Instant::now() + timeout;
loop {
let producer = {
let ctx = h.ctx().lock().await;
let producer_ctx = ctx.producer_context();
let registry = ctx.registry();
let component = registry
.get("direct")
.expect("direct component not registered");
let endpoint = component
.create_endpoint(endpoint_uri, &*ctx)
.expect("failed to create direct endpoint");
endpoint
.create_producer(test_rt(), &producer_ctx)
.expect("failed to create direct producer")
};
match producer.oneshot(exchange.clone()).await {
Ok(_) => return,
Err(_) if tokio::time::Instant::now() < deadline => {
tokio::time::sleep(Duration::from_millis(20)).await;
}
Err(e) => panic!("failed to send exchange within {timeout:?}: {e}"),
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cache_eip_hit_and_miss_end_to_end() {
let h = CamelTestContext::builder()
.with_direct()
.with_mock()
.build()
.await;
let yaml = r#"
routes:
- id: "cache-smoke"
from: "direct:start"
steps:
- cache:
key: "${header.cacheKey}"
on_miss:
- set_body: "fresh"
- to: "mock:result"
"#;
for route in camel_dsl::parse_yaml(yaml).unwrap() {
h.add_route(route).await.unwrap();
}
h.start().await;
let mock = h
.mock()
.get_endpoint("result")
.expect("mock endpoint created during route compilation");
let mut ex1 = Exchange::new(Message::new("original"));
ex1.input.set_header("cacheKey", Value::String("k".into()));
send_to_direct(&h, "direct:start", ex1, Duration::from_secs(2)).await;
mock.await_exchanges(1, Duration::from_secs(2)).await;
{
let received = mock.get_received_exchanges().await;
let body = received[0].input.body.as_text();
assert_eq!(
body,
Some("fresh"),
"exchange 1 (miss): on_miss must set body to 'fresh'"
);
}
let mut ex2 = Exchange::new(Message::new("original"));
ex2.input.set_header("cacheKey", Value::String("k".into()));
send_to_direct(&h, "direct:start", ex2, Duration::from_secs(2)).await;
mock.await_exchanges(2, Duration::from_secs(2)).await;
{
let received = mock.get_received_exchanges().await;
assert_eq!(received.len(), 2, "exactly 2 exchanges must reach the mock");
let body1 = received[0].input.body.as_text();
let body2 = received[1].input.body.as_text();
assert_eq!(
body1,
Some("fresh"),
"exchange 1 (miss): body must be 'fresh'"
);
assert_eq!(
body2,
Some("fresh"),
"exchange 2 (hit): body must be reconstructed from cache as 'fresh'"
);
}
h.stop().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cache_invalidate_step_compiles_and_executes() {
let h = CamelTestContext::builder()
.with_direct()
.with_mock()
.build()
.await;
let yaml = r#"
routes:
- id: "cache-invalidate-smoke"
from: "direct:start"
steps:
- cache:
key: "${header.cacheKey}"
on_miss:
- set_body: "fresh"
- cache_invalidate:
key: "${header.cacheKey}"
- to: "mock:result"
"#;
for route in camel_dsl::parse_yaml(yaml).unwrap() {
h.add_route(route).await.unwrap();
}
h.start().await;
let mock = h.mock().get_endpoint("result").expect("mock endpoint");
let mut ex1 = Exchange::new(Message::new("original"));
ex1.input.set_header("cacheKey", Value::String("k".into()));
send_to_direct(&h, "direct:start", ex1, Duration::from_secs(2)).await;
mock.await_exchanges(1, Duration::from_secs(2)).await;
let mut ex2 = Exchange::new(Message::new("original"));
ex2.input.set_header("cacheKey", Value::String("k".into()));
send_to_direct(&h, "direct:start", ex2, Duration::from_secs(2)).await;
mock.await_exchanges(2, Duration::from_secs(2)).await;
let received = mock.get_received_exchanges().await;
assert_eq!(received.len(), 2);
assert_eq!(received[0].input.body.as_text(), Some("fresh"));
assert_eq!(received[1].input.body.as_text(), Some("fresh"));
h.stop().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cache_peek_stale_step_compiles_and_executes() {
let h = CamelTestContext::builder()
.with_direct()
.with_mock()
.build()
.await;
let yaml = r#"
routes:
- id: "cache-peek-stale-smoke"
from: "direct:start"
steps:
- cache:
key: "${header.cacheKey}"
ttl: "1ms"
on_miss:
- set_body: "stale_data"
- to: "mock:result"
"#;
for route in camel_dsl::parse_yaml(yaml).unwrap() {
h.add_route(route).await.unwrap();
}
h.start().await;
let mock = h.mock().get_endpoint("result").expect("mock endpoint");
let mut ex1 = Exchange::new(Message::new("original"));
ex1.input.set_header("cacheKey", Value::String("k".into()));
send_to_direct(&h, "direct:start", ex1, Duration::from_secs(2)).await;
mock.await_exchanges(1, Duration::from_secs(2)).await;
tokio::time::sleep(Duration::from_millis(50)).await;
{
let ctx = h.ctx().lock().await;
let repo = ctx
.cache_repository("memory")
.expect("memory cache registered");
let stale = repo.peek_stale("k").await.unwrap();
assert!(
stale.is_some(),
"peek_stale should return the expired-but-retained entry"
);
let entry = stale.unwrap();
assert_eq!(
entry.bytes,
b"stale_data".to_vec(),
"stale entry should contain the original data"
);
}
h.stop().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn unregistered_repository_returns_error() {
let h = CamelTestContext::builder()
.with_direct()
.with_mock()
.build()
.await;
let yaml = r#"
routes:
- id: "cache-bad-repo"
from: "direct:start"
steps:
- cache:
repository: "absent"
key: "k"
on_miss:
- set_body: "x"
- to: "mock:result"
"#;
let routes = camel_dsl::parse_yaml(yaml).unwrap();
let result = h.add_route(routes.into_iter().next().unwrap()).await;
assert!(
result.is_err(),
"route with unregistered repository 'absent' must fail to compile, got: {result:?}"
);
h.stop().await;
}