use std::collections::BTreeMap;
use std::collections::VecDeque;
use std::io;
use std::net::SocketAddr;
use std::str::FromStr;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::RwLock;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::time::Duration;
use bytes::Bytes;
use camel_api::Value;
use futures::future::BoxFuture;
use http::HeaderName;
use http::HeaderValue;
use http::Method;
use http::Request;
use http::Response;
use http::StatusCode;
use http::Uri;
use http_body_util::BodyExt;
use http_body_util::Full;
use hyper::body::Incoming;
use hyper::client::conn::http1;
use hyper::server::conn::http1::Builder as ServerBuilder;
use hyper::service::service_fn;
use hyper_util::rt::TokioIo;
use tokio::net::TcpListener;
use tokio::net::TcpStream;
use tokio::sync::Mutex as AsyncMutex;
use tokio::sync::mpsc;
use tokio::sync::oneshot;
use tokio::sync::watch;
use crate::adapters::ArrivalLaneOverflow;
use crate::adapters::IncomingMessage;
use crate::adapters::OutgoingMessage;
use crate::adapters::PartnerAdapter;
use crate::adapters::ReceiveError;
use crate::adapters::ReceiveTimeout;
use crate::adapters::TransportError;
use crate::adapters::lock_through;
use crate::adapters::redact_wire_path;
use crate::document::PartnerFault;
const UNMATCHED_STATUS: u16 = 500;
pub(crate) const ARRIVAL_LANE_CAPACITY: usize = 64;
const LANE_FIFO_CAPACITY: usize = 64;
#[derive(Debug, Clone)]
pub struct ScriptedResponse {
pub method: Option<String>,
pub path: Option<String>,
pub times: u32,
pub delay: Option<Duration>,
pub fault: Option<PartnerFault>,
pub status: u16,
pub headers: BTreeMap<String, String>,
pub body: Vec<u8>,
}
impl Default for ScriptedResponse {
fn default() -> Self {
Self {
method: None,
path: None,
times: 1,
delay: None,
fault: None,
status: 200,
headers: BTreeMap::new(),
body: Vec::new(),
}
}
}
impl ScriptedResponse {
fn matches(&self, request: &HttpWireRequest) -> bool {
let method_ok = self
.method
.as_deref()
.is_none_or(|m| m.eq_ignore_ascii_case(&request.method));
let path_ok = self.path.as_deref().is_none_or(|p| p == request.path);
method_ok && path_ok
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct HttpWireRequest {
pub method: String,
pub path: String,
pub headers: BTreeMap<String, String>,
pub body: Vec<u8>,
}
#[derive(Clone)]
pub struct HttpRecorder {
requests: Arc<Mutex<Vec<HttpWireRequest>>>,
}
impl HttpRecorder {
pub fn recorded_requests(&self) -> Vec<HttpWireRequest> {
lock_through(&self.requests).clone()
}
}
struct ArrivalLane {
tx: mpsc::Sender<IncomingMessage>,
rx: AsyncMutex<mpsc::Receiver<IncomingMessage>>,
dropped: AtomicUsize,
}
struct ServerState {
scripted: Mutex<Vec<ScriptedResponse>>,
fallback_status: Option<u16>,
requests: Arc<Mutex<Vec<HttpWireRequest>>>,
arrivals: Mutex<BTreeMap<String, Arc<ArrivalLane>>>,
}
struct HttpInner {
bound: SocketAddr,
server: Arc<ServerState>,
shutdown: Mutex<Option<watch::Sender<bool>>>,
secret_query_keys: RwLock<Vec<String>>,
}
pub struct HttpPartner {
inner: Arc<HttpInner>,
}
impl HttpPartner {
pub async fn start(scripted: Vec<ScriptedResponse>) -> io::Result<Self> {
Self::start_with(scripted, None).await
}
pub async fn start_permissive(status: u16) -> io::Result<Self> {
Self::start_with(Vec::new(), Some(status)).await
}
async fn start_with(
scripted: Vec<ScriptedResponse>,
fallback_status: Option<u16>,
) -> io::Result<Self> {
let listener = TcpListener::bind(("127.0.0.1", 0)).await?;
let bound = listener.local_addr()?;
let requests = Arc::new(Mutex::new(Vec::new()));
let server = Arc::new(ServerState {
scripted: Mutex::new(scripted),
fallback_status,
requests: Arc::clone(&requests),
arrivals: Mutex::new(BTreeMap::new()),
});
let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
let accept_state = Arc::clone(&server);
tokio::spawn(async move {
loop {
tokio::select! {
_ = shutdown_rx.changed() => break,
accepted = listener.accept() => {
let Ok((stream, _peer)) = accepted else {
continue;
};
let server = Arc::clone(&accept_state);
tokio::spawn(async move {
let service = service_fn(move |request| {
let server = Arc::clone(&server);
async move { serve(server, request).await }
});
let _ = ServerBuilder::new()
.serve_connection(TokioIo::new(stream), service)
.await;
});
}
}
}
});
Ok(Self {
inner: Arc::new(HttpInner {
bound,
server,
shutdown: Mutex::new(Some(shutdown_tx)),
secret_query_keys: RwLock::new(Vec::new()),
}),
})
}
pub fn bound_addr(&self) -> SocketAddr {
self.inner.bound
}
pub fn recorder(&self) -> HttpRecorder {
HttpRecorder {
requests: Arc::clone(&self.inner.server.requests),
}
}
fn stored_secret_query_keys(&self) -> Vec<String> {
self.inner
.secret_query_keys
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clone()
}
fn recorded_lane_paths_redacted(&self) -> Vec<String> {
let secret_keys = self.stored_secret_query_keys();
let mut unique: Vec<String> = Vec::new();
for request in lock_through(&self.inner.server.requests).iter() {
if !unique.contains(&request.path) {
unique.push(request.path.clone());
}
}
unique
.iter()
.map(|path| redact_wire_path(path, &secret_keys))
.collect()
}
async fn await_arrival(
&self,
lane_key: &str,
source_uri: &str,
deadline: Duration,
) -> Result<IncomingMessage, ReceiveError> {
let path = ParsedTarget::parse(lane_key)
.map_err(ReceiveError::Transport)?
.target;
let lane = lane_for(&self.inner.server.arrivals, &path);
let mut rx = lane.rx.lock().await;
let started = tokio::time::Instant::now();
match tokio::time::timeout(deadline, rx.recv()).await {
Ok(Some(message)) => Ok(message),
Ok(None) | Err(_) => {
let dropped = lane.dropped.load(Ordering::Relaxed);
let secret_keys = self.stored_secret_query_keys();
if dropped > 0 {
return Err(ReceiveError::Overflow(ArrivalLaneOverflow {
endpoint: redact_wire_path(source_uri, &secret_keys),
dropped,
}));
}
Err(ReceiveError::Timeout(ReceiveTimeout {
endpoint: redact_wire_path(source_uri, &secret_keys),
deadline,
elapsed: started.elapsed(),
lanes_recorded: self.recorded_lane_paths_redacted(),
}))
}
}
}
}
impl Drop for HttpPartner {
fn drop(&mut self) {
if let Some(shutdown) = lock_through(&self.inner.shutdown).take() {
let _ = shutdown.send(true);
}
}
}
impl PartnerAdapter for HttpPartner {
fn receive<'a>(
&'a self,
lane_key: &'a str,
source_uri: &'a str,
deadline: Duration,
) -> BoxFuture<'a, Result<IncomingMessage, ReceiveError>> {
Box::pin(async move { self.await_arrival(lane_key, source_uri, deadline).await })
}
fn bound_authority(&self) -> Option<String> {
Some(self.bound_addr().to_string())
}
fn recorded_requests(&self) -> Vec<HttpWireRequest> {
self.recorder().recorded_requests()
}
fn set_secret_query_keys(&self, keys: &[String]) {
*self
.inner
.secret_query_keys
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = keys.to_vec();
}
}
struct LaneEntry {
generation: u64,
rx: oneshot::Receiver<Result<IncomingMessage, TransportError>>,
}
pub struct ClientLane {
in_flight: Arc<Mutex<BTreeMap<String, VecDeque<LaneEntry>>>>,
next_generation: AtomicU64,
launched_wire_paths: Mutex<Vec<String>>,
secret_query_keys: RwLock<Vec<String>>,
}
impl ClientLane {
pub(crate) fn new() -> Self {
Self {
in_flight: Arc::new(Mutex::new(BTreeMap::new())),
next_generation: AtomicU64::new(0),
launched_wire_paths: Mutex::new(Vec::new()),
secret_query_keys: RwLock::new(Vec::new()),
}
}
pub(crate) fn set_secret_query_keys(&self, keys: &[String]) {
*self
.secret_query_keys
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = keys.to_vec();
}
pub(crate) async fn launch(
self: Arc<Self>,
lane_key: &str,
target_uri: &str,
msg: OutgoingMessage,
) -> Result<(), TransportError> {
if lock_through(&self.in_flight)
.get(lane_key)
.is_some_and(|fifo| fifo.len() >= LANE_FIFO_CAPACITY)
{
return Err(TransportError::LaneFifoOverflow {
lane_key: lane_key.to_string(),
bound: LANE_FIFO_CAPACITY,
});
}
let target = ParsedTarget::parse(target_uri)?;
let stream = TcpStream::connect((target.host.as_str(), target.port))
.await
.map_err(|e| TransportError::Other {
message: format!("connect to {}:{} failed: {e}", target.host, target.port),
})?;
let generation = self.next_generation.fetch_add(1, Ordering::Relaxed);
let (tx, rx) = oneshot::channel();
{
let mut launched = lock_through(&self.launched_wire_paths);
if !launched.contains(&target.target) {
launched.push(target.target.clone());
}
}
{
let mut lanes = lock_through(&self.in_flight);
let fifo = lanes.entry(lane_key.to_string()).or_default();
if fifo.len() >= LANE_FIFO_CAPACITY {
return Err(TransportError::LaneFifoOverflow {
lane_key: lane_key.to_string(),
bound: LANE_FIFO_CAPACITY,
});
}
fifo.push_back(LaneEntry { generation, rx });
}
let lane = Arc::clone(&self);
let key = lane_key.to_string();
tokio::spawn(async move {
let result = perform_exchange(stream, &target, msg).await;
match result {
Ok(response) => {
let _ = tx.send(Ok(response));
}
Err(error) => {
if !lane.fail_lane_entry(&key, generation, error.clone()) {
let _ = tx.send(Err(error));
}
}
}
});
Ok(())
}
pub(crate) fn take(
&self,
lane_key: &str,
) -> Option<oneshot::Receiver<Result<IncomingMessage, TransportError>>> {
let mut lanes = lock_through(&self.in_flight);
let fifo = lanes.get_mut(lane_key)?;
let entry = fifo.pop_front()?;
if fifo.is_empty() {
lanes.remove(lane_key);
}
Some(entry.rx)
}
pub(crate) fn fail_lane_entry(
&self,
key: &str,
generation: u64,
error: TransportError,
) -> bool {
fail_lane_map_entry(&self.in_flight, key, generation, error)
}
pub(crate) async fn await_parked(
&self,
endpoint: &str,
deadline: Duration,
rx: oneshot::Receiver<Result<IncomingMessage, TransportError>>,
) -> Result<IncomingMessage, ReceiveError> {
let started = tokio::time::Instant::now();
match tokio::time::timeout(deadline, rx).await {
Err(_) => {
let secret_keys = self
.secret_query_keys
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clone();
let lanes_recorded = lock_through(&self.launched_wire_paths)
.iter()
.map(|path| redact_wire_path(path, &secret_keys))
.collect();
Err(ReceiveError::Timeout(ReceiveTimeout {
endpoint: redact_wire_path(endpoint, &secret_keys),
deadline,
elapsed: started.elapsed(),
lanes_recorded,
}))
}
Ok(Ok(result)) => result.map_err(ReceiveError::Transport),
Ok(Err(_cancelled)) => Err(ReceiveError::Transport(TransportError::Other {
message: "http request task ended without delivering a response".to_string(),
})),
}
}
}
#[derive(Debug)]
struct ParsedTarget {
host: String,
port: u16,
target: String,
}
impl ParsedTarget {
fn parse(endpoint: &str) -> Result<Self, TransportError> {
let invalid = |detail: String| TransportError::Other {
message: format!("endpoint {endpoint}: {detail}"),
};
let uri = Uri::try_from(endpoint).map_err(|e| invalid(format!("invalid uri: {e}")))?;
match uri.scheme_str() {
Some("http") => {}
other => {
return Err(invalid(format!(
"unsupported scheme {} (the http partner speaks plain http)",
other.unwrap_or("<none>")
)));
}
}
let host = uri
.host()
.ok_or_else(|| invalid("no host".to_string()))?
.to_string();
let port = uri.port_u16().unwrap_or(80);
let authored_target = endpoint
.split_once("://")
.map(|(_, rest)| {
let start = rest.find(['/', '?', '#']).unwrap_or(rest.len());
&rest[start..]
})
.unwrap_or("");
if !authored_target.starts_with('/') {
return Err(invalid(
"empty or absent path: a harness target must declare a request path, never a silent `/` lane"
.to_string(),
));
}
let target = uri
.path_and_query()
.map(|pq| pq.as_str().to_string())
.unwrap_or_default();
Ok(Self { host, port, target })
}
}
async fn perform_exchange(
stream: TcpStream,
target: &ParsedTarget,
msg: OutgoingMessage,
) -> Result<IncomingMessage, TransportError> {
let transport = |detail: String| TransportError::Other { message: detail };
let (mut sender, connection) = http1::handshake(TokioIo::new(stream))
.await
.map_err(|e| transport(format!("http handshake failed: {e}")))?;
tokio::spawn(async move {
let _ = connection.await;
});
sender
.ready()
.await
.map_err(|e| transport(format!("http connection not ready: {e}")))?;
let body = value_to_wire(&msg.body);
let method = Method::from_str(&msg.method)
.map_err(|e| transport(format!("invalid http method `{}`: {e}", msg.method)))?;
let mut builder = Request::builder()
.method(method)
.uri(target.target.clone())
.header("host", format!("{}:{}", target.host, target.port))
.header("connection", "close");
for (name, value) in &msg.headers {
builder = builder.header(name.as_str(), value_to_header(value));
}
let request = builder
.body(Full::new(Bytes::from(body)))
.map_err(|e| transport(format!("http request build failed: {e}")))?;
let response = sender
.send_request(request)
.await
.map_err(|e| transport(format!("http request failed: {e}")))?;
let (parts, response_body) = response.into_parts();
let bytes = response_body
.collect()
.await
.map_err(|e| transport(format!("http response body failed: {e}")))?
.to_bytes();
let content_type = parts
.headers
.get("content-type")
.and_then(|v| v.to_str().ok())
.map(|s| s.to_ascii_lowercase());
Ok(IncomingMessage {
status: Some(parts.status.as_u16()),
headers: wire_headers_to_value(&parts.headers),
body: wire_body_to_value(content_type.as_deref(), &bytes),
method: None,
path: None,
arrival: std::time::Instant::now(),
})
}
fn fail_lane_map_entry(
in_flight: &Mutex<BTreeMap<String, VecDeque<LaneEntry>>>,
key: &str,
generation: u64,
error: TransportError,
) -> bool {
let mut lanes = lock_through(in_flight);
let Some(entry) = lanes
.get_mut(key)
.and_then(|fifo| fifo.iter_mut().find(|entry| entry.generation == generation))
else {
return false;
};
let (tx, rx) = oneshot::channel();
let _ = tx.send(Err(error));
entry.rx = rx;
true
}
fn lane_for(arrivals: &Mutex<BTreeMap<String, Arc<ArrivalLane>>>, path: &str) -> Arc<ArrivalLane> {
let mut lanes = lock_through(arrivals);
lanes
.entry(path.to_string())
.or_insert_with(|| {
let (tx, rx) = mpsc::channel(ARRIVAL_LANE_CAPACITY);
Arc::new(ArrivalLane {
tx,
rx: AsyncMutex::new(rx),
dropped: AtomicUsize::new(0),
})
})
.clone()
}
fn enqueue_arrival(
arrivals: &Mutex<BTreeMap<String, Arc<ArrivalLane>>>,
wire: &HttpWireRequest,
request_headers: &http::HeaderMap,
bytes: &[u8],
) {
let content_type = request_headers
.get("content-type")
.and_then(|v| v.to_str().ok())
.map(|s| s.to_ascii_lowercase());
let arrival = IncomingMessage {
body: wire_body_to_value(content_type.as_deref(), bytes),
headers: wire_headers_to_value(request_headers),
status: None,
method: Some(wire.method.clone()),
path: Some(wire.path.clone()),
arrival: std::time::Instant::now(),
};
let lane = lane_for(arrivals, &wire.path);
if lane.tx.try_send(arrival).is_err() {
lane.dropped.fetch_add(1, Ordering::Relaxed);
tracing::warn!(
path = %wire.path,
capacity = ARRIVAL_LANE_CAPACITY,
"arrival lane full; arrival recorded but not queued for receive"
);
}
}
async fn serve(
state: Arc<ServerState>,
request: Request<Incoming>,
) -> io::Result<Response<Full<Bytes>>> {
let (parts, body) = request.into_parts();
let bytes = body
.collect()
.await
.map(|collected| collected.to_bytes())
.unwrap_or_default();
let wire = HttpWireRequest {
method: parts.method.as_str().to_ascii_uppercase(),
path: parts
.uri
.path_and_query()
.map(|pq| pq.as_str().to_string())
.unwrap_or_else(|| "/".to_string()),
headers: wire_headers_to_string(&parts.headers),
body: bytes.to_vec(),
};
lock_through(&state.requests).push(wire.clone());
enqueue_arrival(&state.arrivals, &wire, &parts.headers, &bytes);
let scripted = {
let mut queue = lock_through(&state.scripted);
let idx = queue.iter().position(|s| s.matches(&wire) && s.times > 0);
idx.map(|idx| {
queue[idx].times -= 1;
let entry = queue[idx].clone();
if queue[idx].times == 0 {
queue.remove(idx);
}
entry
})
};
let Some(scripted) = scripted else {
return Ok(empty_response(
state.fallback_status.unwrap_or(UNMATCHED_STATUS),
));
};
if let Some(delay) = scripted.delay {
tokio::time::sleep(delay).await;
}
if scripted.fault == Some(PartnerFault::Close) {
return Err(io::Error::new(
io::ErrorKind::ConnectionAborted,
"partner fault: close",
));
}
Ok(build_response(
scripted.status,
scripted.headers,
scripted.body,
))
}
fn build_response(
status: u16,
headers: BTreeMap<String, String>,
body: Vec<u8>,
) -> Response<Full<Bytes>> {
let mut builder = Response::builder()
.status(StatusCode::from_u16(status).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR));
for (name, value) in &headers {
if let (Ok(name), Ok(value)) = (
HeaderName::try_from(name.as_str()),
HeaderValue::from_str(value),
) {
builder = builder.header(name, value);
}
}
builder
.body(Full::new(Bytes::from(body)))
.unwrap_or_else(|_| empty_response(UNMATCHED_STATUS))
}
fn empty_response(status: u16) -> Response<Full<Bytes>> {
Response::builder()
.status(StatusCode::from_u16(status).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR))
.body(Full::new(Bytes::new()))
.unwrap_or_else(|_| {
Response::new(Full::new(Bytes::new()))
})
}
pub(crate) fn value_to_wire(body: &Value) -> Vec<u8> {
match body {
Value::Null => Vec::new(),
Value::String(text) => text.clone().into_bytes(),
other => other.to_string().into_bytes(),
}
}
fn value_to_header(value: &Value) -> String {
match value {
Value::String(text) => text.clone(),
other => other.to_string(),
}
}
fn fold_wire_headers<V>(
headers: &http::HeaderMap,
mut render: impl FnMut(String) -> V,
) -> BTreeMap<String, V> {
let mut folded: BTreeMap<String, Vec<String>> = BTreeMap::new();
for (name, value) in headers.iter() {
folded
.entry(name.as_str().to_string())
.or_default()
.push(String::from_utf8_lossy(value.as_bytes()).into_owned());
}
folded
.into_iter()
.map(|(name, values)| (name, render(values.join(", "))))
.collect()
}
fn wire_headers_to_string(headers: &http::HeaderMap) -> BTreeMap<String, String> {
fold_wire_headers(headers, |joined| joined)
}
fn wire_headers_to_value(headers: &http::HeaderMap) -> BTreeMap<String, Value> {
fold_wire_headers(headers, Value::String)
}
fn wire_body_to_value(content_type: Option<&str>, bytes: &[u8]) -> Value {
if content_type.is_some_and(|ct| ct.contains("application/json")) && !bytes.is_empty() {
return serde_json::from_slice(bytes)
.unwrap_or_else(|_| Value::String(String::from_utf8_lossy(bytes).into_owned()));
}
Value::String(String::from_utf8_lossy(bytes).into_owned())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapters::PartnerRouter;
#[test]
fn scripted_response_default_is_ok_once() {
let scripted = ScriptedResponse::default();
assert_eq!(scripted.status, 200);
assert_eq!(scripted.times, 1);
assert_eq!(scripted.method, None);
assert_eq!(scripted.path, None);
assert_eq!(scripted.delay, None);
assert_eq!(scripted.fault, None);
assert!(scripted.headers.is_empty());
assert!(scripted.body.is_empty());
}
#[test]
fn parsed_target_empty_path_is_apparatus_error() {
let error = ParsedTarget::parse("http://host").expect_err("an empty path must fail");
match error {
TransportError::Other { message } => {
assert!(
message.contains("http://host"),
"must name the declaration: {message}"
);
assert!(
message.contains("path"),
"must name the missing path: {message}"
);
}
other => panic!("expected an apparatus-class transport error, got {other:?}"),
}
}
#[test]
fn fail_lane_entry_is_conditional() {
let lane = ClientLane::new();
let (_, rx) = oneshot::channel();
lock_through(&lane.in_flight).insert(
"K".to_string(),
VecDeque::from([LaneEntry { generation: 2, rx }]),
);
let error = || TransportError::Other {
message: "boom".to_string(),
};
assert!(!lane.fail_lane_entry("K", 1, error()));
assert_eq!(
lock_through(&lane.in_flight)
.get("K")
.and_then(|fifo| fifo.front())
.map(|entry| entry.generation),
Some(2)
);
assert!(lane.fail_lane_entry("K", 2, error()));
let mut rx = lane.take("K").expect("the entry stays present");
match rx.try_recv() {
Ok(Err(TransportError::Other { message })) => assert_eq!(message, "boom"),
other => panic!("the parked error must surface on receive, got {other:?}"),
}
}
async fn raw_request(authority: &str, method: &str, target: &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} {target} 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");
}
async fn partner_router(declared: &str) -> (PartnerRouter, String) {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let authority = partner.bound_addr().to_string();
let mut adapters: BTreeMap<String, Box<dyn PartnerAdapter>> = BTreeMap::new();
adapters.insert(declared.to_string(), Box::new(partner));
(PartnerRouter::new(adapters), authority)
}
async fn expired_receive_message(router: &PartnerRouter, declared: &str) -> String {
match router.receive(declared, declared, Duration::ZERO).await {
Err(ReceiveError::Timeout(timeout)) => timeout.to_string(),
other => panic!("expected a receive timeout, got {other:?}"),
}
}
#[tokio::test]
async fn receive_timeout_message_lists_arrived_wire_paths() {
let (router, authority) = partner_router("http://127.0.0.1:0/other").await;
raw_request(&authority, "GET", "/api?x=1").await;
let message = expired_receive_message(&router, "http://127.0.0.1:0/other").await;
assert!(
message.contains("lanes recorded: [/api?x=1]"),
"must list the arrived wire path: {message}"
);
}
#[tokio::test]
async fn arrivals_lane_timeout_redacts_secrets() {
use camel_api::component_metadata::{ComponentMetadata, OptionKind, UriOption};
let (router, authority) = partner_router("http://127.0.0.1:0/elsewhere").await;
let mut metadata = ComponentMetadata::minimal("http");
metadata.uri_options.push(
UriOption::new(
"authPassword",
"partner authentication password",
OptionKind::String,
)
.secret(),
);
let secret_keys: Vec<String> = metadata
.uri_options
.iter()
.filter(|option| option.secret)
.map(|option| option.name.clone())
.collect();
router.set_secret_query_keys(secret_keys);
raw_request(&authority, "GET", "/login?authPassword=hunter2&x=1").await;
let message = expired_receive_message(&router, "http://127.0.0.1:0/elsewhere").await;
assert!(
message.contains("authPassword=***"),
"the secret value must be masked: {message}"
);
assert!(
!message.contains("hunter2"),
"the secret must never print: {message}"
);
assert!(
message.contains("x=1"),
"non-secret pairs must stay visible: {message}"
);
}
#[tokio::test]
async fn encoded_secret_query_key_redacts() {
let (router, authority) = partner_router("http://127.0.0.1:0/elsewhere").await;
router.set_secret_query_keys(vec!["authPassword".to_string()]);
raw_request(&authority, "GET", "/login?%61uthPassword=hunter2&x=1").await;
let message = expired_receive_message(&router, "http://127.0.0.1:0/elsewhere").await;
assert!(
message.contains("%61uthPassword=***&x=1"),
"the raw key span stays, the value masks: {message}"
);
assert!(
!message.contains("hunter2"),
"the secret must never print: {message}"
);
assert!(
message.contains("x=1"),
"non-secret pairs stay visible: {message}"
);
}
#[tokio::test]
async fn await_parked_timeout_redacts_secrets() {
let scripted = ScriptedResponse {
path: Some("/login?authPassword=hunter2&x=1".to_string()),
delay: Some(Duration::from_secs(30)),
..ScriptedResponse::default()
};
let partner = HttpPartner::start(vec![scripted])
.await
.expect("partner must bind 127.0.0.1:0");
let authority = partner.bound_addr().to_string();
let declared = format!("http://{authority}/login?authPassword=hunter2&x=1");
let mut adapters: BTreeMap<String, Box<dyn PartnerAdapter>> = BTreeMap::new();
adapters.insert(declared.clone(), Box::new(partner));
let router = PartnerRouter::new(adapters);
router.set_secret_query_keys(vec!["authPassword".to_string()]);
router
.send(
&declared,
&declared,
OutgoingMessage {
body: Value::Null,
headers: BTreeMap::new(),
method: "GET".to_string(),
},
)
.await
.expect("the send must dial the partner");
for _ in 0..2000 {
if !router.recorded_requests(&declared).is_empty() {
break;
}
tokio::time::sleep(Duration::from_millis(5)).await;
}
let message = expired_receive_message(&router, &declared).await;
assert!(
message.contains("authPassword=***"),
"the secret value must be masked: {message}"
);
assert!(
!message.contains("hunter2"),
"the secret must never print: {message}"
);
assert!(
message.contains("x=1"),
"non-secret pairs must stay visible: {message}"
);
}
#[tokio::test]
async fn query_bearing_receive_matches_end_to_end() {
let partner = HttpPartner::start_permissive(200)
.await
.expect("partner must bind 127.0.0.1:0");
let authority = partner.bound_addr().to_string();
let declared = format!("http://{authority}/api?flag=a&x=1");
let mut adapters: BTreeMap<String, Box<dyn PartnerAdapter>> = BTreeMap::new();
adapters.insert(declared.clone(), Box::new(partner));
let router = PartnerRouter::new(adapters);
router
.send(
&declared,
&declared,
OutgoingMessage {
body: Value::Null,
headers: BTreeMap::new(),
method: "GET".to_string(),
},
)
.await
.expect("the send must dial the declared endpoint");
let message = router
.receive(&declared, &declared, Duration::from_secs(5))
.await
.expect("the roundtrip must complete");
assert_eq!(
message.status,
Some(200),
"the permissive partner serves 200"
);
let recorded = router.recorded_requests(&declared);
assert_eq!(recorded.len(), 1, "exactly one request crossed the wire");
assert_eq!(
recorded[0].path, "/api?flag=a&x=1",
"the arrival key equals the declared wire path_and_query"
);
assert_eq!(recorded[0].method, "GET");
}
}