use std::path::Path;
use super::runner::EndpointResult;
const FAKE_SCHEME: &str = "fake";
#[cfg(feature = "integration-http")]
const BOOT_SCHEMES: [&str; 3] = [FAKE_SCHEME, "direct", "http"];
#[cfg(not(feature = "integration-http"))]
const PROVIDED_ADAPTERS: &str = "only the `fake:` in-memory adapter";
#[cfg(feature = "integration-http")]
const PROVIDED_ADAPTERS: &str = "the `fake:` in-memory adapter and the `http:` wire partner";
pub(super) struct ScenarioDocResult {
pub action_results: Vec<EndpointResult>,
pub doc_error: Option<String>,
pub apparatus: bool,
}
pub(super) fn is_apparatus(failure: &camel_integration_test::ScenarioFailure) -> bool {
use camel_integration_test::ScenarioFailure as F;
matches!(
failure,
F::ActionTransport { .. } | F::PartnerStartup { .. } | F::ShutdownFailure { .. }
)
}
fn action_label(index: usize, action: &camel_integration_test::ScenarioAction) -> String {
use camel_integration_test::ScenarioAction as A;
let kind = match action {
A::Send { .. } => "send",
A::Receive { .. } => "receive",
A::Sleep { .. } => "sleep",
A::Validate { .. } => "validate",
_ => "action",
};
format!("scenario[{index}] {kind}")
}
fn outcome_rows(
doc: &camel_integration_test::ScenarioDocument,
outcome: &camel_integration_test::DocumentOutcome,
) -> (Vec<EndpointResult>, bool) {
let mut apparatus = false;
let action_results = outcome
.per_action
.iter()
.enumerate()
.map(|(index, result)| {
let label = action_label(index, &doc.scenario[index]);
match result {
Ok(_) => EndpointResult {
endpoint: label,
outcome: Ok(()),
},
Err(failure) => {
apparatus = apparatus || is_apparatus(failure);
EndpointResult {
endpoint: label,
outcome: Err(failure.to_string()),
}
}
}
})
.collect();
(action_results, apparatus)
}
fn scheme_of(endpoint: &str) -> &str {
endpoint.split(':').next().unwrap_or(endpoint)
}
fn wire_endpoint_refs(
doc: &camel_integration_test::ScenarioDocument,
) -> Vec<&camel_integration_test::EndpointRef> {
use camel_integration_test::ScenarioAction as A;
let mut seen = std::collections::BTreeSet::new();
let mut out = Vec::new();
for action in &doc.scenario {
let reference = match action {
A::Send { to, .. } => to,
A::Receive { from, .. } => from,
_ => continue,
};
if seen.insert(reference.endpoint.clone()) {
out.push(reference);
}
}
out
}
#[cfg(feature = "integration-http")]
fn is_harness_http(reference: &camel_integration_test::EndpointRef) -> bool {
use camel_integration_test::Provisioning;
reference.provisioning == Some(Provisioning::Harness)
&& scheme_of(&reference.endpoint) == "http"
}
#[cfg(feature = "integration-http")]
fn partner_scripts_for(
doc: &camel_integration_test::ScenarioDocument,
endpoint: &str,
) -> Option<Vec<camel_integration_test::ScriptedResponse>> {
camel_integration_test::partner_scripts_for(doc, endpoint)
}
#[cfg(feature = "integration-http")]
async fn bind_partners(
doc: &camel_integration_test::ScenarioDocument,
wired: &[&camel_integration_test::EndpointRef],
) -> Result<
(
std::collections::BTreeMap<String, Box<dyn camel_integration_test::PartnerAdapter>>,
std::collections::BTreeMap<String, String>,
),
String,
> {
use camel_integration_test::HttpPartner;
let mut adapters: std::collections::BTreeMap<
String,
Box<dyn camel_integration_test::PartnerAdapter>,
> = std::collections::BTreeMap::new();
let mut harness_provisioned: std::collections::BTreeMap<String, String> =
std::collections::BTreeMap::new();
for reference in wired {
let endpoint = &reference.endpoint;
if !is_harness_http(reference) {
continue;
}
let partner = match partner_scripts_for(doc, endpoint) {
Some(scripts) => HttpPartner::start(scripts).await,
None => HttpPartner::start_permissive(200).await,
};
let partner = match partner {
Ok(partner) => partner,
Err(e) => {
return Err(format!("partner-bind-failure: endpoint {endpoint}: {e}"));
}
};
if let Some(bind_var) = &reference.bind_var {
harness_provisioned
.insert(bind_var.clone(), format!("http://{}", partner.bound_addr()));
}
adapters.insert(endpoint.clone(), Box::new(partner));
}
Ok((adapters, harness_provisioned))
}
pub(super) async fn run_scenario_doc(
doc: &camel_integration_test::ScenarioDocument,
#[cfg_attr(not(feature = "integration-http"), allow(unused_variables))] root: &Path,
) -> ScenarioDocResult {
let wired = wire_endpoint_refs(doc);
#[cfg(feature = "integration-http")]
if wired.iter().any(|r| scheme_of(&r.endpoint) != FAKE_SCHEME)
&& wired
.iter()
.all(|r| BOOT_SCHEMES.contains(&scheme_of(&r.endpoint)))
{
return run_scenario_full_boot(doc, root).await;
}
run_scenario_fake_smoke(doc, &wired).await
}
async fn run_scenario_fake_smoke(
doc: &camel_integration_test::ScenarioDocument,
wired: &[&camel_integration_test::EndpointRef],
) -> ScenarioDocResult {
use camel_integration_test::{FakeAdapter, PartnerRouter, ScenarioVars};
let mut adapters: std::collections::BTreeMap<
String,
Box<dyn camel_integration_test::PartnerAdapter>,
> = std::collections::BTreeMap::new();
let mut missing: Vec<(&str, &str)> = Vec::new();
for reference in wired {
let endpoint = reference.endpoint.as_str();
if scheme_of(endpoint) == FAKE_SCHEME {
adapters.insert(
endpoint.to_string(),
Box::new(FakeAdapter::scripted(Vec::new())),
);
} else {
missing.push((endpoint, scheme_of(endpoint)));
}
}
if let Some((endpoint, scheme)) = missing.first() {
return ScenarioDocResult {
action_results: Vec::new(),
doc_error: Some(format!(
"infra-unavailable: endpoint {endpoint} needs the {scheme} partner adapter; \
this build provides {PROVIDED_ADAPTERS}"
)),
apparatus: true,
};
}
debug_assert!(
wired
.iter()
.all(|reference| adapters.contains_key(&reference.endpoint))
);
let router = PartnerRouter::new(adapters);
let mut vars = ScenarioVars::new();
let outcome = camel_integration_test::run_scenario_document(doc, &router, &mut vars).await;
let (action_results, apparatus) = outcome_rows(doc, &outcome);
ScenarioDocResult {
action_results,
doc_error: None,
apparatus,
}
}
#[cfg(feature = "integration-http")]
async fn http_secret_query_keys(
ctx: &std::sync::Arc<tokio::sync::Mutex<camel_core::CamelContext>>,
) -> Vec<String> {
ctx.lock()
.await
.component_metadata("http")
.map(|metadata| {
metadata
.uri_options
.iter()
.filter(|option| option.secret)
.map(|option| option.name.clone())
.collect()
})
.unwrap_or_default()
}
#[cfg(feature = "integration-http")]
async fn run_scenario_full_boot(
doc: &camel_integration_test::ScenarioDocument,
root: &Path,
) -> ScenarioDocResult {
use camel_integration_test::{
DirectStimulus, EndpointRef, FakeAdapter, LayeredEnv, PartnerRouter, ScenarioFailure,
ScenarioVars, ambient_std, boot_scenario, run_scenario_document,
};
use std::sync::Arc;
use tokio::sync::Mutex;
let wired = wire_endpoint_refs(doc);
if let Some(partners) = &doc.partners {
let harness_http: std::collections::BTreeSet<&str> = wired
.iter()
.copied()
.filter(|r| is_harness_http(r))
.map(|r| r.endpoint.as_str())
.collect();
if let Some(key) = partners
.keys()
.find(|key| !harness_http.contains(key.as_str()))
{
return ScenarioDocResult {
action_results: Vec::new(),
doc_error: Some(format!(
"doc-validation: partners[{key}]: no wired harness `http` endpoint \
reference declares this key"
)),
apparatus: true,
};
}
}
let (mut adapters, harness_provisioned) = match bind_partners(doc, &wired).await {
Ok((adapters, harness_provisioned)) => (adapters, harness_provisioned),
Err(doc_error) => {
return ScenarioDocResult {
action_results: Vec::new(),
doc_error: Some(doc_error),
apparatus: true,
};
}
};
let env = LayeredEnv::new(
doc.env.clone().unwrap_or_default(),
harness_provisioned,
doc.env_passthrough.clone().unwrap_or_default(),
ambient_std(),
);
let run = match boot_scenario(doc, root, &env).await {
Ok(run) => run,
Err(e) => {
let class = matches!(e, camel_api::CamelError::AuthProviderUnavailable(_))
.then_some("infra-unavailable")
.unwrap_or("full-boot-failure");
return ScenarioDocResult {
action_results: Vec::new(),
doc_error: Some(format!("{class}: scenario boot failed: {e}")),
apparatus: true,
};
}
};
let ctx = Arc::new(Mutex::new(run.ctx));
for reference in &wired {
match scheme_of(&reference.endpoint) {
"direct" => {
adapters.insert(
reference.endpoint.clone(),
Box::new(DirectStimulus::new(Arc::clone(&ctx))),
);
}
FAKE_SCHEME => {
adapters.insert(
reference.endpoint.clone(),
Box::new(FakeAdapter::scripted(Vec::new())),
);
}
_ => {}
}
}
let router = PartnerRouter::new(adapters);
router.set_secret_query_keys(http_secret_query_keys(&ctx).await);
let mut vars = ScenarioVars::new();
let wired_refs: Vec<EndpointRef> = wired.iter().map(|r| (*r).clone()).collect();
camel_integration_test::runner::fill_bind_vars(&wired_refs, &router, &mut vars);
let mut outcome = run_scenario_document(doc, &router, &mut vars).await;
let (mut action_results, mut apparatus) = outcome_rows(doc, &outcome);
{
let mut guard = ctx.lock().await;
if let Err(e) = run.boot.shutdown(&mut guard).await {
tracing::error!("scenario boot shutdown failed: {e}");
outcome.final_failure = Some(ScenarioFailure::ShutdownFailure {
message: e.to_string(),
});
}
}
if let Some(final_failure) = &outcome.final_failure {
apparatus = true;
action_results.push(EndpointResult {
endpoint: "shutdown".to_string(),
outcome: Err(final_failure.to_string()),
});
}
ScenarioDocResult {
action_results,
doc_error: None,
apparatus,
}
}
#[cfg(all(test, feature = "integration-http"))]
#[path = "scenario_tests.rs"]
mod tests;