#![cfg(feature = "http")]
use std::collections::BTreeMap;
use std::time::Duration;
use camel_api::Value;
use camel_integration_test::adapters::{
IncomingMessage, OutgoingMessage, ReceiveError, TransportError,
};
use camel_integration_test::{HttpPartner, PartnerAdapter, PartnerRouter, ScriptedResponse};
fn send_msg(body: &str) -> OutgoingMessage {
OutgoingMessage {
body: Value::String(body.to_string()),
headers: BTreeMap::new(),
method: "POST".to_string(),
}
}
async fn orders_partner(body: &str) -> HttpPartner {
HttpPartner::start(vec![ScriptedResponse {
method: Some("POST".to_string()),
path: Some("/orders".to_string()),
status: 200,
headers: BTreeMap::new(),
body: body.as_bytes().to_vec(),
..Default::default()
}])
.await
.expect("partner binds 127.0.0.1:0")
}
fn orders_uri(partner: &HttpPartner) -> String {
format!("http://{}/orders", partner.bound_addr())
}
#[tokio::test]
async fn failed_send_does_not_poison_later_receive() {
let partner = orders_partner("b-response").await;
let lane_key = orders_uri(&partner);
let bound = partner.bound_addr();
let router = PartnerRouter::new(BTreeMap::from([(
lane_key.clone(),
Box::new(partner) as Box<dyn PartnerAdapter>,
)]));
let dead = tokio::net::TcpListener::bind(("127.0.0.1", 0))
.await
.expect("dead listener binds");
let dead_addr = dead.local_addr().expect("dead listener has an address");
drop(dead);
let failed = router
.send(
&lane_key,
&format!("http://{dead_addr}/orders"),
send_msg("a"),
)
.await;
assert!(
matches!(&failed, Err(TransportError::Other { message }) if message.contains("connect")),
"send A must fail inline at the transport, got {failed:?}"
);
let empty = router
.receive(&lane_key, &lane_key, Duration::from_millis(200))
.await;
assert!(
matches!(empty, Err(ReceiveError::Timeout(_))),
"no lane entry may exist after a pre-wire failure, got {empty:?}"
);
router
.send(&lane_key, &format!("http://{bound}/orders"), send_msg("b"))
.await
.expect("send B dials the live partner");
let response: IncomingMessage = router
.receive(&lane_key, &lane_key, Duration::from_secs(5))
.await
.expect("B's roundtrip is parked on the lane key");
assert_eq!(response.body, Value::String("b-response".to_string()));
}
#[tokio::test]
async fn post_connect_failure_still_parks() {
let listener = tokio::net::TcpListener::bind(("127.0.0.1", 0))
.await
.expect("listener binds");
let addr = listener.local_addr().expect("listener has an address");
tokio::spawn(async move {
if let Ok((stream, _)) = listener.accept().await {
drop(stream);
}
});
let router = PartnerRouter::new(BTreeMap::new());
let lane_key = format!("http://{addr}/orders");
router
.send(&lane_key, &lane_key, send_msg("a"))
.await
.expect("the dial succeeds against the accepting listener");
let parked = router
.receive(&lane_key, &lane_key, Duration::from_secs(5))
.await;
assert!(
matches!(&parked, Err(ReceiveError::Transport(_))),
"the post-connect failure must park on the lane key, got {parked:?}"
);
}
#[tokio::test]
async fn same_key_sends_park_fifo() {
let partner = HttpPartner::start(vec![
ScriptedResponse {
method: Some("POST".to_string()),
path: Some("/orders".to_string()),
body: b"a-response".to_vec(),
delay: Some(Duration::from_millis(50)),
..Default::default()
},
ScriptedResponse {
method: Some("POST".to_string()),
path: Some("/orders".to_string()),
body: b"b-response".to_vec(),
delay: Some(Duration::from_millis(50)),
..Default::default()
},
ScriptedResponse {
method: Some("POST".to_string()),
path: Some("/orders".to_string()),
body: b"c-response".to_vec(),
delay: Some(Duration::from_millis(50)),
..Default::default()
},
])
.await
.expect("partner binds 127.0.0.1:0");
let lane_key = orders_uri(&partner);
let router = PartnerRouter::new(BTreeMap::from([(
lane_key.clone(),
Box::new(partner) as Box<dyn PartnerAdapter>,
)]));
for body in ["a", "b", "c"] {
router
.send(&lane_key, &lane_key, send_msg(body))
.await
.expect("send parks its roundtrip on the lane key");
}
for expected in ["a-response", "b-response", "c-response"] {
let response: IncomingMessage = router
.receive(&lane_key, &lane_key, Duration::from_secs(5))
.await
.expect("the oldest parked roundtrip resolves first");
assert_eq!(
response.body,
Value::String(expected.to_string()),
"receives must drain the same-key park in wire order"
);
}
}
#[tokio::test]
async fn lane_fifo_overflow_is_apparatus() {
let partner = HttpPartner::start(vec![ScriptedResponse {
method: Some("POST".to_string()),
path: Some("/orders".to_string()),
body: b"late".to_vec(),
delay: Some(Duration::from_secs(30)),
times: 65,
..Default::default()
}])
.await
.expect("partner binds 127.0.0.1:0");
let lane_key = orders_uri(&partner);
let recorder = partner.recorder();
let router = PartnerRouter::new(BTreeMap::from([(
lane_key.clone(),
Box::new(partner) as Box<dyn PartnerAdapter>,
)]));
for _ in 0..64 {
router
.send(&lane_key, &lane_key, send_msg("bulk"))
.await
.expect("the first 64 sends book inside the FIFO");
}
let overflow = router
.send(&lane_key, &lane_key, send_msg("one-too-many"))
.await;
let Err(error) = &overflow else {
panic!("the 65th send must fail at the transport, got {overflow:?}");
};
let TransportError::LaneFifoOverflow { lane_key, bound } = error else {
panic!("the overflow must be the apparatus variant, got {error:?}");
};
assert_eq!(*bound, 64, "the overflow carries the FIFO bound");
let rendered = error.to_string();
assert!(
rendered.contains(lane_key.as_str()),
"the overflow names the lane key: {rendered}"
);
assert!(
rendered.contains("64"),
"the overflow names the bound: {rendered}"
);
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
while recorder.recorded_requests().len() < 64 {
assert!(
tokio::time::Instant::now() < deadline,
"the 64 booked sends must reach the wire"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(
recorder.recorded_requests().len(),
64,
"the refused send must reach no wire"
);
}
#[tokio::test]
async fn lane_fifo_overflow_redacts_secret_query() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner binds 127.0.0.1:0");
let uri = format!("http://{}/orders?authPassword=sekrit", partner.bound_addr());
let router = PartnerRouter::new(BTreeMap::new());
router.set_secret_query_keys(vec!["authPassword".to_string()]);
for _ in 0..64 {
router
.send(&uri, &uri, send_msg("bulk"))
.await
.expect("the first 64 sends book inside the FIFO");
}
let overflow = router.send(&uri, &uri, send_msg("one-too-many")).await;
let Err(error) = &overflow else {
panic!("the 65th send must fail at the transport, got {overflow:?}");
};
let TransportError::LaneFifoOverflow { lane_key, .. } = error else {
panic!("the overflow must be the apparatus variant, got {error:?}");
};
let rendered = error.to_string();
assert!(
rendered.contains("/orders"),
"the overflow names the path: {rendered}"
);
assert!(
lane_key.contains("authPassword=***") && rendered.contains("authPassword=***"),
"both key and path halves mask the secret: {rendered}"
);
assert!(
!rendered.contains("sekrit"),
"the overflow never echoes the secret: {rendered}"
);
}