#![cfg(feature = "http")]
use std::collections::BTreeMap;
use std::sync::Arc;
use camel_integration_test::{
DirectStimulus, DocumentOutcome, HttpPartner, LayeredEnv, PartnerAdapter, PartnerRouter,
ScenarioDocument, ScenarioFailure, ScenarioVerdict, ambient_std, boot_scenario,
parse_scenario_document, partner_scripts_for, run_scenario_document,
};
use tokio::sync::Mutex;
const UPSTREAM: &str = "http://127.0.0.1:0/upstream";
const TTL: &str = "100ms";
const SLEEP: &str = "400ms";
fn routes_yaml() -> String {
r#"
routes:
- id: cb-stale
from: direct:fetch
circuit_breaker:
failure_threshold: 1
open_duration_ms: 60000
fallback:
- cache_peek_stale:
repository: memory
key: tile-xyz
steps:
- cache:
repository: memory
key: tile-xyz
ttl: "TTL"
on_miss:
- set_body: "tile-stale-7f"
- to: ${env:UPSTREAM}
- id: cb-miss
from: direct:miss
circuit_breaker:
failure_threshold: 1
open_duration_ms: 60000
fallback:
- cache_peek_stale:
repository: memory
key: never-seeded
steps:
- to: ${env:UPSTREAM}
"#
.replace("TTL", TTL)
}
async fn run_two_docs(
opener_yaml: &str,
consumer_yaml: &str,
) -> (DocumentOutcome, DocumentOutcome) {
let dir = tempfile::tempdir().expect("temp dir");
let root = dir.path();
std::fs::write(
root.join("Camel.toml"),
"log_level = \"info\"\n\n[components.http]\nallow_internal = true\n",
)
.expect("write Camel.toml");
std::fs::write(root.join("routes.yaml"), routes_yaml()).expect("write routes.yaml");
let docs: Vec<ScenarioDocument> = [opener_yaml, consumer_yaml]
.iter()
.enumerate()
.map(|(n, yaml)| {
let path = root.join(format!("case-{n}.test.yaml"));
std::fs::write(&path, yaml).expect("write case file");
parse_scenario_document(&path).expect("document must load")
})
.collect();
let scripts = partner_scripts_for(&docs[0], UPSTREAM);
let partner = match scripts {
Some(scripts) => HttpPartner::start(scripts).await,
None => HttpPartner::start_permissive(200).await,
}
.expect("partner must bind 127.0.0.1:0");
let bound = partner.bound_addr().to_string();
let harness_provisioned =
BTreeMap::from([("UPSTREAM".to_string(), format!("http://{bound}/upstream"))]);
let env = LayeredEnv::new(
docs[0].env.clone().unwrap_or_default(),
harness_provisioned,
docs[0].env_passthrough.clone().unwrap_or_default(),
ambient_std(),
);
let run = boot_scenario(&docs[0], root, &env)
.await
.expect("the full boot must succeed");
let ctx = Arc::new(Mutex::new(run.ctx));
let mut adapters: BTreeMap<String, Box<dyn PartnerAdapter>> = BTreeMap::new();
adapters.insert(
"direct:fetch".to_string(),
Box::new(DirectStimulus::new(Arc::clone(&ctx))),
);
adapters.insert(
"direct:miss".to_string(),
Box::new(DirectStimulus::new(Arc::clone(&ctx))),
);
adapters.insert(UPSTREAM.to_string(), Box::new(partner));
let router = PartnerRouter::new(adapters);
let mut outcomes = Vec::with_capacity(docs.len());
for doc in &docs {
let mut vars = camel_integration_test::ScenarioVars::new();
let outcome = run_scenario_document(doc, &router, &mut vars, None).await;
outcomes.push(outcome);
}
let [first, second, ..] = &outcomes[..] else {
panic!("exactly two documents must run");
};
let (first, second) = (first.clone(), second.clone());
let mut second = second;
if let Err(e) = run.boot.shutdown(&mut *ctx.lock().await).await {
second.final_failure = Some(ScenarioFailure::ShutdownFailure {
message: e.to_string(),
});
}
(first, second)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn circuit_fallback_serves_stale_on_upstream_failure() {
let (opener, consumer) = run_two_docs(
r#"
routeFiles: [routes.yaml]
scenario:
- send:
to: direct:fetch
body: open
- validate:
target:
partner:
endpoint: http://127.0.0.1:0/upstream
provisioning: harness
expectation:
count: 1
partners:
http://127.0.0.1:0/upstream:
- fault: close
"#,
r#"
routeFiles: [routes.yaml]
scenario:
- sleep:
duration: SLEEP
- send:
to: direct:fetch
body: fetch
expectReply:
contains: tile-stale-7f
- validate:
target:
partner:
endpoint: http://127.0.0.1:0/upstream
provisioning: harness
expectation:
count: 1
partners:
http://127.0.0.1:0/upstream:
- fault: close
"#
.replace("SLEEP", SLEEP)
.as_str(),
)
.await;
assert_eq!(
opener.verdict, None,
"the opening dial must fail its send: {opener:?}"
);
assert!(
matches!(
opener.per_action.first(),
Some(Err(ScenarioFailure::ActionTransport { .. }))
),
"the opening failure must be ActionTransport: {:?}",
opener.per_action
);
assert_eq!(
consumer.verdict,
Some(ScenarioVerdict::Pass),
"the stale fallback must serve the aged value: {consumer:?}"
);
assert!(
consumer.per_action.iter().all(|result| result.is_ok()),
"no consumer action may fail: {:?}",
consumer.per_action
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn circuit_fallback_miss_stops_cleanly() {
let (opener, consumer) = run_two_docs(
r#"
routeFiles: [routes.yaml]
scenario:
- send:
to: direct:miss
body: open
- validate:
target:
partner:
endpoint: http://127.0.0.1:0/upstream
provisioning: harness
expectation:
count: 1
partners:
http://127.0.0.1:0/upstream:
- fault: close
"#,
r#"
routeFiles: [routes.yaml]
scenario:
- send:
to: direct:miss
body: fallback
- validate:
target:
partner:
endpoint: http://127.0.0.1:0/upstream
provisioning: harness
expectation:
count: 1
partners:
http://127.0.0.1:0/upstream:
- fault: close
"#,
)
.await;
assert_eq!(
opener.verdict, None,
"the opening dial must fail its send: {opener:?}"
);
assert!(
matches!(
opener.per_action.first(),
Some(Err(ScenarioFailure::ActionTransport { .. }))
),
"the opening failure must be ActionTransport: {:?}",
opener.per_action
);
assert_eq!(
consumer.verdict,
Some(ScenarioVerdict::Pass),
"a fallback MISS is a clean Stopped outcome, not a failure: {consumer:?}"
);
assert!(
consumer.per_action.iter().all(|result| result.is_ok()),
"the MISS must never surface as an action Err: {:?}",
consumer.per_action
);
}