#![allow(dead_code)]
use std::collections::BTreeMap;
use std::sync::Arc;
use camel_integration_test::runner::fill_bind_vars;
use camel_integration_test::{
DirectStimulus, DocumentOutcome, EndpointRef, LayeredEnv, PartnerAdapter, PartnerRouter,
ScenarioAction, ScenarioTarget, ScenarioVars, ambient_std, boot_scenario,
ensure_capture_subscriber, parse_scenario_document, run_scenario_document,
};
pub static RUN_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
pub fn lock_run() -> std::sync::MutexGuard<'static, ()> {
RUN_LOCK
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub async fn run_logs_document(doc_yaml: &str, fixture: &str) -> DocumentOutcome {
ensure_capture_subscriber();
let dir = tempfile::tempdir().expect("temp dir");
let root = dir.path();
std::fs::write(root.join("Camel.toml"), "log_level = \"info\"\n").expect("write Camel.toml");
let fixture_path = format!(
"{}/tests/fixtures/logs/{fixture}",
env!("CARGO_MANIFEST_DIR")
);
std::fs::copy(&fixture_path, root.join("routes.yaml")).expect("copy route fixture");
let path = root.join("case.test.yaml");
std::fs::write(&path, doc_yaml).expect("write case file");
let doc = parse_scenario_document(&path).expect("document must load");
let env = LayeredEnv::new(
doc.env.clone().unwrap_or_default(),
BTreeMap::new(),
doc.env_passthrough.clone().unwrap_or_default(),
ambient_std(),
);
let run = boot_scenario(&doc, root, &env)
.await
.expect("the full boot must succeed");
let ctx = Arc::new(tokio::sync::Mutex::new(run.ctx));
let mut adapters: BTreeMap<String, Box<dyn PartnerAdapter>> = BTreeMap::new();
adapters.insert(
"direct:start".to_string(),
Box::new(DirectStimulus::new(Arc::clone(&ctx))),
);
let router = PartnerRouter::new(adapters);
let wired = wired_refs(&doc);
let mut vars = ScenarioVars::new();
fill_bind_vars(&wired, &router, &mut vars);
let outcome = run_scenario_document(&doc, &router, &mut vars, None).await;
let mut guard = ctx.lock().await;
run.boot
.shutdown(&mut guard)
.await
.expect("clean shutdown must complete");
outcome
}
pub fn wired_refs(doc: &camel_integration_test::ScenarioDocument) -> Vec<EndpointRef> {
doc.scenario
.iter()
.filter_map(|action| match action {
ScenarioAction::Send { to, .. } => Some(to.clone()),
ScenarioAction::Receive { from, .. } => Some(from.clone()),
ScenarioAction::Validate {
target: ScenarioTarget::LastReceived(endpoint),
..
} => Some(endpoint.clone()),
_ => None,
})
.collect()
}
#[cfg(feature = "http")]
pub async fn bind_doc_partners(
doc: &camel_integration_test::ScenarioDocument,
) -> (
PartnerRouter,
BTreeMap<String, camel_integration_test::HttpRecorder>,
BTreeMap<String, String>,
) {
use camel_integration_test::{HttpPartner, Provisioning, partner_scripts_for};
let mut adapters: BTreeMap<String, Box<dyn PartnerAdapter>> = BTreeMap::new();
let mut recorders: BTreeMap<String, camel_integration_test::HttpRecorder> = BTreeMap::new();
let mut authorities: BTreeMap<String, String> = BTreeMap::new();
for reference in wired_refs(doc) {
if reference.provisioning != Some(Provisioning::Harness)
|| !reference.endpoint.starts_with("http://")
|| adapters.contains_key(&reference.endpoint)
{
continue;
}
let partner = match partner_scripts_for(doc, &reference.endpoint) {
Some(scripts) => HttpPartner::start(scripts).await,
None => HttpPartner::start_permissive(200).await,
}
.expect("partner must bind 127.0.0.1:0");
authorities.insert(reference.endpoint.clone(), partner.bound_addr().to_string());
recorders.insert(reference.endpoint.clone(), partner.recorder());
adapters.insert(reference.endpoint.clone(), Box::new(partner));
}
(PartnerRouter::new(adapters), recorders, authorities)
}
#[cfg(feature = "http")]
pub async fn run_doc(
yaml: &str,
) -> (
DocumentOutcome,
BTreeMap<String, camel_integration_test::HttpRecorder>,
) {
let (outcome, recorders, _authorities) = run_doc_with_authorities(yaml).await;
(outcome, recorders)
}
#[cfg(feature = "http")]
pub async fn run_doc_with_authorities(
yaml: &str,
) -> (
DocumentOutcome,
BTreeMap<String, camel_integration_test::HttpRecorder>,
BTreeMap<String, String>,
) {
let dir = tempfile::tempdir().expect("temp dir");
let path = dir.path().join("case.test.yaml");
std::fs::write(&path, yaml).expect("write case file");
let doc = parse_scenario_document(&path).expect("document must load");
let (router, recorders, authorities) = bind_doc_partners(&doc).await;
let wired = wired_refs(&doc);
let mut vars = ScenarioVars::new();
fill_bind_vars(&wired, &router, &mut vars);
let outcome = run_scenario_document(&doc, &router, &mut vars, None).await;
(outcome, recorders, authorities)
}