use std::fs::{File, OpenOptions};
use std::io::Write;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use camel_api::{BoxProcessor, BoxProcessorExt, Exchange, OpaqueProcessor};
use camel_core::{BuilderStep, RouteDefinition};
use percent_encoding::{AsciiSet, CONTROLS, utf8_percent_encode};
const BENCH_URI_SAFE: &AsciiSet = &CONTROLS.add(b' ');
const BENCH_START: &str = "BenchStart";
const BENCH_MODE_ROUTE: &str = "route";
fn is_route_mode(mode: &str) -> bool {
mode.trim().eq_ignore_ascii_case(BENCH_MODE_ROUTE)
}
pub fn maybe_instrument_routes(defs: Vec<RouteDefinition>) -> Vec<RouteDefinition> {
let Ok(path) = std::env::var("BENCH_LATENCY_FILE") else {
return defs;
};
let file = match OpenOptions::new().create(true).append(true).open(&path) {
Ok(f) => f,
Err(e) => {
tracing::error!("bench_instrument: cannot open BENCH_LATENCY_FILE '{path}': {e}");
return defs;
}
};
let shared_file = Arc::new(Mutex::new(file));
let mode_raw = std::env::var("BENCH_LATENCY_MODE").ok();
let route_mode = mode_raw.as_deref().map(is_route_mode).unwrap_or(false);
if let Some(mode) = &mode_raw
&& !mode.is_empty()
&& !is_route_mode(mode)
{
tracing::warn!(
"bench_instrument: unrecognized BENCH_LATENCY_MODE '{mode}' \
(expected 'route'); falling back to pair mode"
);
}
if route_mode {
tracing::info!(
"bench_instrument: route-bracket mode, timer-sourced routes only (file={path})"
);
route_bracket_defs(defs, &shared_file)
} else {
tracing::info!("bench_instrument: wrapping top-level To steps (file={path})");
defs.into_iter()
.map(|def| {
let sf = Arc::clone(&shared_file);
let route_id = def.route_id().to_string();
def.map_steps(|steps| inject_timing(steps, route_id, sf))
})
.collect()
}
}
fn route_bracket_defs(defs: Vec<RouteDefinition>, file: &Arc<Mutex<File>>) -> Vec<RouteDefinition> {
defs.into_iter()
.map(|def| {
if !def.from_uri().starts_with("timer:") {
return def;
}
let sf = Arc::clone(file);
let route_id = def.route_id().to_string();
let from_uri = def.from_uri().to_string();
def.map_steps(|steps| inject_route_bracket(steps, route_id, from_uri, sf))
})
.collect()
}
fn inject_timing(
steps: Vec<BuilderStep>,
route_id: String,
file: Arc<Mutex<File>>,
) -> Vec<BuilderStep> {
let mut result = Vec::with_capacity(steps.len() * 3);
for step in steps {
if let BuilderStep::To(uri) = step {
let counter = Arc::new(AtomicU64::new(0)); result.push(make_start_processor());
result.push(BuilderStep::To(uri.clone()));
result.push(make_end_processor(
route_id.clone(),
uri,
counter,
Arc::clone(&file),
));
} else {
result.push(step);
}
}
result
}
fn inject_route_bracket(
steps: Vec<BuilderStep>,
route_id: String,
from_uri: String,
file: Arc<Mutex<File>>,
) -> Vec<BuilderStep> {
let start_slot = Arc::new(Mutex::new(None::<Instant>));
let mut result = Vec::with_capacity(steps.len() + 2);
result.push(make_route_start_processor(Arc::clone(&start_slot)));
result.extend(steps);
result.push(make_route_end_processor(
route_id, from_uri, start_slot, file,
));
result
}
fn make_route_start_processor(start_slot: Arc<Mutex<Option<Instant>>>) -> BuilderStep {
BuilderStep::Processor(OpaqueProcessor(BoxProcessor::from_fn(
move |exchange: Exchange| {
let start_slot = Arc::clone(&start_slot);
Box::pin(async move {
if let Ok(mut slot) = start_slot.lock() {
*slot = Some(Instant::now());
}
Ok(exchange)
})
},
)))
}
fn make_route_end_processor(
route_id: String,
from_uri: String,
start_slot: Arc<Mutex<Option<Instant>>>,
file: Arc<Mutex<File>>,
) -> BuilderStep {
let counter = Arc::new(AtomicU64::new(0)); BuilderStep::Processor(OpaqueProcessor(BoxProcessor::from_fn(
move |exchange: Exchange| {
let counter = Arc::clone(&counter);
let file = Arc::clone(&file);
let route_id = route_id.clone();
let from_uri = from_uri.clone();
let start_slot = Arc::clone(&start_slot);
Box::pin(async move {
let id = counter.fetch_add(1, Ordering::Relaxed) + 1;
let duration_ns = start_slot
.lock()
.ok()
.and_then(|slot| *slot)
.map(|t| t.elapsed().as_nanos() as u64)
.unwrap_or(0);
let line = format_bench_line(&route_id, &from_uri, id, duration_ns);
if let Ok(mut f) = file.lock() {
let _ = f.write_all(line.as_bytes());
}
Ok(exchange)
})
},
)))
}
fn make_start_processor() -> BuilderStep {
BuilderStep::Processor(OpaqueProcessor(BoxProcessor::from_fn(
|mut exchange: Exchange| {
Box::pin(async move {
exchange.set_extension(BENCH_START, Arc::new(Instant::now()));
Ok(exchange)
})
},
)))
}
fn make_end_processor(
route_id: String,
uri: String,
counter: Arc<AtomicU64>,
file: Arc<Mutex<File>>,
) -> BuilderStep {
BuilderStep::Processor(OpaqueProcessor(BoxProcessor::from_fn(
move |exchange: Exchange| {
let counter = Arc::clone(&counter);
let file = Arc::clone(&file);
let route_id = route_id.clone();
let uri = uri.clone();
Box::pin(async move {
let id = counter.fetch_add(1, Ordering::Relaxed) + 1;
let duration_ns = exchange
.get_extension::<Instant>(BENCH_START)
.map(|t| t.elapsed().as_nanos() as u64)
.unwrap_or(0);
let line = format_bench_line(&route_id, &uri, id, duration_ns);
if let Ok(mut f) = file.lock() {
let _ = f.write_all(line.as_bytes());
}
Ok(exchange)
})
},
)))
}
fn format_bench_line(route_id: &str, uri: &str, id: u64, duration_ns: u64) -> String {
let route = if route_id.trim().is_empty() {
"-"
} else {
route_id
};
let encoded = utf8_percent_encode(uri, BENCH_URI_SAFE);
format!("BENCH_LATENCY {id} {duration_ns} {route} {encoded}\n")
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs::File;
#[test]
fn inject_timing_wraps_each_top_level_to() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let file = Arc::new(Mutex::new(File::create(tmp.path()).unwrap()));
let steps = vec![
BuilderStep::To("xslt:a".into()),
BuilderStep::Stop,
BuilderStep::To("xslt:b".into()),
];
let out = inject_timing(steps, "test-route".to_string(), file);
assert_eq!(out.len(), 7);
assert!(matches!(out[0], BuilderStep::Processor(_)));
assert!(matches!(out[1], BuilderStep::To(_)));
assert!(matches!(out[2], BuilderStep::Processor(_)));
assert!(matches!(out[3], BuilderStep::Stop));
assert!(matches!(out[4], BuilderStep::Processor(_)));
assert!(matches!(out[5], BuilderStep::To(_)));
assert!(matches!(out[6], BuilderStep::Processor(_)));
}
#[test]
fn inject_timing_skips_non_to_steps() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let file = Arc::new(Mutex::new(File::create(tmp.path()).unwrap()));
let steps = vec![BuilderStep::Stop, BuilderStep::Stop];
let out = inject_timing(steps, "test-route".to_string(), file);
assert_eq!(out.len(), 2);
}
#[test]
fn maybe_instrument_routes_noop_when_env_unset() {
unsafe {
std::env::remove_var("BENCH_LATENCY_FILE");
}
let def = camel_core::RouteDefinition::new(
"direct:test".to_string(),
vec![BuilderStep::To("mock:a".into())],
);
let defs = vec![def];
let out = maybe_instrument_routes(defs);
assert_eq!(out[0].steps().len(), 1);
}
#[test]
fn format_bench_line_emits_route_id_and_percent_encoded_uri() {
assert_eq!(
format_bench_line("nacional-chain", "sql:noop?ds=cartodb", 1, 821_852_517),
"BENCH_LATENCY 1 821852517 nacional-chain sql:noop?ds=cartodb\n"
);
assert_eq!(
format_bench_line("r2", "http:host?q=a b", 3, 1000),
"BENCH_LATENCY 3 1000 r2 http:host?q=a%20b\n"
);
assert_eq!(
format_bench_line("", "direct:foo", 2, 500),
"BENCH_LATENCY 2 500 - direct:foo\n"
);
}
#[test]
fn is_route_mode_matches_only_route() {
assert!(is_route_mode("route"));
assert!(is_route_mode(" Route "));
assert!(is_route_mode("ROUTE"));
assert!(!is_route_mode(""));
assert!(!is_route_mode("to"));
assert!(!is_route_mode("pair"));
assert!(!is_route_mode("routes"));
}
#[test]
fn route_bracket_defs_brackets_only_timer_sourced_routes() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let file = Arc::new(Mutex::new(File::create(tmp.path()).unwrap()));
let timer = camel_core::RouteDefinition::new(
"timer:bench?period=10&repeatCount=10000".to_string(),
vec![BuilderStep::Stop],
);
let consumer = camel_core::RouteDefinition::new(
"direct:agg-in".to_string(),
vec![BuilderStep::To("mock:a".into())],
);
let out = route_bracket_defs(vec![timer, consumer], &file);
assert_eq!(out.len(), 2);
assert_eq!(out[0].steps().len(), 3);
assert!(matches!(
out[0].steps().first(),
Some(BuilderStep::Processor(_))
));
assert!(matches!(
out[0].steps().last(),
Some(BuilderStep::Processor(_))
));
assert_eq!(out[1].steps().len(), 1);
assert!(matches!(out[1].steps().first(), Some(BuilderStep::To(_))));
}
#[test]
fn inject_route_bracket_wraps_whole_step_list() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let file = Arc::new(Mutex::new(File::create(tmp.path()).unwrap()));
let steps = vec![BuilderStep::Stop, BuilderStep::Stop];
let out = inject_route_bracket(
steps,
"bench-route".to_string(),
"timer:bench".to_string(),
file,
);
assert_eq!(out.len(), 4);
assert!(matches!(out[0], BuilderStep::Processor(_)));
assert!(matches!(out[1], BuilderStep::Stop));
assert!(matches!(out[2], BuilderStep::Stop));
assert!(matches!(out[3], BuilderStep::Processor(_)));
}
#[tokio::test]
async fn route_bracket_emits_one_record_per_pass() {
use tower::Service as _;
use tower::ServiceExt as _;
let tmp = tempfile::NamedTempFile::new().unwrap();
let file = Arc::new(Mutex::new(File::create(tmp.path()).unwrap()));
let steps = vec![BuilderStep::Stop];
let out = inject_route_bracket(
steps,
"bench-route".to_string(),
"timer:bench?period=10&repeatCount=10000".to_string(),
file,
);
let (start, end) = match (out.first(), out.last()) {
(Some(BuilderStep::Processor(s)), Some(BuilderStep::Processor(e))) => {
(s.0.clone(), e.0.clone())
}
_ => panic!("expected processor bracket, got {out:?}"),
};
for _pass in 1..=3u64 {
let ex = Exchange::new(camel_api::Message::default());
let ex = start.clone().ready().await.unwrap().call(ex).await.unwrap();
let _ = end.clone().ready().await.unwrap().call(ex).await.unwrap();
}
let content = std::fs::read_to_string(tmp.path()).unwrap();
let lines: Vec<&str> = content.lines().collect();
assert_eq!(lines.len(), 3, "one record per pass, got: {content}");
for (i, line) in lines.iter().enumerate() {
let id = i + 1;
let rest = line
.strip_prefix(&format!("BENCH_LATENCY {id} "))
.unwrap_or_else(|| panic!("record {id} has bad id: {line}"));
let ns = rest.split_whitespace().next().unwrap();
assert!(
ns.parse::<u64>().unwrap() > 0,
"record {id} duration must be positive ns: {line}"
);
assert_eq!(
rest,
format!("{ns} bench-route timer:bench?period=10&repeatCount=10000"),
"record {id} exact format: {line}"
);
}
}
}