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";
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));
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 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 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"
);
}
}