use std::fs;
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::Duration;
use camel_cli::commands::test::document::parse_test_document;
use camel_cli::commands::test::run_tests;
use camel_cli::commands::test::runner::run_test_doc;
mod common;
use common::{drain_to_buffer, send_term, spawn_camel_run, wait_exit_bounded, wait_for_marker};
fn temp_dir(tag: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"camel-test-intercepts-{tag}-{}",
std::process::id()
));
fs::create_dir_all(&dir).expect("create temp dir"); dir
}
#[tokio::test(flavor = "multi_thread")]
async fn skip_to_unregistered_component_passes() {
let dir = temp_dir("skip-unregistered");
let yaml = r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "kafka:orders"
inputs:
- to: "direct:start"
body: "x"
intercepts:
kafka:orders:
skipTo: mock:orders
expects:
mock:orders:
count: 1
bodies: ["x"]
"#;
let doc = parse_test_document(yaml).expect("document should parse"); let (result, _mock) = run_test_doc(&doc, &dir).await;
assert!(
result.doc_error.is_none(),
"doc_error: {:?}",
result.doc_error
);
assert_eq!(result.endpoint_results.len(), 1);
for er in &result.endpoint_results {
assert!(
er.outcome.is_ok(),
"endpoint {} failed: {:?}",
er.endpoint,
er.outcome
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn intercept_target_and_expects_meet_on_endpoint() {
let dir = temp_dir("naming-bridge");
let yaml = r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "kafka:orders"
inputs:
- to: "direct:start"
body: "x"
intercepts:
kafka:orders:
skipTo: mock:orders
expects:
mock:orders:
count: 1
"#;
let doc = parse_test_document(yaml).expect("document should parse"); let (result, mock) = run_test_doc(&doc, &dir).await;
assert!(
result.doc_error.is_none(),
"doc_error: {:?}",
result.doc_error
);
assert_eq!(result.endpoint_results.len(), 1);
for er in &result.endpoint_results {
assert!(
er.outcome.is_ok(),
"endpoint {} failed: {:?}",
er.endpoint,
er.outcome
);
}
let inner = mock
.get_endpoint("orders")
.expect("mock:orders endpoint exists"); assert_eq!(inner.received_count().await, 1);
}
#[tokio::test(flavor = "multi_thread")]
async fn divert_copies_while_real_seda_receives() {
let dir = temp_dir("divert-seda");
let yaml = r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "seda:audit"
- to: "mock:sink"
- id: r2
from: "seda:audit"
steps:
- to: "mock:drained"
inputs:
- to: "direct:start"
body: "x"
intercepts:
seda:audit:
divertCopyTo: mock:audit
expects:
mock:audit:
count: 1
mock:drained:
count: 1
mock:sink:
count: 1
"#;
let doc = parse_test_document(yaml).expect("document should parse"); let (result, _mock) = run_test_doc(&doc, &dir).await;
assert!(
result.doc_error.is_none(),
"doc_error: {:?}",
result.doc_error
);
assert_eq!(result.endpoint_results.len(), 3);
for er in &result.endpoint_results {
assert!(
er.outcome.is_ok(),
"endpoint {} failed: {:?}",
er.endpoint,
er.outcome
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn divert_unregistered_fails_route_load() {
let dir = temp_dir("divert-unregistered");
let yaml = r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "kafka:orders"
inputs:
- to: "direct:start"
body: "x"
intercepts:
kafka:orders:
divertCopyTo: mock:orders
expects:
mock:orders:
count: 1
"#;
let doc = parse_test_document(yaml).expect("document should parse"); let (result, _mock) = run_test_doc(&doc, &dir).await;
let err = result.doc_error.expect("doc_error must be Some"); assert!(
err.contains("kafka"),
"doc_error must name the unresolvable component, got: {err}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn source_query_params_are_significant() {
let dir = temp_dir("query-significant");
let yaml = r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "kafka:orders?x=1"
inputs:
- to: "direct:start"
body: "x"
intercepts:
kafka:orders:
skipTo: mock:orders
expects:
mock:orders:
count: 1
"#;
let doc = parse_test_document(yaml).expect("document should parse"); let (result, _mock) = run_test_doc(&doc, &dir).await;
let err = result.doc_error.expect("doc_error must be Some"); assert!(
err.contains("kafka"),
"doc_error must name the unresolvable component, got: {err}"
);
}
#[test]
fn camel_run_ignores_intercepts_block() {
let dir = tempfile::tempdir().expect("tempdir");
std::fs::write(
dir.path().join("Camel.toml"),
r#"[default]
routes = ["config/*.yaml"]
log_level = "INFO"
"#,
)
.expect("write Camel.toml"); let config_dir = dir.path().join("config");
std::fs::create_dir_all(&config_dir).expect("create config/"); std::fs::write(
config_dir.join("routes.yaml"),
r#"routes:
- id: "demo"
from: "direct:start"
steps:
- to: "log:done"
"#,
)
.expect("write config/routes.yaml"); std::fs::write(
config_dir.join("orders.test.yaml"),
r#"routeFiles: [routes.yaml]
inputs: []
intercepts:
mock:result:
skipTo: mock:intercepted
expects:
mock:result:
count: 1
"#,
)
.expect("write config/orders.test.yaml");
let mut child = spawn_camel_run(dir.path());
let out_buf: Arc<Mutex<String>> = Arc::new(Mutex::new(String::new()));
let err_buf: Arc<Mutex<String>> = Arc::new(Mutex::new(String::new()));
let stdout = child
.stdout
.take()
.expect("child stdout was configured as piped"); let stderr = child
.stderr
.take()
.expect("child stderr was configured as piped");
let out_thread_buf = Arc::clone(&out_buf);
let err_thread_buf = Arc::clone(&err_buf);
let out_handle = thread::spawn(move || drain_to_buffer(stdout, out_thread_buf));
let err_handle = thread::spawn(move || drain_to_buffer(stderr, err_thread_buf));
let alive = wait_for_marker(
&mut child,
&[Arc::clone(&out_buf), Arc::clone(&err_buf)],
"camel-cli: running",
Duration::from_secs(30),
);
let captured_so_far = || {
format!(
"stdout:\n{}\nstderr:\n{}",
out_buf.lock().expect("stdout buffer lock poisoned"),
err_buf.lock().expect("stderr buffer lock poisoned")
)
};
assert!(
alive,
"camel run did not reach the running state within 30 s;\n{}",
captured_so_far()
);
thread::sleep(Duration::from_secs(1));
let status = child.try_wait().expect("try_wait after liveness window"); assert!(
status.is_none(),
"camel run died after startup (status: {status:?}); the run must ignore \
the colocated test document;\n{}",
captured_so_far()
);
send_term(&child);
let exited = wait_exit_bounded(&mut child, Duration::from_secs(10));
let _ = out_handle.join();
let _ = err_handle.join();
assert!(
exited,
"camel run did not exit within 10 s after SIGTERM;\n{}",
captured_so_far()
);
let stdout_text = out_buf.lock().expect("stdout buffer lock poisoned").clone(); let stderr_text = err_buf.lock().expect("stderr buffer lock poisoned").clone(); for (stream, text) in [("stdout", &stdout_text), ("stderr", &stderr_text)] {
assert!(
!text.contains("orders.test.yaml"),
"{stream} names the test document; camel run must skip *.test.yaml:\n{text}"
);
assert!(
!text.contains("intercepts"),
"{stream} names the intercepts block; camel run must not read it:\n{text}"
);
assert!(
!text.contains("start with 'mock:'") && !text.contains("start with `mock:`"),
"{stream} carries an intercept validation error; camel run must not parse \
test documents:\n{text}"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn mixed_multi_doc_run_isolated() {
let dir = temp_dir("mixed-isolated");
let a_path = dir.join("a.test.yaml");
let b_path = dir.join("b.test.yaml");
fs::write(
&a_path,
r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "kafka:orders"
inputs:
- to: "direct:start"
body: "a"
intercepts:
kafka:orders:
skipTo: mock:orders
expects:
mock:orders:
count: 1
bodies: ["a"]
"#,
)
.expect("write a.test.yaml"); fs::write(
&b_path,
r#"
routes:
- id: r1
from: "direct:start"
steps:
- to: "mock:plain"
inputs:
- to: "direct:start"
body: "b"
expects:
mock:plain:
count: 1
bodies: ["b"]
"#,
)
.expect("write b.test.yaml");
let mut out = Vec::new();
let mut err = Vec::new();
let summary = run_tests(&[a_path.clone(), b_path.clone()], &mut out, &mut err).await;
let out_str = String::from_utf8(out).expect("out utf8"); let err_str = String::from_utf8(err).expect("err utf8"); assert!(
err_str.is_empty(),
"no doc_error expected; err was: {err_str}\nout was: {out_str}"
);
assert_eq!(
summary.exit_code, 0,
"both docs must pass (exit 0); err: {err_str} out: {out_str}"
);
assert_eq!(
summary.passed, 2,
"both endpoints must pass; out: {out_str} err: {err_str}"
);
assert_eq!(
summary.failed, 0,
"no failures expected; out: {out_str} err: {err_str}"
);
assert!(
out_str.contains("a.test.yaml#orders") || out_str.contains("a.test.yaml#mock:orders"),
"out must contain PASS for a.test.yaml mock:orders; out: {out_str}"
);
assert!(
out_str.contains("b.test.yaml#plain") || out_str.contains("b.test.yaml#mock:plain"),
"out must contain PASS for b.test.yaml mock:plain; out: {out_str}"
);
assert!(
out_str.contains("2 passed, 0 failed"),
"out summary must be 2 passed, 0 failed; out: {out_str}"
);
}