use std::sync::{
Arc,
atomic::{AtomicU32, AtomicUsize, Ordering},
};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use apimock_routing::util::http::percent_decode_url_path;
use serde::Serialize;
use tokio::io::AsyncWriteExt;
use tokio::sync::broadcast;
pub const TRACE_CHANNEL_CAPACITY: usize = 1_024;
pub const MAX_SUBSCRIBERS: usize = 4;
#[derive(Clone, Debug, Serialize)]
pub struct MatchTraceEvent {
pub event_id: u64,
pub schema_version: u8,
pub received_at_ms: u64,
pub duration_ms: u32,
pub request: RequestSummary,
pub outcome: Outcome,
pub dropped_count: u32,
}
#[derive(Clone, Debug, Serialize)]
#[non_exhaustive]
pub struct RequestSummary {
pub method: String,
pub url_path: String,
pub headers: Vec<(String, String)>,
#[serde(skip_serializing_if = "Option::is_none")]
pub body_json: Option<serde_json::Value>,
#[serde(skip_serializing_if = "std::ops::Not::not")]
pub body_truncated: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub body_len: Option<usize>,
}
impl RequestSummary {
pub fn new(
method: String,
url_path: String,
headers: Vec<(String, String)>,
body_len: Option<usize>,
config: &TraceConfig,
) -> Self {
Self {
method,
url_path,
headers: config.redact_headers(headers),
body_json: None,
body_truncated: false,
body_len,
}
}
}
pub const REDACTED_HEADER_VALUE: &str = "[redacted]";
pub const DEFAULT_HEADER_DENYLIST: &[&str] = &[
"authorization",
"cookie",
"set-cookie",
"proxy-authorization",
"x-api-key",
"token",
"access_token",
"refresh_token",
"password",
"secret",
"client_secret",
"api_key",
];
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum HeaderRedactionMode {
Denylist,
Allowlist,
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct TraceConfig {
pub capture_body: bool,
pub max_body_bytes: usize,
pub header_redaction: HeaderRedactionMode,
pub header_denylist: Vec<String>,
pub header_allowlist: Vec<String>,
}
impl Default for TraceConfig {
fn default() -> Self {
Self {
capture_body: false,
max_body_bytes: 8_192,
header_redaction: HeaderRedactionMode::Denylist,
header_denylist: DEFAULT_HEADER_DENYLIST
.iter()
.map(|s| s.to_string())
.collect(),
header_allowlist: Vec::new(),
}
}
}
impl TraceConfig {
pub(crate) fn is_redacted_key(&self, name: &str) -> bool {
match self.header_redaction {
HeaderRedactionMode::Denylist => self
.header_denylist
.iter()
.any(|denied| denied.eq_ignore_ascii_case(name)),
HeaderRedactionMode::Allowlist => !self
.header_allowlist
.iter()
.any(|allowed| allowed.eq_ignore_ascii_case(name)),
}
}
fn redact_headers(&self, headers: Vec<(String, String)>) -> Vec<(String, String)> {
headers
.into_iter()
.map(|(name, value)| {
if self.is_redacted_key(&name) {
(name, REDACTED_HEADER_VALUE.to_string())
} else {
(name, value)
}
})
.collect()
}
pub(crate) fn redact_query_string(&self, query: &str) -> String {
query
.split('&')
.map(|pair| match pair.split_once('=') {
Some((key, _value)) if self.is_redacted_key(&percent_decode_url_path(key)) => {
format!("{key}={REDACTED_HEADER_VALUE}")
}
_ => pair.to_string(),
})
.collect::<Vec<String>>()
.join("&")
}
pub(crate) fn redact_json_value(&self, value: &serde_json::Value) -> serde_json::Value {
match value {
serde_json::Value::Object(map) => serde_json::Value::Object(
map.iter()
.map(|(key, val)| {
let redacted_val = if self.is_redacted_key(key) {
serde_json::Value::String(REDACTED_HEADER_VALUE.to_string())
} else {
self.redact_json_value(val)
};
(key.clone(), redacted_val)
})
.collect(),
),
serde_json::Value::Array(items) => {
serde_json::Value::Array(items.iter().map(|v| self.redact_json_value(v)).collect())
}
scalar => scalar.clone(),
}
}
}
#[derive(Clone, Debug, Serialize)]
#[serde(tag = "type", rename_all = "snake_case")]
#[non_exhaustive]
pub enum Outcome {
Matched {
rule_set_index: usize,
rule_index: usize,
},
Middleware {
file_path: String,
status: u16,
},
Fallback {
file_path: String,
status: u16,
},
Miss {
status: u16,
},
Error {
kind: String,
message: String,
},
}
#[derive(Clone)]
pub struct TraceEmitter {
sender: broadcast::Sender<MatchTraceEvent>,
event_counter: Arc<AtomicU32>,
dropped_counter: Arc<AtomicU32>,
pub config: Arc<TraceConfig>,
}
impl TraceEmitter {
pub fn new() -> Self {
Self::with_config(TraceConfig::default())
}
pub fn with_config(config: TraceConfig) -> Self {
let (sender, _) = broadcast::channel(TRACE_CHANNEL_CAPACITY);
Self {
sender,
event_counter: Arc::new(AtomicU32::new(0)),
dropped_counter: Arc::new(AtomicU32::new(0)),
config: Arc::new(config),
}
}
pub fn subscribe(&self) -> broadcast::Receiver<MatchTraceEvent> {
self.sender.subscribe()
}
pub fn enrich_with_body(
&self,
summary: &mut RequestSummary,
body_json: Option<&serde_json::Value>,
) {
if !self.config.capture_body {
return;
}
match body_json {
None => {} Some(v) => {
let redacted = self.config.redact_json_value(v);
match serde_json::to_string(&redacted) {
Ok(s) if s.len() <= self.config.max_body_bytes => {
summary.body_json = Some(redacted);
}
Ok(_) => {
summary.body_truncated = true;
}
Err(_) => {} }
}
}
}
pub fn emit(
&self,
received_at_ms: u64,
duration_ms: u32,
request: RequestSummary,
outcome: Outcome,
) {
let event_id = self.event_counter.fetch_add(1, Ordering::Relaxed) as u64;
let dropped_count = self.dropped_counter.swap(0, Ordering::Relaxed);
let event = MatchTraceEvent {
event_id,
schema_version: 1,
received_at_ms,
duration_ms,
request,
outcome,
dropped_count,
};
if self.sender.send(event).is_err() {
self.dropped_counter.fetch_add(1, Ordering::Relaxed);
}
}
pub fn has_subscribers(&self) -> bool {
self.sender.receiver_count() > 0
}
}
impl Default for TraceEmitter {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone, Debug, Default)]
pub enum TraceTransportConfig {
#[cfg(unix)]
Uds { path: String },
Tcp { addr: String },
#[default]
Disabled,
}
pub struct TraceTransport;
impl TraceTransport {
pub async fn accept_loop(config: TraceTransportConfig, emitter: TraceEmitter) {
match config {
#[cfg(unix)]
TraceTransportConfig::Uds { path } => Self::uds_accept_loop(path, emitter).await,
TraceTransportConfig::Tcp { addr } => Self::tcp_accept_loop(addr, emitter).await,
TraceTransportConfig::Disabled => {
}
}
}
async fn tcp_accept_loop(addr: String, emitter: TraceEmitter) {
let listener = match tokio::net::TcpListener::bind(&addr).await {
Ok(l) => {
let bound = l
.local_addr()
.map(|a| a.to_string())
.unwrap_or_else(|_| addr.clone());
log::info!("trace transport: TCP listening on {}", bound);
if !l.local_addr().map(|a| a.ip().is_loopback()).unwrap_or(true) {
log::warn!(
"trace transport: TCP listening on a non-loopback address ({}) — \
this transport has no authentication; anything that can reach it \
receives the live request trace feed",
bound
);
}
l
}
Err(e) => {
log::error!("trace transport: failed to bind TCP {}: {}", addr, e);
return;
}
};
let active = Arc::new(AtomicUsize::new(0));
loop {
match listener.accept().await {
Ok((stream, peer)) => {
let count = active.fetch_add(1, Ordering::Relaxed) + 1;
if count > MAX_SUBSCRIBERS {
active.fetch_sub(1, Ordering::Relaxed);
tokio::spawn(async move {
let (_, mut writer) = tokio::io::split(stream);
let _ = writer
.write_all(b"{\"error\":\"max_subscribers_reached\"}\n")
.await;
});
continue;
}
log::debug!("trace: TCP subscriber connected from {}", peer);
let rx = emitter.subscribe();
let active_clone = active.clone();
tokio::spawn(async move {
let (_, writer) = tokio::io::split(stream);
Self::forward_events(writer, rx).await;
active_clone.fetch_sub(1, Ordering::Relaxed);
log::debug!("trace: TCP subscriber {} disconnected", peer);
});
}
Err(e) => {
log::error!("trace: TCP accept error: {}", e);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}
}
#[cfg(unix)]
async fn uds_accept_loop(path: String, emitter: TraceEmitter) {
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::remove_file(&path);
let listener = match tokio::net::UnixListener::bind(&path) {
Ok(l) => {
if let Err(e) =
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))
{
log::error!(
"trace transport: failed to restrict UDS permissions on {}: {}",
path,
e
);
}
log::info!("trace transport: UDS listening at {}", path);
l
}
Err(e) => {
log::error!("trace transport: failed to bind UDS {}: {}", path, e);
return;
}
};
let active = Arc::new(AtomicUsize::new(0));
loop {
match listener.accept().await {
Ok((stream, _)) => {
let count = active.fetch_add(1, Ordering::Relaxed) + 1;
if count > MAX_SUBSCRIBERS {
active.fetch_sub(1, Ordering::Relaxed);
tokio::spawn(async move {
let (_, mut writer) = tokio::io::split(stream);
let _ = writer
.write_all(b"{\"error\":\"max_subscribers_reached\"}\n")
.await;
});
continue;
}
log::debug!("trace: UDS subscriber connected");
let rx = emitter.subscribe();
let active_clone = active.clone();
tokio::spawn(async move {
let (_, writer) = tokio::io::split(stream);
Self::forward_events(writer, rx).await;
active_clone.fetch_sub(1, Ordering::Relaxed);
log::debug!("trace: UDS subscriber disconnected");
});
}
Err(e) => {
log::error!("trace: UDS accept error: {}", e);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}
}
async fn forward_events<W>(mut writer: W, mut rx: broadcast::Receiver<MatchTraceEvent>)
where
W: tokio::io::AsyncWrite + Unpin,
{
let mut lagged_events: u32 = 0;
loop {
let mut event = match rx.recv().await {
Ok(e) => e,
Err(broadcast::error::RecvError::Lagged(n)) => {
lagged_events =
lagged_events.saturating_add(u32::try_from(n).unwrap_or(u32::MAX));
log::debug!("trace: subscriber lagged, {} events dropped", n);
continue;
}
Err(broadcast::error::RecvError::Closed) => break,
};
event.dropped_count = event.dropped_count.saturating_add(lagged_events);
lagged_events = 0;
let mut line = match serde_json::to_string(&event) {
Ok(s) => s,
Err(e) => {
log::error!("trace: serialise error: {}", e);
continue;
}
};
line.push('\n');
if writer.write_all(line.as_bytes()).await.is_err() {
break; }
}
}
}
pub fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or(Duration::ZERO)
.as_millis() as u64
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn emit_received_by_subscriber() {
let emitter = TraceEmitter::new();
let mut rx = emitter.subscribe();
emitter.emit(
1_000_000,
5,
RequestSummary {
method: "GET".into(),
url_path: "/api/test".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
},
Outcome::Miss { status: 404 },
);
let event = rx.try_recv().expect("event in channel");
assert_eq!(event.event_id, 0);
assert_eq!(event.schema_version, 1);
assert_eq!(event.request.method, "GET");
assert_eq!(event.duration_ms, 5);
assert_eq!(event.dropped_count, 0);
assert!(matches!(event.outcome, Outcome::Miss { status: 404 }));
}
#[tokio::test]
async fn emit_no_subscriber_increments_dropped() {
let emitter = TraceEmitter::new();
emitter.emit(
0,
0,
RequestSummary {
method: "GET".into(),
url_path: "/".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
},
Outcome::Miss { status: 404 },
);
let mut rx = emitter.subscribe();
emitter.emit(
0,
0,
RequestSummary {
method: "GET".into(),
url_path: "/".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
},
Outcome::Miss { status: 200 },
);
let event = rx.try_recv().expect("second event visible");
assert_eq!(
event.dropped_count, 1,
"first event should be counted dropped"
);
}
#[test]
fn has_subscribers_reflects_state() {
let emitter = TraceEmitter::new();
assert!(!emitter.has_subscribers());
let _rx = emitter.subscribe();
assert!(emitter.has_subscribers());
}
#[tokio::test]
async fn outcome_serialises_correctly() {
let event = MatchTraceEvent {
event_id: 7,
schema_version: 1,
received_at_ms: 0,
duration_ms: 0,
request: RequestSummary {
method: "POST".into(),
url_path: "/x".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
},
outcome: Outcome::Matched {
rule_set_index: 0,
rule_index: 2,
},
dropped_count: 0,
};
let json = serde_json::to_string(&event).unwrap();
assert!(json.contains("\"type\":\"matched\""));
assert!(json.contains("\"rule_index\":2"));
assert!(json.contains("\"schema_version\":1"));
}
#[tokio::test]
async fn tcp_transport_delivers_events() {
let emitter = TraceEmitter::new();
let emitter_clone = emitter.clone();
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let bound_addr = listener.local_addr().unwrap();
tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let rx = emitter_clone.subscribe();
let (_, writer) = tokio::io::split(stream);
TraceTransport::forward_events(writer, rx).await;
});
let mut client = tokio::net::TcpStream::connect(bound_addr).await.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
emitter.emit(
42,
3,
RequestSummary {
method: "GET".into(),
url_path: "/ping".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
},
Outcome::Miss { status: 404 },
);
use tokio::io::AsyncBufReadExt;
let mut reader = tokio::io::BufReader::new(&mut client);
let mut line = String::new();
tokio::time::timeout(
std::time::Duration::from_secs(2),
reader.read_line(&mut line),
)
.await
.expect("timeout")
.expect("read ok");
let value: serde_json::Value = serde_json::from_str(line.trim()).expect("valid JSON");
assert_eq!(value["request"]["url_path"], "/ping");
assert_eq!(value["outcome"]["type"], "miss");
assert_eq!(value["schema_version"], 1);
}
fn dummy_summary() -> RequestSummary {
RequestSummary {
method: "GET".into(),
url_path: "/".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
}
}
#[tokio::test]
async fn a_lagging_subscriber_reports_dropped_count_on_its_next_event() {
let emitter = TraceEmitter::new();
let rx = emitter.subscribe();
for _ in 0..(TRACE_CHANNEL_CAPACITY + 10) {
emitter.emit(0, 0, dummy_summary(), Outcome::Miss { status: 404 });
}
drop(emitter);
let mut buf: Vec<u8> = Vec::new();
TraceTransport::forward_events(&mut buf, rx).await;
let text = String::from_utf8(buf).expect("valid utf8");
let first_line = text.lines().next().expect("at least one forwarded event");
let event: serde_json::Value = serde_json::from_str(first_line).expect("valid JSON");
assert!(
event["dropped_count"].as_u64().unwrap_or(0) > 0,
"the first event surviving a lag must report it: {first_line}"
);
}
#[test]
fn enrich_with_body_disabled_by_default() {
let emitter = TraceEmitter::new(); let mut summary = RequestSummary {
method: "POST".into(),
url_path: "/".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
};
let body = serde_json::json!({"action": "create"});
emitter.enrich_with_body(&mut summary, Some(&body));
assert!(
summary.body_json.is_none(),
"body should not be captured when disabled"
);
assert!(!summary.body_truncated);
}
#[test]
fn enrich_with_body_enabled_captures_small_body() {
let emitter = TraceEmitter::with_config(TraceConfig {
capture_body: true,
max_body_bytes: 8_192,
..Default::default()
});
let mut summary = RequestSummary {
method: "POST".into(),
url_path: "/".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
};
let body = serde_json::json!({"action": "create", "user_id": 42});
emitter.enrich_with_body(&mut summary, Some(&body));
assert!(
summary.body_json.is_some(),
"body should be captured when enabled"
);
assert_eq!(summary.body_json.unwrap()["action"], "create");
assert!(!summary.body_truncated);
}
#[test]
fn enrich_with_body_truncates_oversized_body() {
let emitter = TraceEmitter::with_config(TraceConfig {
capture_body: true,
max_body_bytes: 10,
..Default::default()
});
let mut summary = RequestSummary {
method: "POST".into(),
url_path: "/".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
};
let body = serde_json::json!({"data": "this is longer than 10 bytes"});
emitter.enrich_with_body(&mut summary, Some(&body));
assert!(
summary.body_json.is_none(),
"oversized body should be omitted"
);
assert!(summary.body_truncated, "body_truncated flag should be set");
}
#[test]
fn request_summary_body_json_not_in_serialised_output_when_none() {
let summary = RequestSummary {
method: "GET".into(),
url_path: "/api".into(),
headers: vec![],
body_json: None,
body_truncated: false,
body_len: None,
};
let json = serde_json::to_string(&summary).unwrap();
assert!(
!json.contains("body_json"),
"absent body_json must be skipped"
);
assert!(
!json.contains("body_truncated"),
"false body_truncated must be skipped"
);
}
fn headers_with_credentials() -> Vec<(String, String)> {
vec![
("authorization".into(), "Bearer secret-token".into()),
("cookie".into(), "session=abc123".into()),
("x-api-key".into(), "sk-live-very-secret".into()),
("content-type".into(), "application/json".into()),
]
}
#[test]
fn default_config_redacts_credential_headers_in_serialised_output() {
let config = TraceConfig::default();
let summary = RequestSummary::new(
"POST".into(),
"/login".into(),
headers_with_credentials(),
None,
&config,
);
let json = serde_json::to_string(&summary).unwrap();
assert!(!json.contains("Bearer secret-token"), "json was: {json}");
assert!(!json.contains("session=abc123"), "json was: {json}");
assert!(!json.contains("sk-live-very-secret"), "json was: {json}");
assert!(
json.contains("application/json"),
"a non-credential header must survive: {json}"
);
}
#[test]
fn redacted_headers_are_present_and_marked_not_absent() {
let config = TraceConfig::default();
let summary = RequestSummary::new(
"POST".into(),
"/login".into(),
headers_with_credentials(),
None,
&config,
);
assert_eq!(summary.headers.len(), 4, "no header should be dropped");
let authorization = summary
.headers
.iter()
.find(|(name, _)| name == "authorization")
.expect("authorization header must still be present");
assert_eq!(authorization.1, REDACTED_HEADER_VALUE);
let json = serde_json::to_string(&summary).unwrap();
assert!(
json.contains("\"authorization\""),
"redacted header name must still appear: {json}"
);
assert!(json.contains(REDACTED_HEADER_VALUE), "json was: {json}");
}
#[test]
fn denylist_matches_case_insensitively() {
let config = TraceConfig::default();
let headers = vec![
("Authorization".into(), "Bearer secret-token".into()),
("COOKIE".into(), "session=abc123".into()),
];
let summary = RequestSummary::new("GET".into(), "/".into(), headers, None, &config);
let json = serde_json::to_string(&summary).unwrap();
assert!(!json.contains("Bearer secret-token"), "json was: {json}");
assert!(!json.contains("session=abc123"), "json was: {json}");
assert!(json.contains(REDACTED_HEADER_VALUE), "json was: {json}");
}
#[test]
fn allowlist_mode_redacts_everything_not_listed() {
let config = TraceConfig {
header_redaction: HeaderRedactionMode::Allowlist,
header_allowlist: vec!["content-type".into()],
..Default::default()
};
let headers = vec![
("content-type".into(), "application/json".into()),
("authorization".into(), "Bearer secret-token".into()),
("x-request-id".into(), "not-a-credential".into()),
];
let summary = RequestSummary::new("GET".into(), "/".into(), headers, None, &config);
let by_name = |name: &str| {
summary
.headers
.iter()
.find(|(n, _)| n == name)
.map(|(_, v)| v.as_str())
};
assert_eq!(by_name("content-type"), Some("application/json"));
assert_eq!(by_name("authorization"), Some(REDACTED_HEADER_VALUE));
assert_eq!(
by_name("x-request-id"),
Some(REDACTED_HEADER_VALUE),
"an unlisted, non-credential header must still be redacted in allowlist mode"
);
}
#[test]
fn allowlist_mode_with_no_entries_redacts_everything() {
let config = TraceConfig {
header_redaction: HeaderRedactionMode::Allowlist,
..Default::default()
};
let summary = RequestSummary::new(
"GET".into(),
"/".into(),
vec![("content-type".into(), "application/json".into())],
None,
&config,
);
assert_eq!(summary.headers[0].1, REDACTED_HEADER_VALUE);
}
#[test]
fn three_body_states_are_distinguishable_in_the_serialised_form() {
let config = TraceConfig::default();
let no_body = RequestSummary::new("GET".into(), "/".into(), vec![], None, &config);
let no_body_json = serde_json::to_string(&no_body).unwrap();
assert!(!no_body_json.contains("body_json"), "{no_body_json}");
assert!(!no_body_json.contains("body_len"), "{no_body_json}");
let mut json_captured =
RequestSummary::new("POST".into(), "/".into(), vec![], Some(11), &config);
let emitter = TraceEmitter::with_config(TraceConfig {
capture_body: true,
..Default::default()
});
emitter.enrich_with_body(&mut json_captured, Some(&serde_json::json!({"a": 1})));
let json_captured_str = serde_json::to_string(&json_captured).unwrap();
assert!(
json_captured_str.contains("\"body_json\""),
"{json_captured_str}"
);
assert!(
json_captured_str.contains("\"body_len\":11"),
"a JSON-captured body must still report its length: {json_captured_str}"
);
let body_present_not_captured =
RequestSummary::new("POST".into(), "/".into(), vec![], Some(27), &config);
let not_captured_str = serde_json::to_string(&body_present_not_captured).unwrap();
assert!(
!not_captured_str.contains("body_json"),
"{not_captured_str}"
);
assert!(
not_captured_str.contains("\"body_len\":27"),
"{not_captured_str}"
);
}
#[test]
fn non_json_body_reports_length_but_never_content() {
let config = TraceConfig::default();
let summary = RequestSummary::new("POST".into(), "/".into(), vec![], Some(32), &config);
let json = serde_json::to_string(&summary).unwrap();
assert!(json.contains("\"body_len\":32"), "json was: {json}");
assert!(
!json.contains("username") && !json.contains("hunter2"),
"no fragment of a body — captured or not — should appear: {json}"
);
}
#[test]
fn a_query_string_token_is_redacted_by_default() {
let config = TraceConfig::default();
let redacted = config.redact_query_string("token=secret&page=2");
assert_eq!(redacted, "token=[redacted]&page=2");
}
#[test]
fn a_non_denied_query_parameter_survives() {
let config = TraceConfig::default();
let redacted = config.redact_query_string("page=2&access_token=abc123&sort=asc");
assert_eq!(redacted, "page=2&access_token=[redacted]&sort=asc");
}
#[test]
fn a_bare_flag_parameter_is_left_alone() {
let config = TraceConfig::default();
let redacted = config.redact_query_string("verbose&token=secret");
assert_eq!(redacted, "verbose&token=[redacted]");
}
#[test]
fn a_query_string_key_is_matched_case_insensitively() {
let config = TraceConfig::default();
let redacted = config.redact_query_string("TOKEN=secret");
assert_eq!(redacted, "TOKEN=[redacted]");
}
#[test]
fn a_percent_encoded_query_key_does_not_bypass_redaction() {
let config = TraceConfig::default();
let redacted = config.redact_query_string("%74oken=secret");
assert_eq!(redacted, "%74oken=[redacted]");
}
#[test]
fn a_top_level_body_secret_is_redacted_by_default() {
let config = TraceConfig::default();
let body = serde_json::json!({"username": "alice", "password": "hunter2"});
let redacted = config.redact_json_value(&body);
assert_eq!(redacted["username"], "alice");
assert_eq!(redacted["password"], REDACTED_HEADER_VALUE);
}
#[test]
fn a_nested_body_secret_is_redacted_too() {
let config = TraceConfig::default();
let body = serde_json::json!({
"user": {"name": "alice", "api_key": "sk-live-very-secret"},
"items": [{"id": 1}, {"token": "should-not-appear"}],
});
let redacted = config.redact_json_value(&body);
assert_eq!(redacted["user"]["name"], "alice");
assert_eq!(redacted["user"]["api_key"], REDACTED_HEADER_VALUE);
assert_eq!(redacted["items"][0]["id"], 1);
assert_eq!(redacted["items"][1]["token"], REDACTED_HEADER_VALUE);
let json = serde_json::to_string(&redacted).unwrap();
assert!(
!json.contains("sk-live-very-secret") && !json.contains("should-not-appear"),
"no redacted value should survive serialisation: {json}"
);
}
#[test]
fn enrich_with_body_redacts_a_captured_body() {
let emitter = TraceEmitter::with_config(TraceConfig {
capture_body: true,
..Default::default()
});
let mut summary = dummy_summary();
let body = serde_json::json!({"action": "login", "password": "hunter2"});
emitter.enrich_with_body(&mut summary, Some(&body));
let captured = summary.body_json.expect("body should be captured");
assert_eq!(captured["action"], "login");
assert_eq!(captured["password"], REDACTED_HEADER_VALUE);
}
#[test]
fn outcome_middleware_serialises_with_file_path_and_status() {
let outcome = Outcome::Middleware {
file_path: "middleware/auth.rhai".into(),
status: 200,
};
let json = serde_json::to_string(&outcome).unwrap();
assert!(json.contains("\"type\":\"middleware\""), "json was: {json}");
assert!(
json.contains("\"file_path\":\"middleware/auth.rhai\""),
"json was: {json}"
);
assert!(json.contains("\"status\":200"), "json was: {json}");
}
#[cfg(unix)]
#[tokio::test]
async fn uds_socket_is_created_with_owner_only_permissions() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("trace.sock").to_str().unwrap().to_owned();
let emitter = TraceEmitter::new();
let accept_loop = tokio::spawn(TraceTransport::accept_loop(
TraceTransportConfig::Uds { path: path.clone() },
emitter,
));
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while !std::path::Path::new(&path).exists() {
assert!(
std::time::Instant::now() < deadline,
"socket never appeared"
);
tokio::time::sleep(Duration::from_millis(5)).await;
}
let mode = std::fs::metadata(&path)
.expect("stat socket")
.permissions()
.mode();
assert_eq!(
mode & 0o777,
0o600,
"socket permissions should be owner-only, got {mode:o}"
);
accept_loop.abort();
}
}