use std::net::TcpListener;
use std::time::Duration;
use camel_api::CircuitBreakerConfig;
use camel_api::{Exchange, Message};
use camel_builder::{RouteBuilder, StepAccumulator};
use camel_component_api::{NoOpComponentContext, RuntimeObservability};
use camel_component_direct::DirectComponent;
use camel_config::config::CamelConfig;
use camel_core::CamelContext;
use std::sync::Arc;
use tower::ServiceExt;
fn test_rt() -> Arc<dyn RuntimeObservability> {
Arc::new(NoOpComponentContext)
}
fn prealloc_port() -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind 127.0.0.1:0");
listener
.local_addr()
.expect("read pre-allocated local addr")
.port()
}
async fn context_from_toml(toml: &str) -> CamelContext {
let config: CamelConfig = toml::from_str(toml).expect("test TOML parses into CamelConfig");
let mut ctx = CamelConfig::configure_context(&config)
.await
.expect("configure_context succeeds");
ctx.register_component(DirectComponent::new());
ctx
}
async fn route_started(ctx: &CamelContext, route_id: &str) -> bool {
matches!(
ctx.runtime_route_status(route_id).await,
Ok(Some(status)) if status == "Started"
)
}
async fn wait_for_started(ctx: &CamelContext, route_id: &str) {
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
while !route_started(ctx, route_id).await {
assert!(
tokio::time::Instant::now() < deadline,
"route {route_id} did not reach Started within 5s"
);
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
async fn send_one(ctx: &CamelContext) -> Result<Exchange, camel_api::CamelError> {
let producer = {
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("direct:cb-entry", ctx)
.expect("failed to create direct endpoint");
endpoint
.create_producer(test_rt(), &producer_ctx)
.expect("failed to create direct producer")
};
producer
.oneshot(Exchange::new_in_out(Message::new("cb-probe")))
.await
}
async fn poll_metrics_body(port: u16) -> String {
let url = format!("http://127.0.0.1:{port}/metrics");
let client = reqwest::Client::new();
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
loop {
if let Ok(resp) = client.get(&url).send().await
&& resp.status().is_success()
&& let Ok(body) = resp.text().await
&& body.contains("camel_exchanges_total")
&& body.contains("camel_errors_total")
{
return body;
}
assert!(
tokio::time::Instant::now() < deadline,
"prometheus /metrics never exposed exchange+error families at {url}"
);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
async fn poll_rejection_body(port: u16, route: &str) -> String {
let url = format!("http://127.0.0.1:{port}/metrics");
let client = reqwest::Client::new();
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
loop {
if let Ok(resp) = client.get(&url).send().await
&& resp.status().is_success()
&& let Ok(body) = resp.text().await
&& rejection_value(&body, route).is_some_and(|v| v > 0)
{
return body;
}
assert!(
tokio::time::Instant::now() < deadline,
"prometheus /metrics never exposed a positive rejection counter for {route} at {url}"
);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
fn errors_series(body: &str) -> Vec<String> {
let mut lines: Vec<String> = body
.lines()
.filter(|l| l.starts_with("camel_errors_total"))
.map(str::to_string)
.collect();
lines.sort();
lines
}
fn rejection_value(body: &str, route: &str) -> Option<u64> {
let needle = format!("camel_circuit_breaker_rejections_total{{route=\"{route}\"}} ");
body.lines().find_map(|l| {
l.strip_prefix(&needle)
.and_then(|v| v.trim().parse::<u64>().ok())
})
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn open_breaker_end_to_end() {
let port = prealloc_port();
let toml_cfg = format!(
r#"[observability.prometheus]
enabled = true
host = "127.0.0.1"
port = {port}
"#
);
let ctx = context_from_toml(&toml_cfg).await;
let route = RouteBuilder::from("direct:cb-entry")
.route_id("cb-entry")
.circuit_breaker(
CircuitBreakerConfig::new()
.failure_threshold(1)
.open_duration(Duration::from_secs(60)),
)
.to("direct:missing?failIfNoConsumers=false")
.build()
.expect("breaker route builds");
let mut ctx = ctx;
ctx.add_route_definition(route)
.await
.expect("breaker route registers");
ctx.start().await.expect("context starts");
wait_for_started(&ctx, "cb-entry").await;
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
if let Ok(Err(_)) = tokio::time::timeout(Duration::from_secs(5), send_one(&ctx)).await {
break; }
assert!(
tokio::time::Instant::now() < deadline,
"breaker route never reported the downstream failure"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
let baseline = poll_metrics_body(port).await;
let errors_before = errors_series(&baseline);
for _ in 0..3 {
let _ = tokio::time::timeout(Duration::from_millis(50), send_one(&ctx)).await;
}
let after = poll_rejection_body(port, "cb-entry").await;
let rejections = rejection_value(&after, "cb-entry")
.expect("camel_circuit_breaker_rejections_total{route=\"cb-entry\"} present");
assert!(
rejections > 0,
"open-breaker sends must increment the rejection counter:\n{after}"
);
let errors_after = errors_series(&after);
assert_eq!(
errors_after, errors_before,
"camel_errors_total must not grow during open-phase sends\nbefore:\n{errors_before:?}\nafter:\n{errors_after:?}"
);
ctx.stop().await.expect("context stops");
}