use std::collections::BTreeMap;
use std::time::Duration;
use camel_api::Value;
use futures::future::BoxFuture;
#[cfg(feature = "http")]
use crate::adapters::ReceiveTimeout;
use crate::adapters::{
FakeAdapter, IncomingMessage, OutgoingMessage, PartnerAdapter, PartnerRouter, ReceiveError,
TransportError,
};
use crate::document::{
EndpointRef, Expectation, Provisioning, RouteSource, ScenarioAction, ScenarioDocument,
ScenarioTarget, ValidateExpectation,
};
use crate::runner::{
DocumentOutcome, ScenarioFailure, ScenarioVars, ScenarioVerdict, fill_bind_vars,
interpolate_value, resolve_placeholders, run_scenario, run_scenario_document,
};
#[cfg(feature = "http")]
use crate::adapters::http::{HttpPartner, HttpWireRequest};
#[cfg(feature = "http")]
use crate::document::{CountBound, PartnerExpectation, PathFilter};
#[cfg(feature = "http")]
use crate::runner::{matching_requests, partner_mismatch_detail, render_bound, render_filters};
fn endpoint(uri: &str) -> EndpointRef {
EndpointRef {
endpoint: uri.to_string(),
provisioning: None,
bind_var: None,
}
}
fn doc_with(actions: Vec<ScenarioAction>) -> ScenarioDocument {
ScenarioDocument {
route_source: RouteSource::RouteFiles(vec!["routes.yaml".into()]),
scenario: actions,
partners: None,
env: None,
env_passthrough: None,
profile: None,
}
}
fn router_for(uri: &str, fake: FakeAdapter) -> PartnerRouter {
PartnerRouter::new(BTreeMap::from([(
uri.to_string(),
Box::new(fake) as Box<dyn PartnerAdapter>,
)]))
}
fn text_message(body: &str) -> IncomingMessage {
IncomingMessage {
body: Value::String(body.to_string()),
headers: BTreeMap::new(),
status: None,
method: None,
path: None,
arrival: std::time::Instant::now(),
}
}
#[tokio::test]
async fn send_then_receive_within_deadline() {
let fake = FakeAdapter::scripted(vec![text_message("hello")]);
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![
ScenarioAction::Send {
to: endpoint("partner://fake"),
body: Some(Value::String("hello".to_string())),
headers: None,
method: "POST".to_string(),
},
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: None,
},
ScenarioAction::Validate {
target: ScenarioTarget::LastReceived(endpoint("partner://fake")),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"hello".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
]);
let mut vars = ScenarioVars::new();
let verdict = run_scenario(&doc, &router, &mut vars).await;
assert_eq!(verdict, Ok(ScenarioVerdict::Pass));
}
#[tokio::test]
async fn receive_timeout_is_verdict_failure() {
let fake = FakeAdapter::scripted(Vec::new());
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_millis(50),
extract: None,
}]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("empty queue must time out");
assert!(
matches!(failure, ScenarioFailure::ReceiveTimeout { .. }),
"expected ReceiveTimeout, got {failure:?}"
);
assert!(
failure.to_string().starts_with("receive-timeout"),
"error must name the receive-timeout class: {failure}"
);
}
#[tokio::test]
async fn variable_extraction_flows_forward() {
fn scripted_with_id(id: &str) -> FakeAdapter {
FakeAdapter::scripted(vec![IncomingMessage {
body: Value::String("payload".to_string()),
headers: BTreeMap::from([("X-Id".to_string(), Value::String(id.to_string()))]),
status: None,
method: None,
path: None,
arrival: std::time::Instant::now(),
}])
}
fn extraction_doc() -> ScenarioDocument {
doc_with(vec![
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: Some(BTreeMap::from([(
"id".to_string(),
"headers.X-Id".to_string(),
)])),
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("id".to_string()),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"abc-123".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
])
}
let router = router_for("partner://fake", scripted_with_id("abc-123"));
let mut vars = ScenarioVars::new();
let verdict = run_scenario(&extraction_doc(), &router, &mut vars).await;
assert_eq!(verdict, Ok(ScenarioVerdict::Pass));
assert_eq!(
vars.get("id"),
Some(&Value::String("abc-123".to_string())),
"extraction must persist the variable for later actions"
);
let router = router_for("partner://fake", scripted_with_id("nope"));
let mut vars = ScenarioVars::new();
let failure = run_scenario(&extraction_doc(), &router, &mut vars)
.await
.expect_err("mismatched header must fail validation");
assert!(
matches!(
failure,
ScenarioFailure::ValidationMismatch { action: 1, .. }
),
"expected ValidationMismatch on action 1, got {failure:?}"
);
}
#[tokio::test]
async fn transport_error_is_apparatus_failure() {
let fake = FakeAdapter::failing_send("connection refused");
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![ScenarioAction::Send {
to: endpoint("partner://fake"),
body: None,
headers: None,
method: "GET".to_string(),
}]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("failing send must fail the scenario");
assert!(
matches!(failure, ScenarioFailure::ActionTransport { action: 0, .. }),
"expected ActionTransport on action 0, got {failure:?}"
);
assert!(
failure.to_string().starts_with("action-transport-failure"),
"error must name the action-transport-failure class: {failure}"
);
}
#[tokio::test]
async fn receive_transport_error_is_apparatus_failure() {
let fake = FakeAdapter::failing_receive("connection reset");
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: None,
}]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("failing receive must fail the scenario");
assert!(
matches!(failure, ScenarioFailure::ActionTransport { action: 0, .. }),
"expected ActionTransport on action 0, got {failure:?}"
);
assert!(
failure.to_string().starts_with("action-transport-failure"),
"error must name the action-transport-failure class: {failure}"
);
}
#[tokio::test]
async fn router_dispatches_and_fake_records_sends() {
let fake = FakeAdapter::scripted(Vec::new());
let handle = fake.recorder();
let router = router_for("partner://fake", fake);
let sent = OutgoingMessage {
body: Value::String("recorded".to_string()),
headers: BTreeMap::from([("X-Trace".to_string(), Value::String("t1".to_string()))]),
method: "POST".to_string(),
};
router
.send("partner://fake", "partner://fake", sent)
.await
.expect("send must succeed");
let recorded = handle.sent_messages();
assert_eq!(recorded.len(), 1);
assert_eq!(recorded[0].endpoint, "partner://fake");
assert_eq!(
recorded[0].message.body,
Value::String("recorded".to_string())
);
let err = router
.send(
"partner://other",
"partner://other",
OutgoingMessage {
body: Value::Null,
headers: BTreeMap::new(),
method: "GET".to_string(),
},
)
.await
.expect_err("unbound endpoint must fail");
assert!(matches!(err, TransportError::Unbound { .. }));
let failure = router
.receive(
"partner://other",
"partner://other",
Duration::from_secs(30),
)
.await
.expect_err("unbound endpoint must never deliver");
assert!(
matches!(
failure,
ReceiveError::Transport(TransportError::Unbound { .. })
),
"expected Transport(Unbound), got {failure:?}"
);
}
#[tokio::test]
async fn selector_extracts_status_method_and_path() {
let fake = FakeAdapter::scripted(vec![IncomingMessage {
body: Value::String("payload".to_string()),
headers: BTreeMap::new(),
status: Some(201),
method: Some("POST".to_string()),
path: Some("/orders".to_string()),
arrival: std::time::Instant::now(),
}]);
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: Some(BTreeMap::from([
("status".to_string(), "status".to_string()),
("method".to_string(), "method".to_string()),
("path".to_string(), "path".to_string()),
])),
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("status".to_string()),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::Number(
201.into(),
))),
deadline: None,
elapsed_at_least: None,
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("method".to_string()),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"POST".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("path".to_string()),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"/orders".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
]);
let mut vars = ScenarioVars::new();
let verdict = run_scenario(&doc, &router, &mut vars).await;
assert_eq!(verdict, Ok(ScenarioVerdict::Pass));
}
#[tokio::test]
async fn selector_header_lookup_is_case_insensitive() {
fn scripted_header(header_key: &str) -> FakeAdapter {
FakeAdapter::scripted(vec![IncomingMessage {
body: Value::Null,
headers: BTreeMap::from([(header_key.to_string(), Value::String("t-42".to_string()))]),
status: None,
method: None,
path: None,
arrival: std::time::Instant::now(),
}])
}
fn doc(selector: &str) -> ScenarioDocument {
doc_with(vec![
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: Some(BTreeMap::from([(
"trace".to_string(),
selector.to_string(),
)])),
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("trace".to_string()),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"t-42".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
])
}
let router = router_for("partner://fake", scripted_header("X-Trace"));
let mut vars = ScenarioVars::new();
let verdict = run_scenario(&doc("headers.X-Trace"), &router, &mut vars).await;
assert_eq!(verdict, Ok(ScenarioVerdict::Pass));
let router = router_for("partner://fake", scripted_header("x-trace"));
let mut vars = ScenarioVars::new();
let verdict = run_scenario(&doc("headers.X-Trace"), &router, &mut vars).await;
assert_eq!(verdict, Ok(ScenarioVerdict::Pass));
let router = router_for("partner://fake", scripted_header("X-Trace"));
let mut vars = ScenarioVars::new();
let verdict = run_scenario(&doc("headers.x-trace"), &router, &mut vars).await;
assert_eq!(verdict, Ok(ScenarioVerdict::Pass));
}
#[tokio::test]
async fn document_run_all_pass_records_verdict() {
let fake = FakeAdapter::scripted(vec![text_message("one"), text_message("two")]);
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: None,
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("unset".to_string()),
expectation: ValidateExpectation::Message(Expectation::Exists),
deadline: None,
elapsed_at_least: None,
},
]);
let mut vars = ScenarioVars::new();
vars.set("unset", Value::String("set".to_string()));
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert_eq!(
outcome,
DocumentOutcome {
per_action: vec![Ok(ScenarioVerdict::Pass), Ok(ScenarioVerdict::Pass)],
verdict: Some(ScenarioVerdict::Pass),
final_failure: None,
}
);
}
#[tokio::test]
async fn document_run_stops_at_first_failure() {
let fake = FakeAdapter::scripted(vec![text_message("one")]);
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: None,
},
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_millis(50),
extract: None,
},
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: None,
},
]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert_eq!(outcome.per_action.len(), 2, "only two actions ran");
assert_eq!(outcome.per_action[0], Ok(ScenarioVerdict::Pass));
assert!(matches!(
outcome.per_action[1],
Err(ScenarioFailure::ReceiveTimeout { .. })
));
assert_eq!(outcome.verdict, None, "no verdict after a failure");
assert_eq!(outcome.final_failure, None);
assert!(
vars.last_received("partner://fake").is_some(),
"executed actions' side effects must persist"
);
}
#[tokio::test]
async fn variable_mismatch_names_the_variable() {
let fake = FakeAdapter::scripted(vec![IncomingMessage {
body: Value::Null,
headers: BTreeMap::from([(
"X-Order-Type".to_string(),
Value::String("priority".to_string()),
)]),
status: None,
method: None,
path: None,
arrival: std::time::Instant::now(),
}]);
let router = router_for("partner://fake", fake);
let doc = doc_with(vec![
ScenarioAction::Receive {
from: endpoint("partner://fake"),
deadline: Duration::from_secs(1),
extract: Some(BTreeMap::from([(
"orderType".to_string(),
"headers.X-Order-Type".to_string(),
)])),
},
ScenarioAction::Validate {
target: ScenarioTarget::Variable("orderType".to_string()),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"express".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert_eq!(outcome.verdict, None);
match &outcome.per_action[1] {
Err(ScenarioFailure::ValidationMismatch { action: 1, detail }) => {
assert!(
detail.contains("orderType"),
"mismatch must name the variable: {detail}"
);
assert!(
detail.contains("express") && detail.contains("priority"),
"mismatch must show expected and actual: {detail}"
);
}
other => panic!("expected ValidationMismatch on action 1, got {other:?}"),
}
}
#[test]
fn partner_adapter_trait_object_is_send_sync() {
fn assert_send_sync<T: Send + Sync + ?Sized>() {}
assert_send_sync::<dyn PartnerAdapter>();
assert_send_sync::<Box<dyn PartnerAdapter>>();
}
#[test]
fn resolve_substitutes_known_var() {
let mut vars = ScenarioVars::new();
vars.set("PARTNER", Value::String("127.0.0.1:9".to_string()));
assert_eq!(
resolve_placeholders("http://${PARTNER}/orders", &vars),
Ok("http://127.0.0.1:9/orders".to_string())
);
}
#[test]
fn resolve_escape_yields_literal() {
let vars = ScenarioVars::new();
assert_eq!(
resolve_placeholders("$${not_a_var}", &vars),
Ok("${not_a_var}".to_string())
);
}
#[test]
fn resolve_unset_var_names_it() {
let vars = ScenarioVars::new();
assert_eq!(
resolve_placeholders("${missing}", &vars),
Err(ScenarioFailure::VarUnresolved {
name: "missing".to_string()
})
);
}
#[test]
fn resolve_non_string_stringifies() {
let mut vars = ScenarioVars::new();
vars.set("N", Value::Number(42.into()));
assert_eq!(resolve_placeholders("${N}", &vars), Ok("42".to_string()));
}
#[test]
fn resolve_invalid_name_stays_literal() {
let vars = ScenarioVars::new();
assert_eq!(
resolve_placeholders("${a-b}", &vars),
Ok("${a-b}".to_string())
);
}
#[test]
fn resolve_env_placeholder_stays_literal() {
let vars = ScenarioVars::new();
assert_eq!(
resolve_placeholders("${env:FOO}", &vars),
Ok("${env:FOO}".to_string())
);
}
#[test]
fn interpolate_walks_nested_leaves() {
let mut vars = ScenarioVars::new();
vars.set("x", Value::String("1".to_string()));
vars.set("y", Value::String("2".to_string()));
let body = Value::Object(
[
(
"a".to_string(),
Value::Array(vec![
Value::String("${x}".to_string()),
Value::Number(1.into()),
]),
),
(
"b".to_string(),
Value::Object(
[("c".to_string(), Value::String("${y}".to_string()))]
.into_iter()
.collect(),
),
),
]
.into_iter()
.collect(),
);
let expected = Value::Object(
[
(
"a".to_string(),
Value::Array(vec![
Value::String("1".to_string()),
Value::Number(1.into()),
]),
),
(
"b".to_string(),
Value::Object(
[("c".to_string(), Value::String("2".to_string()))]
.into_iter()
.collect(),
),
),
]
.into_iter()
.collect(),
);
assert_eq!(interpolate_value(&body, &vars), Ok(expected));
}
#[test]
fn interpolate_unset_in_body_propagates() {
let vars = ScenarioVars::new();
let body = Value::Object(
[(
"a".to_string(),
Value::Object(
[("b".to_string(), Value::String("${missing}".to_string()))]
.into_iter()
.collect(),
),
)]
.into_iter()
.collect(),
);
assert_eq!(
interpolate_value(&body, &vars),
Err(ScenarioFailure::VarUnresolved {
name: "missing".to_string()
})
);
}
struct StaticAuthority(&'static str);
impl PartnerAdapter for StaticAuthority {
fn receive<'a>(
&'a self,
_lane_key: &'a str,
source_uri: &'a str,
_deadline: Duration,
) -> BoxFuture<'a, Result<IncomingMessage, ReceiveError>> {
Box::pin(async move {
Err(ReceiveError::Transport(TransportError::Other {
message: format!("{source_uri} has no receive role in this test"),
}))
})
}
fn bound_authority(&self) -> Option<String> {
Some(self.0.to_string())
}
}
#[tokio::test]
#[cfg(feature = "http")]
async fn send_interpolates_endpoint() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let recorder = partner.recorder();
let authority = partner.bound_addr().to_string();
let router = PartnerRouter::new(BTreeMap::from([(
"http://127.0.0.1:0/orders".to_string(),
Box::new(partner) as Box<dyn PartnerAdapter>,
)]));
let doc = doc_with(vec![
ScenarioAction::Send {
to: endpoint("http://${PARTNER}/orders"),
body: None,
headers: None,
method: "POST".to_string(),
},
ScenarioAction::Receive {
from: endpoint("http://${PARTNER}/orders"),
deadline: Duration::from_secs(5),
extract: None,
},
]);
let mut vars = ScenarioVars::new();
vars.set("PARTNER", Value::String(authority));
run_scenario(&doc, &router, &mut vars)
.await
.expect("the interpolated send must reach the partner");
let recorded = recorder.recorded_requests();
assert_eq!(recorded.len(), 1, "exactly one request must reach the wire");
assert_eq!(recorded[0].path, "/orders");
}
#[tokio::test]
#[cfg(feature = "http")]
async fn send_interpolates_body_and_headers() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let recorder = partner.recorder();
let uri = format!("http://{}/orders", partner.bound_addr());
let router = PartnerRouter::new(BTreeMap::from([(
uri.clone(),
Box::new(partner) as Box<dyn PartnerAdapter>,
)]));
let doc = doc_with(vec![
ScenarioAction::Send {
to: endpoint(&uri),
body: Some(Value::Object(
[("sku".to_string(), Value::String("${SKU}".to_string()))]
.into_iter()
.collect(),
)),
headers: Some(BTreeMap::from([(
"X-Trace".to_string(),
Value::String("${SKU}".to_string()),
)])),
method: "POST".to_string(),
},
ScenarioAction::Receive {
from: endpoint(&uri),
deadline: Duration::from_secs(5),
extract: None,
},
]);
let mut vars = ScenarioVars::new();
vars.set("SKU", Value::String("x1".to_string()));
run_scenario(&doc, &router, &mut vars)
.await
.expect("the send must reach the partner");
let recorded = recorder.recorded_requests();
assert_eq!(recorded.len(), 1);
assert!(
String::from_utf8_lossy(&recorded[0].body).contains("x1"),
"recorded body must carry the substituted SKU: {:?}",
recorded[0].body,
);
assert_eq!(
recorded[0].headers.get("x-trace").map(String::as_str),
Some("x1"),
"recorded header must carry the substituted value"
);
}
#[test]
fn fill_bind_vars_sets_authority_without_scheme() {
let uri = "http://127.0.0.1:0/orders";
let router = PartnerRouter::new(BTreeMap::from([(
uri.to_string(),
Box::new(StaticAuthority("127.0.0.1:45678")) as Box<dyn PartnerAdapter>,
)]));
let wired = vec![EndpointRef {
endpoint: uri.to_string(),
provisioning: Some(Provisioning::Harness),
bind_var: Some("PARTNER".to_string()),
}];
let mut vars = ScenarioVars::new();
fill_bind_vars(&wired, &router, &mut vars);
assert_eq!(
vars.get("PARTNER"),
Some(&Value::String("127.0.0.1:45678".to_string())),
"the bind variable must carry the bare host:port authority"
);
}
#[cfg(feature = "http")]
const ORDERS: &str = "http://127.0.0.1:0/orders";
#[cfg(feature = "http")]
fn orders_router(partner: HttpPartner) -> PartnerRouter {
PartnerRouter::new(BTreeMap::from([(
ORDERS.to_string(),
Box::new(partner) as Box<dyn PartnerAdapter>,
)]))
}
#[cfg(feature = "http")]
fn orders_send() -> ScenarioAction {
ScenarioAction::Send {
to: endpoint(ORDERS),
body: None,
headers: None,
method: "POST".to_string(),
}
}
#[cfg(feature = "http")]
fn partner_validate(
count: u64,
method: Option<&str>,
path: Option<&str>,
deadline: Option<Duration>,
) -> ScenarioAction {
ScenarioAction::Validate {
target: ScenarioTarget::Partner(endpoint(ORDERS)),
expectation: ValidateExpectation::Partner(PartnerExpectation {
bound: CountBound::Exact(count),
method: method.map(str::to_string),
path: path.map(|path| PathFilter::Exact(path.to_string())),
query: None,
}),
deadline,
elapsed_at_least: None,
}
}
#[cfg(feature = "http")]
async fn raw_request(authority: &str, method: &str, path: &str) {
use tokio::io::AsyncReadExt;
use tokio::io::AsyncWriteExt;
let mut stream = tokio::net::TcpStream::connect(authority)
.await
.expect("the partner's bound address must accept");
let request = format!(
"{method} {path} HTTP/1.1\r\nhost: {authority}\r\nconnection: close\r\ncontent-length: 0\r\n\r\n"
);
stream
.write_all(request.as_bytes())
.await
.expect("the raw request must leave");
let mut sink = Vec::new();
stream
.read_to_end(&mut sink)
.await
.expect("the partner must close after its response");
}
#[cfg(feature = "http")]
fn first_failure(outcome: &DocumentOutcome) -> &ScenarioFailure {
outcome
.per_action
.iter()
.find_map(|result| result.as_ref().err())
.expect("the document must have failed")
}
#[cfg(feature = "http")]
fn wire(method: &str, path: &str) -> HttpWireRequest {
HttpWireRequest {
method: method.to_string(),
path: path.to_string(),
headers: BTreeMap::new(),
body: Vec::new(),
}
}
#[cfg(feature = "http")]
fn wire_get(path: &str) -> HttpWireRequest {
wire("GET", path)
}
#[test]
#[cfg(feature = "http")]
fn matching_requests_filters_method_case_insensitive_and_exact_path() {
let requests = vec![
wire("POST", "/orders"),
wire("GET", "/orders"),
wire("GET", "/orders?page=2"),
wire("GET", "/health"),
wire("delete", "/orders"),
];
assert_eq!(matching_requests(&requests, None, None, None), 5);
assert_eq!(matching_requests(&requests, Some("get"), None, None), 3);
assert_eq!(matching_requests(&requests, Some("DELETE"), None, None), 1);
assert_eq!(
matching_requests(
&requests,
None,
Some(&PathFilter::Exact("/orders".to_string())),
None
),
3
);
assert_eq!(
matching_requests(
&requests,
None,
Some(&PathFilter::Exact("/orders?page=2".to_string())),
None
),
1
);
assert_eq!(
matching_requests(
&requests,
Some("get"),
Some(&PathFilter::Exact("/orders".to_string())),
None
),
1
);
}
#[test]
#[cfg(feature = "http")]
fn matching_exact_path_is_byte_strict() {
let requests = vec![wire_get("/q?bbox=1.5%2C2.5"), wire_get("/q?bbox=1.5,2.5")];
assert_eq!(
matching_requests(
&requests,
None,
Some(&PathFilter::Exact("/q?bbox=1.5%2C2.5".to_string())),
None
),
1
);
}
#[test]
#[cfg(feature = "http")]
fn matching_contains_tolerates_encoding() {
let requests = vec![wire_get("/q?bbox=1.5%2C2.5"), wire_get("/q?bbox=1.5,2.5")];
assert_eq!(
matching_requests(
&requests,
None,
Some(&PathFilter::Contains("bbox=".to_string())),
None
),
2
);
}
#[test]
#[cfg(feature = "http")]
fn matching_regex_narrows() {
let requests = vec![wire_get("/orders/42"), wire_get("/health")];
assert_eq!(
matching_requests(
&requests,
None,
Some(&PathFilter::Matches("^/orders/\\d+$".to_string())),
None
),
1
);
}
#[test]
#[cfg(feature = "http")]
fn matching_query_subset_decodes_and_ignores_order() {
let requests = vec![wire_get("/q?b=2&a=1%2B1")];
let query = BTreeMap::from([
("a".to_string(), "1+1".to_string()),
("b".to_string(), "2".to_string()),
]);
assert_eq!(matching_requests(&requests, None, None, Some(&query)), 1);
}
#[test]
#[cfg(feature = "http")]
fn matching_query_subset_absent_pair_excludes() {
let requests = vec![wire_get("/q?a=1")];
let query = BTreeMap::from([
("a".to_string(), "1".to_string()),
("c".to_string(), "3".to_string()),
]);
assert_eq!(matching_requests(&requests, None, None, Some(&query)), 0);
}
#[test]
#[cfg(feature = "http")]
fn matching_method_composes_with_query() {
let requests = vec![wire("POST", "/q?a=1"), wire("GET", "/q?a=1")];
let query = BTreeMap::from([("a".to_string(), "1".to_string())]);
assert_eq!(
matching_requests(&requests, Some("post"), None, Some(&query)),
1
);
}
#[tokio::test]
#[cfg(feature = "http")]
async fn immediate_count_passes_and_mismatch_names_counts() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let router = orders_router(partner);
let doc = doc_with(vec![
orders_send(),
ScenarioAction::Receive {
from: endpoint(ORDERS),
deadline: Duration::from_secs(5),
extract: None,
},
partner_validate(1, None, None, None),
partner_validate(2, None, None, None),
]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert!(
matches!(outcome.per_action[2], Ok(ScenarioVerdict::Pass)),
"the exact immediate count must pass: {outcome:?}"
);
assert_eq!(outcome.verdict, None, "the count: 2 validate must fail");
let ScenarioFailure::ValidationMismatch { detail, .. } = first_failure(&outcome) else {
panic!(
"expected ValidationMismatch, got {:?}",
first_failure(&outcome)
);
};
assert!(
detail.contains("partner http://127.0.0.1:0/orders"),
"the mismatch must name the partner URI: {detail}"
);
assert!(
detail.contains("expected 2, actual 1"),
"the mismatch must name both counts: {detail}"
);
}
#[tokio::test]
#[cfg(feature = "http")]
async fn filtered_mismatch_names_method_and_path_clauses() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let router = orders_router(partner);
let doc = doc_with(vec![
orders_send(),
ScenarioAction::Receive {
from: endpoint(ORDERS),
deadline: Duration::from_secs(5),
extract: None,
},
partner_validate(2, Some("post"), Some("/orders"), None),
]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert_eq!(outcome.verdict, None, "the filtered count must fail");
let ScenarioFailure::ValidationMismatch { detail, .. } = first_failure(&outcome) else {
panic!(
"expected ValidationMismatch, got {:?}",
first_failure(&outcome)
);
};
assert!(
detail.contains("method post"),
"the mismatch must name the method filter: {detail}"
);
assert!(
detail.contains("path /orders"),
"the mismatch must name the path filter: {detail}"
);
assert!(
detail.contains("expected 2, actual 1"),
"the mismatch must name both counts: {detail}"
);
}
#[tokio::test]
#[cfg(feature = "http")]
async fn poll_passes_once_count_settles() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let authority = partner.bound_addr().to_string();
raw_request(&authority, "POST", "/orders").await;
let router = orders_router(partner);
let settling = tokio::spawn({
let authority = authority.clone();
async move {
tokio::time::sleep(Duration::from_millis(300)).await;
raw_request(&authority, "POST", "/orders").await;
raw_request(&authority, "POST", "/orders").await;
}
});
let doc = doc_with(vec![partner_validate(
3,
None,
None,
Some(Duration::from_secs(5)),
)]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
settling.await.expect("the settling task must finish");
assert_eq!(
outcome.verdict,
Some(ScenarioVerdict::Pass),
"the polled count must settle to 3: {outcome:?}"
);
}
#[tokio::test]
#[cfg(feature = "http")]
async fn overshoot_never_passes() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let authority = partner.bound_addr().to_string();
for _ in 0..4 {
raw_request(&authority, "POST", "/orders").await;
}
let router = orders_router(partner);
let doc = doc_with(vec![partner_validate(
3,
None,
None,
Some(Duration::from_secs(1)),
)]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert_eq!(
outcome.verdict, None,
"a count above the expectation must never pass: {outcome:?}"
);
let ScenarioFailure::ValidationMismatch { detail, .. } = first_failure(&outcome) else {
panic!(
"expected ValidationMismatch, got {:?}",
first_failure(&outcome)
);
};
assert!(
detail.contains("expected 3, actual 4"),
"the mismatch must name the final counts: {detail}"
);
}
#[tokio::test]
#[cfg(feature = "http")]
async fn deadline_expiry_reports_final_actual() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let authority = partner.bound_addr().to_string();
raw_request(&authority, "POST", "/orders").await;
let router = orders_router(partner);
let doc = doc_with(vec![partner_validate(
3,
None,
None,
Some(Duration::from_secs(1)),
)]);
let mut vars = ScenarioVars::new();
let outcome = run_scenario_document(&doc, &router, &mut vars).await;
assert_eq!(outcome.verdict, None, "the count must never reach 3");
let ScenarioFailure::ValidationMismatch { detail, .. } = first_failure(&outcome) else {
panic!(
"expected ValidationMismatch, got {:?}",
first_failure(&outcome)
);
};
assert!(
detail.contains("actual 1"),
"the mismatch must report the final snapshot's count: {detail}"
);
}
#[test]
#[cfg(feature = "http")]
fn partner_mismatch_detail_lists_recorded_paths() {
let expected = PartnerExpectation {
bound: CountBound::Exact(2),
method: None,
path: None,
query: None,
};
let detail = partner_mismatch_detail(
"http://127.0.0.1:0/a",
&expected,
1,
&["/a?b=1".to_string(), "/c".to_string()],
&[],
);
assert!(
detail.contains("expected 2, actual 1"),
"the mismatch must name both counts: {detail}"
);
assert!(
detail.contains("/a?b=1"),
"must list the first path: {detail}"
);
assert!(detail.contains("/c"), "must list the second path: {detail}");
}
#[test]
#[cfg(feature = "http")]
fn count_mismatch_redacts_secrets() {
let expected = PartnerExpectation {
bound: CountBound::Exact(2),
method: None,
path: Some(PathFilter::Exact(
"/login?authPassword=hunter2&x=1".to_string(),
)),
query: None,
};
let detail = partner_mismatch_detail(
"http://127.0.0.1:0/login?authPassword=hunter2&x=1",
&expected,
1,
&["/login?authPassword=hunter2&x=1".to_string()],
&["authPassword".to_string()],
);
assert!(
detail.contains("authPassword=***"),
"the secret value must be masked: {detail}"
);
assert!(
!detail.contains("hunter2"),
"the secret must never print: {detail}"
);
assert!(
detail.contains("x=1"),
"non-secret pairs must stay visible: {detail}"
);
assert!(
detail.contains("partner http://127.0.0.1:0/login?authPassword=***&x=1"),
"the partner URI header must mask the secret too: {detail}"
);
assert!(
detail.contains("path /login?authPassword=***"),
"the path filter echo must mask the secret too: {detail}"
);
}
#[test]
#[cfg(feature = "http")]
fn render_bound_grammar() {
assert_eq!(render_bound(&CountBound::Exact(3)), "expected 3");
assert_eq!(render_bound(&CountBound::AtLeast(3)), "expected at least 3");
assert_eq!(render_bound(&CountBound::AtMost(2)), "expected at most 2");
assert_eq!(
render_bound(&CountBound::Range(2, 4)),
"expected between 2 and 4"
);
}
#[test]
#[cfg(feature = "http")]
fn render_filters_redacts_secret_query_and_elides_patterns() {
let expected = PartnerExpectation {
bound: CountBound::AtLeast(1),
method: Some("GET".to_string()),
path: Some(PathFilter::Contains("secret".to_string())),
query: Some(BTreeMap::from([
("bbox".to_string(), "1,2".to_string()),
("token".to_string(), "abc".to_string()),
])),
};
let rendered = render_filters(&expected, &["token".to_string()]);
assert!(
rendered.contains("token=<redacted>"),
"the secret pair must mask its value: {rendered}"
);
assert!(
rendered.contains("bbox=1,2"),
"the non-secret pair must stay visible: {rendered}"
);
assert!(
rendered.contains("method GET"),
"the method clause must render: {rendered}"
);
assert!(
rendered.contains("pathContains <pattern elided>"),
"the pattern must render by kind only: {rendered}"
);
assert!(
!rendered.contains("abc"),
"the secret value must never print: {rendered}"
);
assert!(
!rendered.contains("secret"),
"the pattern payload must never print: {rendered}"
);
}
#[cfg(feature = "http")]
struct CannedTimeout {
endpoint: String,
lanes_recorded: Vec<String>,
}
#[cfg(feature = "http")]
impl PartnerAdapter for CannedTimeout {
fn receive<'a>(
&'a self,
_lane_key: &'a str,
_source_uri: &'a str,
deadline: Duration,
) -> BoxFuture<'a, Result<IncomingMessage, ReceiveError>> {
Box::pin(async move {
Err(ReceiveError::Timeout(ReceiveTimeout {
endpoint: self.endpoint.clone(),
deadline,
elapsed: Duration::ZERO,
lanes_recorded: self.lanes_recorded.clone(),
}))
})
}
}
#[tokio::test]
#[cfg(feature = "http")]
async fn receive_timeout_failure_carries_redacted_endpoint() {
let declared = "http://host/login?authPassword=hunter2&x=1";
let router = PartnerRouter::new(BTreeMap::from([(
declared.to_string(),
Box::new(CannedTimeout {
endpoint: "http://host/login?authPassword=***&x=1".to_string(),
lanes_recorded: vec!["/login?authPassword=***&x=1".to_string()],
}) as Box<dyn PartnerAdapter>,
)]));
let doc = doc_with(vec![ScenarioAction::Receive {
from: endpoint(declared),
deadline: Duration::from_millis(50),
extract: None,
}]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("the receive must time out");
let text = failure.to_string();
assert!(
!text.contains("hunter2"),
"the raw secret must never print: {text}"
);
assert!(
text.contains("authPassword=***"),
"the redacted form must print: {text}"
);
assert!(
text.contains("x=1"),
"non-secret query keys must stay visible: {text}"
);
}
#[tokio::test]
#[cfg(feature = "http")]
async fn raw_adapter_timeout_redacts_at_the_mapping() {
let declared = "http://host/login?authPassword=hunter2&x=1";
let router = PartnerRouter::new(BTreeMap::from([(
declared.to_string(),
Box::new(CannedTimeout {
endpoint: "http://host/login?authPassword=hunter2&x=1".to_string(),
lanes_recorded: vec!["/login?authPassword=hunter2&x=1".to_string()],
}) as Box<dyn PartnerAdapter>,
)]));
router.set_secret_query_keys(vec!["authPassword".to_string()]);
let doc = doc_with(vec![ScenarioAction::Receive {
from: endpoint(declared),
deadline: Duration::from_millis(50),
extract: None,
}]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("the receive must time out");
let text = failure.to_string();
assert!(
!text.contains("hunter2"),
"the raw secret must never print: {text}"
);
assert!(
text.contains("authPassword=***"),
"the redacted form must print: {text}"
);
assert!(
text.contains("x=1"),
"non-secret query keys must stay visible: {text}"
);
}
#[tokio::test]
async fn body_validation_failure_carries_redacted_subject() {
let declared = "http://host/login?authPassword=hunter2&x=1";
let fake = FakeAdapter::scripted(vec![text_message("mismatch-me")]);
let router = router_for(declared, fake);
router.set_secret_query_keys(vec!["authPassword".to_string()]);
let doc = doc_with(vec![
ScenarioAction::Receive {
from: endpoint(declared),
deadline: Duration::from_secs(1),
extract: None,
},
ScenarioAction::Validate {
target: ScenarioTarget::LastReceived(endpoint(declared)),
expectation: ValidateExpectation::Message(Expectation::Equals(Value::String(
"expected".to_string(),
))),
deadline: None,
elapsed_at_least: None,
},
]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("the body validation must fail");
let text = failure.to_string();
assert!(
!text.contains("hunter2"),
"the raw secret must never print: {text}"
);
assert!(
text.contains("authPassword=***"),
"the redacted form must print: {text}"
);
assert!(
text.contains("x=1"),
"non-secret query keys must stay visible: {text}"
);
}
#[tokio::test]
async fn unreceived_validate_carries_redacted_subject() {
let declared = "http://host/login?authPassword=hunter2&x=1";
let router = router_for(declared, FakeAdapter::scripted(vec![]));
router.set_secret_query_keys(vec!["authPassword".to_string()]);
let doc = doc_with(vec![ScenarioAction::Validate {
target: ScenarioTarget::LastReceived(endpoint(declared)),
expectation: ValidateExpectation::Message(Expectation::Exists),
deadline: None,
elapsed_at_least: None,
}]);
let mut vars = ScenarioVars::new();
let failure = run_scenario(&doc, &router, &mut vars)
.await
.expect_err("the validation must find no message");
let text = failure.to_string();
assert!(
!text.contains("hunter2"),
"the raw secret must never print: {text}"
);
assert!(
text.contains("authPassword=***"),
"the redacted form must print: {text}"
);
assert!(
text.contains("x=1"),
"non-secret keys must stay visible: {text}"
);
}
#[tokio::test]
async fn unbound_receive_failure_carries_redacted_endpoint() {
let declared = "http://host/login?authPassword=hunter2&x=1";
let router = PartnerRouter::new(BTreeMap::new());
router.set_secret_query_keys(vec!["authPassword".to_string()]);
let error = router
.receive(declared, declared, Duration::from_millis(1))
.await
.expect_err("no adapter is registered");
let ReceiveError::Transport(TransportError::Unbound { endpoint }) = error else {
panic!("expected an Unbound transport failure, got {error:?}");
};
assert!(
!endpoint.contains("hunter2"),
"the raw secret must never print: {endpoint}"
);
assert!(
endpoint.contains("authPassword=***"),
"the redacted form must print: {endpoint}"
);
}
#[tokio::test]
async fn unbound_send_failure_carries_redacted_endpoint() {
let declared = "http://host/login?authPassword=hunter2&x=1";
let router = PartnerRouter::new(BTreeMap::new());
router.set_secret_query_keys(vec!["authPassword".to_string()]);
let error = router
.send(
declared,
"partner://nowhere",
OutgoingMessage {
body: Value::Null,
headers: BTreeMap::new(),
method: "GET".to_string(),
},
)
.await
.expect_err("no adapter is registered");
let TransportError::Unbound { endpoint } = error else {
panic!("expected an Unbound transport failure, got {error:?}");
};
assert!(
!endpoint.contains("hunter2"),
"the raw secret must never print: {endpoint}"
);
assert!(
endpoint.contains("authPassword=***"),
"the redacted form must print: {endpoint}"
);
}