use crate::{
agent::cancellation::AgentCancellation, output::redact_sensitive_text,
providers::error::ProviderError,
};
use serde_json::Value;
use std::{
collections::BTreeMap,
io::Read,
sync::{
OnceLock,
atomic::{AtomicUsize, Ordering},
mpsc,
},
thread,
time::{Duration, Instant},
};
const PROVIDER_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
const PROVIDER_READ_IDLE_TIMEOUT: Duration = Duration::from_secs(120);
const PROVIDER_STREAM_WALL_TIMEOUT: Duration = Duration::from_secs(300);
pub(crate) const PROVIDER_ERROR_BODY_MAX_BYTES: usize = 16 * 1024;
const PROVIDER_ERROR_BODY_REDACTION_SLACK_BYTES: usize = 4 * 1024;
const PROVIDER_ERROR_BODY_IDLE_TIMEOUT: Duration = Duration::from_secs(5);
const PROVIDER_WORKER_SHUTDOWN_TIMEOUT: Duration = Duration::from_millis(100);
const PROVIDER_WORKER_SHUTDOWN_POLL_INTERVAL: Duration = Duration::from_millis(5);
const PROVIDER_HEADER_WAIT_WORKER_LIMIT: usize = 32;
const PROVIDER_STREAM_READER_WORKER_LIMIT: usize = 32;
const PROVIDER_ERROR_BODY_READER_WORKER_LIMIT: usize = 32;
static PROVIDER_ERROR_BODY_READER_WORKERS: AtomicUsize = AtomicUsize::new(0);
static PROVIDER_STREAM_READER_WORKERS: AtomicUsize = AtomicUsize::new(0);
static PROVIDER_HEADER_WAIT_WORKERS: AtomicUsize = AtomicUsize::new(0);
struct ProviderWorkerSlot<'a> {
workers: &'a AtomicUsize,
}
impl Drop for ProviderWorkerSlot<'_> {
fn drop(&mut self) {
self.workers.fetch_sub(1, Ordering::SeqCst);
}
}
fn reserve_provider_worker_slot(
workers: &AtomicUsize,
limit: usize,
) -> Option<ProviderWorkerSlot<'_>> {
let mut current = workers.load(Ordering::SeqCst);
loop {
if current >= limit {
return None;
}
match workers.compare_exchange(current, current + 1, Ordering::SeqCst, Ordering::SeqCst) {
Ok(_) => return Some(ProviderWorkerSlot { workers }),
Err(next) => current = next,
}
}
}
fn reserve_provider_error_body_reader_worker_slot() -> Option<ProviderWorkerSlot<'static>> {
reserve_provider_worker_slot(
&PROVIDER_ERROR_BODY_READER_WORKERS,
PROVIDER_ERROR_BODY_READER_WORKER_LIMIT,
)
}
fn reserve_provider_stream_reader_worker_slot() -> Option<ProviderWorkerSlot<'static>> {
reserve_provider_worker_slot(
&PROVIDER_STREAM_READER_WORKERS,
PROVIDER_STREAM_READER_WORKER_LIMIT,
)
}
fn reserve_provider_header_wait_worker_slot() -> Option<ProviderWorkerSlot<'static>> {
reserve_provider_worker_slot(
&PROVIDER_HEADER_WAIT_WORKERS,
PROVIDER_HEADER_WAIT_WORKER_LIMIT,
)
}
#[derive(Clone, PartialEq, Eq)]
pub struct HttpRequest {
pub method: String,
pub url: String,
pub headers: BTreeMap<String, String>,
pub body: Value,
}
impl std::fmt::Debug for HttpRequest {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let headers = self
.headers
.iter()
.map(|(name, value)| {
(
name,
if is_sensitive_header_name(name) {
"<redacted>"
} else {
value.as_str()
},
)
})
.collect::<BTreeMap<_, _>>();
formatter
.debug_struct("HttpRequest")
.field("method", &self.method)
.field("url", &self.url)
.field("headers", &headers)
.field("body", &self.body)
.finish()
}
}
fn is_sensitive_header_name(name: &str) -> bool {
let name = name.to_ascii_lowercase();
matches!(
name.as_str(),
"authorization" | "x-api-key" | "api-key" | "anthropic-api-key"
) || name.contains("token")
|| name.contains("key")
}
pub trait HttpTransport {
#[allow(dead_code)]
fn stream_json(
&self,
request: HttpRequest,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()>;
fn stream_json_cancellable(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()>;
fn stream_json_cancellable_with_semantic_deadline(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
semantic_deadline: &std::sync::atomic::AtomicU64,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()>;
fn stream_json_cancellable_with_response_metadata(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
semantic_deadline: &std::sync::atomic::AtomicU64,
on_response_id: &mut dyn FnMut(Option<String>),
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let _ = on_response_id;
self.stream_json_cancellable_with_semantic_deadline(
request,
cancellation,
semantic_deadline,
on_chunk,
)
}
}
#[derive(Debug, Clone, Default)]
pub struct ReqwestHttpTransport;
impl HttpTransport for ReqwestHttpTransport {
fn stream_json(
&self,
request: HttpRequest,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable(request, &AgentCancellation::default(), on_chunk)
}
fn stream_json_cancellable(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable_with_semantic_deadline(
request,
cancellation,
&std::sync::atomic::AtomicU64::new(0),
on_chunk,
)
}
fn stream_json_cancellable_with_semantic_deadline(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
semantic_deadline: &std::sync::atomic::AtomicU64,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable_with_response_metadata(
request,
cancellation,
semantic_deadline,
&mut |_| {},
on_chunk,
)
}
fn stream_json_cancellable_with_response_metadata(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
semantic_deadline: &std::sync::atomic::AtomicU64,
on_response_id: &mut dyn FnMut(Option<String>),
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let HttpRequest {
method,
url,
headers,
body,
} = request;
let method = reqwest::Method::from_bytes(method.as_bytes()).map_err(|error| {
ProviderError::transport(format!("invalid provider HTTP method '{method}': {error}"))
})?;
let client = provider_client()?;
let mut builder = client.request(method, &url);
for (name, value) in &headers {
builder = builder.header(name, value);
}
let response = send_provider_request_with_header_timeout(
builder,
body,
&url,
PROVIDER_READ_IDLE_TIMEOUT,
cancellation,
)?;
let request_id = response
.headers()
.get("x-request-id")
.and_then(|value| value.to_str().ok())
.map(crate::providers::error::bounded_response_identity_string);
on_response_id(request_id);
if !response.status().is_success() {
let status = response.status();
let body = read_bounded_redacted_response_body(response);
let endpoint = sanitize_provider_error_url(&url);
return Err(ProviderError::http_status(
status.as_u16(),
format!("provider request failed for {endpoint} with status {status}: {body}"),
)
.into());
}
stream_reader_with_idle_timeout(
response,
PROVIDER_READ_IDLE_TIMEOUT,
cancellation,
Some(semantic_deadline),
on_chunk,
)
}
}
fn send_provider_request_with_header_timeout(
builder: reqwest::blocking::RequestBuilder,
body: Value,
url: &str,
idle_timeout: Duration,
cancellation: &AgentCancellation,
) -> anyhow::Result<reqwest::blocking::Response> {
let endpoint = sanitize_provider_error_url(url);
run_provider_header_wait_with_idle_timeout(
move || {
builder.json(&body).send().map_err(|error| {
let error = error.without_url();
ProviderError::transport(format!(
"provider request send failed for {endpoint}: {error}"
))
.into()
})
},
idle_timeout,
cancellation,
)
}
pub(crate) fn run_provider_header_wait_with_idle_timeout<T, F>(
operation: F,
idle_timeout: Duration,
cancellation: &AgentCancellation,
) -> anyhow::Result<T>
where
T: Send + 'static,
F: FnOnce() -> anyhow::Result<T> + Send + 'static,
{
let (sender, receiver) = mpsc::sync_channel(1);
let worker_slot = reserve_provider_header_wait_worker_slot().ok_or_else(|| {
ProviderError::transport(format!(
"provider response header wait worker limit reached ({PROVIDER_HEADER_WAIT_WORKER_LIMIT}); retry after stalled provider requests finish"
))
})?;
let worker = thread::Builder::new()
.name("provider-header-wait".to_string())
.spawn(move || {
let _worker_slot = worker_slot;
let _ = sender.send(operation());
})
.map_err(|error| ProviderError::transport(error.to_string()))?;
match recv_timeout_cancellable(&receiver, idle_timeout, cancellation) {
Ok(result) => {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
result
}
Err(mpsc::RecvTimeoutError::Timeout) => {
drop(receiver);
let joined =
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
let cleanup = if joined {
""
} else {
"; provider header worker detached after 100ms grace period; it still occupies capacity until it exits"
};
Err(ProviderError::transport(format!(
"provider response header timeout after {}s without bytes{}",
idle_timeout.as_secs(),
cleanup
))
.into())
}
Err(mpsc::RecvTimeoutError::Disconnected) if cancellation.is_canceled() => {
drop(receiver);
let joined =
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
let error: anyhow::Error = cancellation.check().unwrap_err();
Err(if joined {
error
} else {
provider_cleanup_context(
error,
"provider header worker cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
)
})
}
Err(mpsc::RecvTimeoutError::Disconnected) => {
drop(receiver);
let joined =
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
let error: anyhow::Error = ProviderError::transport(
"provider response header wait worker disconnected".to_string(),
)
.into();
Err(if joined {
error
} else {
provider_cleanup_context(
error,
"provider header worker cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
)
})
}
}
}
pub(crate) fn join_provider_worker_with_timeout(
handle: thread::JoinHandle<()>,
timeout: Duration,
) -> bool {
let deadline = Instant::now() + timeout;
loop {
if handle.is_finished() {
let _ = handle.join();
return true;
}
if Instant::now() >= deadline {
return false;
}
thread::sleep(PROVIDER_WORKER_SHUTDOWN_POLL_INTERVAL);
}
}
fn recv_timeout_cancellable<T>(
receiver: &mpsc::Receiver<T>,
timeout: Duration,
cancellation: &AgentCancellation,
) -> Result<T, mpsc::RecvTimeoutError> {
const CANCEL_POLL_INTERVAL: Duration = Duration::from_millis(25);
let start = Instant::now();
loop {
match receiver.try_recv() {
Ok(value) => return Ok(value),
Err(mpsc::TryRecvError::Disconnected) => {
return Err(mpsc::RecvTimeoutError::Disconnected);
}
Err(mpsc::TryRecvError::Empty) => {}
}
if cancellation.is_canceled() {
return Err(mpsc::RecvTimeoutError::Disconnected);
}
let elapsed = start.elapsed();
if elapsed >= timeout {
return Err(mpsc::RecvTimeoutError::Timeout);
}
let remaining = timeout.saturating_sub(elapsed);
let wait = remaining.min(CANCEL_POLL_INTERVAL);
match receiver.recv_timeout(wait) {
Ok(value) => return Ok(value),
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
return Err(mpsc::RecvTimeoutError::Disconnected);
}
}
}
}
fn provider_cleanup_context(error: anyhow::Error, detail: &str) -> anyhow::Error {
let summary = error.to_string();
error.context(format!("{detail}: {summary}"))
}
pub(crate) fn stream_reader_with_idle_timeout<R>(
mut reader: R,
idle_timeout: Duration,
cancellation: &AgentCancellation,
semantic_deadline: Option<&std::sync::atomic::AtomicU64>,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()>
where
R: Read + Send + 'static,
{
let worker_slot = reserve_provider_stream_reader_worker_slot().ok_or_else(|| {
ProviderError::transport(format!(
"provider stream reader worker limit reached ({PROVIDER_STREAM_READER_WORKER_LIMIT}); retry after stalled provider streams finish"
))
})?;
let (sender, receiver) = mpsc::sync_channel::<std::io::Result<Vec<u8>>>(1);
let worker = thread::Builder::new()
.name("provider-stream-reader".to_string())
.spawn(move || {
let _worker_slot = worker_slot;
let mut buffer = [0_u8; 8192];
loop {
let message = match reader.read(&mut buffer) {
Ok(0) => Ok(Vec::new()),
Ok(read) => Ok(buffer[..read].to_vec()),
Err(error) => Err(error),
};
let done = matches!(&message, Ok(bytes) if bytes.is_empty()) || message.is_err();
if sender.send(message).is_err() || done {
break;
}
}
})
.map_err(|error| ProviderError::transport(error.to_string()))?;
let mut cleanup = StreamReaderCleanup {
receiver: Some(receiver),
worker: Some(worker),
};
let mut pending_utf8 = Vec::new();
loop {
let timeout = stream_receive_timeout(idle_timeout, semantic_deadline);
match recv_timeout_cancellable(cleanup.receiver.as_ref().unwrap(), timeout, cancellation) {
Ok(Ok(bytes)) if bytes.is_empty() => break,
Ok(Ok(bytes)) => {
if let Err(error) = dispatch_utf8_bytes(&mut pending_utf8, &bytes, on_chunk) {
let joined = cleanup.finish();
return Err(if joined {
error
} else {
provider_cleanup_context(
error,
"provider stream reader cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
)
});
}
}
Ok(Err(error)) => {
let error: anyhow::Error = ProviderError::transport(error.to_string()).into();
let joined = cleanup.finish();
return Err(if joined {
error
} else {
provider_cleanup_context(
error,
"provider stream reader cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
)
});
}
Err(mpsc::RecvTimeoutError::Timeout) => {
let error: anyhow::Error = if semantic_deadline.is_some_and(deadline_reached) {
ProviderError::stream_terminal(
"provider stream no semantic progress before timeout",
)
.into()
} else {
ProviderError::stream_terminal(format!(
"provider stream idle timeout after {}s without bytes",
idle_timeout.as_secs()
))
.into()
};
let joined = cleanup.finish();
return Err(if joined {
error
} else {
provider_cleanup_context(
error,
"provider stream reader cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
)
});
}
Err(mpsc::RecvTimeoutError::Disconnected) if cancellation.is_canceled() => {
let joined = cleanup.finish();
let error: anyhow::Error = cancellation.check().unwrap_err();
if joined {
return Err(error);
}
return Err(provider_cleanup_context(
error,
"provider stream reader cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
));
}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
}
}
if !pending_utf8.is_empty() {
let chunk = std::str::from_utf8(&pending_utf8)?;
if let Err(error) = on_chunk(chunk) {
let joined = cleanup.finish();
return Err(if joined {
error
} else {
provider_cleanup_context(
error,
"provider stream reader cleanup incomplete after 100ms grace period; detached worker still occupies capacity until it exits",
)
});
}
}
let _ = cleanup.finish();
Ok(())
}
struct StreamReaderCleanup {
receiver: Option<mpsc::Receiver<std::io::Result<Vec<u8>>>>,
worker: Option<thread::JoinHandle<()>>,
}
impl StreamReaderCleanup {
fn finish(&mut self) -> bool {
self.receiver.take();
self.worker.take().is_none_or(|worker| {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT)
})
}
}
impl Drop for StreamReaderCleanup {
fn drop(&mut self) {
self.receiver.take();
if let Some(worker) = self.worker.take() {
let _ = join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
}
}
}
fn stream_receive_timeout(
idle_timeout: Duration,
semantic_deadline: Option<&std::sync::atomic::AtomicU64>,
) -> Duration {
semantic_deadline
.and_then(|deadline| {
let deadline = deadline.load(Ordering::SeqCst);
(deadline > 0).then(|| {
Duration::from_millis(deadline.saturating_sub(monotonic_millis())).min(idle_timeout)
})
})
.unwrap_or(idle_timeout)
}
fn deadline_reached(deadline: &std::sync::atomic::AtomicU64) -> bool {
let deadline = deadline.load(Ordering::SeqCst);
deadline > 0 && deadline <= monotonic_millis()
}
pub(crate) fn provider_stream_deadline_after(duration: Duration) -> u64 {
monotonic_millis().saturating_add(duration.as_millis() as u64)
}
fn monotonic_millis() -> u64 {
static START: OnceLock<Instant> = OnceLock::new();
START.get_or_init(Instant::now).elapsed().as_millis() as u64
}
pub(crate) fn dispatch_utf8_bytes(
pending: &mut Vec<u8>,
bytes: &[u8],
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
if pending.is_empty() {
return dispatch_utf8_bytes_without_pending(pending, bytes, on_chunk);
}
pending.extend_from_slice(bytes);
match std::str::from_utf8(pending) {
Ok(chunk) => {
on_chunk(chunk)?;
pending.clear();
Ok(())
}
Err(error) => {
let valid_up_to = error.valid_up_to();
if valid_up_to > 0 {
let valid = std::str::from_utf8(&pending[..valid_up_to])?;
on_chunk(valid)?;
pending.drain(..valid_up_to);
}
if error.error_len().is_some() {
anyhow::bail!("provider stream contained invalid UTF-8");
}
Ok(())
}
}
}
fn dispatch_utf8_bytes_without_pending(
pending: &mut Vec<u8>,
bytes: &[u8],
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
match std::str::from_utf8(bytes) {
Ok(chunk) => on_chunk(chunk),
Err(error) => {
let valid_up_to = error.valid_up_to();
if valid_up_to > 0 {
let valid = std::str::from_utf8(&bytes[..valid_up_to])?;
on_chunk(valid)?;
}
if error.error_len().is_some() {
anyhow::bail!("provider stream contained invalid UTF-8");
}
pending.extend_from_slice(&bytes[valid_up_to..]);
Ok(())
}
}
}
static PROVIDER_CLIENT: OnceLock<reqwest::blocking::Client> = OnceLock::new();
pub(crate) fn provider_client() -> anyhow::Result<reqwest::blocking::Client> {
if let Some(client) = PROVIDER_CLIENT.get() {
return Ok(client.clone());
}
let new_client = provider_client_builder().build()?;
Ok(PROVIDER_CLIENT.get_or_init(|| new_client).clone())
}
#[cfg(test)]
fn provider_client_is_cached() -> bool {
PROVIDER_CLIENT.get().is_some()
}
pub(crate) fn provider_client_builder() -> reqwest::blocking::ClientBuilder {
reqwest::blocking::Client::builder()
.connect_timeout(PROVIDER_CONNECT_TIMEOUT)
.timeout(PROVIDER_STREAM_WALL_TIMEOUT)
}
pub(crate) fn read_bounded_redacted_response_body<R>(reader: R) -> String
where
R: Read + Send + 'static,
{
read_bounded_redacted_response_body_with_limit(
reader,
PROVIDER_ERROR_BODY_MAX_BYTES,
PROVIDER_ERROR_BODY_IDLE_TIMEOUT,
)
}
fn read_bounded_redacted_response_body_with_limit<R>(
reader: R,
max_bytes: usize,
idle_timeout: Duration,
) -> String
where
R: Read + Send + 'static,
{
let read_limit = max_bytes.saturating_add(PROVIDER_ERROR_BODY_REDACTION_SLACK_BYTES);
match read_body_prefix_with_idle_timeout(reader, read_limit, idle_timeout) {
Ok((bytes, truncated)) => format_redacted_response_body(bytes, max_bytes, truncated),
Err(error) => format!(
"<failed to read provider error body: {}>",
redact_sensitive_text(&error.to_string())
),
}
}
fn read_body_prefix_with_idle_timeout<R>(
mut reader: R,
max_bytes: usize,
idle_timeout: Duration,
) -> anyhow::Result<(Vec<u8>, bool)>
where
R: Read + Send + 'static,
{
let worker_slot = reserve_provider_error_body_reader_worker_slot().ok_or_else(|| {
ProviderError::transport(format!("provider error-body reader worker limit reached ({PROVIDER_ERROR_BODY_READER_WORKER_LIMIT}); retry after stalled provider requests finish"))
})?;
let (sender, receiver) = mpsc::sync_channel::<std::io::Result<Vec<u8>>>(1);
let worker = thread::Builder::new()
.name("provider-error-body-reader".to_string())
.spawn(move || {
let _worker_slot = worker_slot;
let mut buffer = [0_u8; 8192];
loop {
let message = match reader.read(&mut buffer) {
Ok(0) => Ok(Vec::new()),
Ok(read) => Ok(buffer[..read].to_vec()),
Err(error) => Err(error),
};
let done = matches!(&message, Ok(bytes) if bytes.is_empty()) || message.is_err();
if sender.send(message).is_err() || done {
break;
}
}
})
.map_err(|error| ProviderError::transport(error.to_string()))?;
let mut worker = Some(worker);
let mut body = Vec::with_capacity(max_bytes.min(8192));
loop {
match receiver.recv_timeout(idle_timeout) {
Ok(Ok(bytes)) if bytes.is_empty() => {
if let Some(worker) = worker.take() {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
}
return Ok((body, false));
}
Ok(Ok(bytes)) => {
let remaining = max_bytes.saturating_sub(body.len());
if bytes.len() > remaining {
body.extend_from_slice(&bytes[..remaining]);
if let Some(worker) = worker.take() {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
}
return Ok((body, true));
}
body.extend_from_slice(&bytes);
}
Ok(Err(error)) => {
if let Some(worker) = worker.take() {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
}
return Err(ProviderError::transport(error.to_string()).into());
}
Err(mpsc::RecvTimeoutError::Timeout) => {
drop(receiver);
if let Some(worker) = worker.take() {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
}
return Err(ProviderError::transport(format!(
"provider error body read idle timeout after {}s",
idle_timeout.as_secs()
))
.into());
}
Err(mpsc::RecvTimeoutError::Disconnected) => {
if let Some(worker) = worker.take() {
join_provider_worker_with_timeout(worker, PROVIDER_WORKER_SHUTDOWN_TIMEOUT);
}
return Ok((body, false));
}
}
}
}
fn format_redacted_response_body(bytes: Vec<u8>, max_bytes: usize, truncated: bool) -> String {
let body = redact_response_body(&String::from_utf8_lossy(&bytes));
truncate_redacted_response_body(body, max_bytes, truncated || bytes.len() > max_bytes)
}
fn truncate_redacted_response_body(mut body: String, max_bytes: usize, truncated: bool) -> String {
if body.len() > max_bytes {
let mut end = max_bytes;
while !body.is_char_boundary(end) {
end -= 1;
}
body.truncate(end);
}
if truncated {
if !body.is_empty() {
body.push('\n');
}
body.push_str("<truncated provider error body>");
}
body
}
fn redact_response_body(body: &str) -> String {
redact_sensitive_text(body)
}
pub(crate) fn sanitize_provider_error_url(url: &str) -> String {
let Ok(mut parsed) = reqwest::Url::parse(url) else {
return redact_sensitive_text(url);
};
let _ = parsed.set_username("");
let _ = parsed.set_password(None);
parsed.set_query(None);
parsed.set_fragment(None);
parsed.to_string()
}
#[cfg(test)]
mod tests {
use super::{
HttpRequest, HttpTransport, PROVIDER_ERROR_BODY_READER_WORKER_LIMIT,
PROVIDER_HEADER_WAIT_WORKER_LIMIT, PROVIDER_READ_IDLE_TIMEOUT, provider_client,
provider_client_builder, provider_client_is_cached,
read_bounded_redacted_response_body_with_limit, redact_response_body,
reserve_provider_worker_slot, run_provider_header_wait_with_idle_timeout,
sanitize_provider_error_url, stream_reader_with_idle_timeout,
};
use crate::agent::cancellation::{AgentCancellation, AgentCancellationHandle, is_run_canceled};
use std::{
io::{self, Read},
sync::{
Arc, Mutex, OnceLock,
atomic::{AtomicBool, AtomicUsize, Ordering},
mpsc,
},
thread,
time::{Duration, Instant},
};
static HEADER_WAIT_TEST_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
static STREAM_READER_TEST_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
static ERROR_BODY_TEST_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
fn header_wait_test_lock() -> std::sync::MutexGuard<'static, ()> {
HEADER_WAIT_TEST_LOCK
.get_or_init(|| Mutex::new(()))
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn error_body_test_lock() -> std::sync::MutexGuard<'static, ()> {
ERROR_BODY_TEST_LOCK
.get_or_init(|| Mutex::new(()))
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn stream_reader_test_lock() -> std::sync::MutexGuard<'static, ()> {
STREAM_READER_TEST_LOCK
.get_or_init(|| Mutex::new(()))
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
#[test]
fn provider_error_body_reader_worker_slots_are_hard_capped_and_released() {
let workers = AtomicUsize::new(0);
let mut slots = Vec::new();
while let Some(slot) =
reserve_provider_worker_slot(&workers, PROVIDER_ERROR_BODY_READER_WORKER_LIMIT)
{
slots.push(slot);
assert!(slots.len() <= PROVIDER_ERROR_BODY_READER_WORKER_LIMIT);
}
assert!(
reserve_provider_worker_slot(&workers, PROVIDER_ERROR_BODY_READER_WORKER_LIMIT)
.is_none()
);
drop(slots);
assert!(
reserve_provider_worker_slot(&workers, PROVIDER_ERROR_BODY_READER_WORKER_LIMIT)
.is_some()
);
}
struct CooperativeReader {
released: Arc<AtomicBool>,
sent_chunk: bool,
}
impl Read for CooperativeReader {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
if !self.sent_chunk {
self.sent_chunk = true;
buf[0] = b'x';
return Ok(1);
}
while !self.released.load(Ordering::SeqCst) {
std::thread::sleep(Duration::from_millis(5));
}
Ok(0)
}
}
struct KeepaliveReader;
impl Read for KeepaliveReader {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
std::thread::sleep(Duration::from_millis(2));
buf[0] = b'\n';
Ok(1)
}
}
#[test]
fn provider_client_reuses_static_client_without_headers() {
let _first = provider_client().unwrap();
assert!(provider_client_is_cached());
let _second = provider_client().unwrap();
assert!(provider_client_is_cached());
}
#[test]
fn provider_client_uses_connect_and_read_idle_timeouts() {
provider_client_builder().build().unwrap();
assert_eq!(PROVIDER_READ_IDLE_TIMEOUT, Duration::from_secs(120));
}
#[test]
fn http_request_debug_redacts_auth_sensitive_header_values() {
let request = HttpRequest {
method: "POST".to_string(),
url: "https://provider.test/v1".to_string(),
headers: std::collections::BTreeMap::from([
(
"Authorization".to_string(),
"Bearer fake-token-FAKE".to_string(),
),
("X-Api-Key".to_string(), "sk-test-FAKE".to_string()),
(
"custom-token-header".to_string(),
"token-secret-FAKE".to_string(),
),
("accept".to_string(), "application/json".to_string()),
]),
body: serde_json::Value::Null,
};
let debug = format!("{request:?}");
assert!(debug.contains("Authorization"));
assert!(debug.contains("X-Api-Key"));
assert!(debug.contains("custom-token-header"));
assert!(debug.contains("application/json"));
assert!(!debug.contains("Bearer fake-token-FAKE"));
assert!(!debug.contains("sk-test-FAKE"));
assert!(!debug.contains("token-secret-FAKE"));
}
#[test]
fn reqwest_transport_rejects_invalid_http_method() {
let request = HttpRequest {
method: "bad method".to_string(),
url: "https://provider.test/v1".to_string(),
headers: std::collections::BTreeMap::new(),
body: serde_json::Value::Null,
};
let error = super::ReqwestHttpTransport
.stream_json(request, &mut |_| Ok(()))
.unwrap_err()
.to_string();
assert!(error.contains("invalid provider HTTP method"), "{error}");
}
#[test]
fn reqwest_transport_honors_http_request_method() {
let _lock = header_wait_test_lock();
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!(
"http://127.0.0.1:{}/v1/models",
listener.local_addr().unwrap().port()
);
let server = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let mut request_bytes = [0_u8; 1024];
let read = stream.read(&mut request_bytes).unwrap();
let request = String::from_utf8_lossy(&request_bytes[..read]).to_string();
let response = "HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n";
std::io::Write::write_all(&mut stream, response.as_bytes()).unwrap();
request
});
let request = HttpRequest {
method: "GET".to_string(),
url,
headers: std::collections::BTreeMap::new(),
body: serde_json::Value::Null,
};
super::ReqwestHttpTransport
.stream_json_cancellable(request, &AgentCancellation::default(), &mut |_| Ok(()))
.unwrap();
let raw_request = server.join().unwrap();
assert!(
raw_request.starts_with("GET /v1/models HTTP/1.1"),
"{raw_request}"
);
}
struct StalledReader;
impl Read for StalledReader {
fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> {
std::thread::sleep(Duration::from_secs(5));
Ok(0)
}
}
#[test]
fn reqwest_transport_send_error_sanitizes_credential_url() {
let _lock = header_wait_test_lock();
let secret = "sk-test-FAKE";
let request = HttpRequest {
method: "POST".to_string(),
url: format!("http://user:password@127.0.0.1:9/v1/responses?api_key={secret}#frag"),
headers: std::collections::BTreeMap::new(),
body: serde_json::Value::Null,
};
let error = super::ReqwestHttpTransport
.stream_json_cancellable(request, &AgentCancellation::default(), &mut |_| Ok(()))
.unwrap_err()
.to_string();
assert!(error.contains("provider request send failed"), "{error}");
assert!(error.contains("http://127.0.0.1:9/v1/responses"), "{error}");
assert!(!error.contains(secret), "{error}");
assert!(!error.contains("user"), "{error}");
assert!(!error.contains("password"), "{error}");
assert!(!error.contains("api_key"), "{error}");
assert!(!error.contains("frag"), "{error}");
}
#[test]
fn provider_error_body_timeout_returns_placeholder_promptly() {
let _lock = error_body_test_lock();
let started = Instant::now();
let body = read_bounded_redacted_response_body_with_limit(
StalledReader,
1024,
Duration::from_millis(100),
);
assert!(started.elapsed() < Duration::from_secs(1));
assert!(
body.contains("<failed to read provider error body"),
"{body}"
);
assert!(body.contains("idle timeout"), "{body}");
}
#[test]
fn provider_error_body_normal_read_still_returns_text() {
let _lock = error_body_test_lock();
let body = read_bounded_redacted_response_body_with_limit(
std::io::Cursor::new(b"plain failure body".to_vec()),
1024,
Duration::from_secs(1),
);
assert_eq!(body, "plain failure body");
}
#[test]
fn provider_header_worker_capacity_is_hard_capped_and_released() {
let workers = AtomicUsize::new(0);
let mut slots = Vec::new();
while let Some(slot) =
reserve_provider_worker_slot(&workers, PROVIDER_HEADER_WAIT_WORKER_LIMIT)
{
slots.push(slot);
assert!(slots.len() <= PROVIDER_HEADER_WAIT_WORKER_LIMIT);
}
assert_eq!(slots.len(), PROVIDER_HEADER_WAIT_WORKER_LIMIT);
assert!(
reserve_provider_worker_slot(&workers, PROVIDER_HEADER_WAIT_WORKER_LIMIT).is_none()
);
drop(slots);
assert!(
reserve_provider_worker_slot(&workers, PROVIDER_HEADER_WAIT_WORKER_LIMIT).is_some()
);
}
#[test]
fn provider_header_timeout_returns_before_blocking_send_finishes() {
let _lock = header_wait_test_lock();
let finished = Arc::new(AtomicBool::new(false));
let operation_finished = Arc::clone(&finished);
let (release_tx, release_rx) = mpsc::channel();
let started = Instant::now();
let error = run_provider_header_wait_with_idle_timeout(
move || {
release_rx.recv().unwrap();
operation_finished.store(true, Ordering::SeqCst);
Ok(())
},
Duration::from_millis(100),
&AgentCancellation::default(),
)
.unwrap_err()
.to_string();
assert!(started.elapsed() < Duration::from_secs(1));
assert!(
error.contains("provider response header timeout"),
"{error}"
);
assert!(error.contains("detached"), "{error}");
assert!(
error.contains("occupies capacity until it exits"),
"{error}"
);
assert!(!finished.load(Ordering::SeqCst));
release_tx.send(()).unwrap();
let deadline = Instant::now() + Duration::from_secs(1);
while !finished.load(Ordering::SeqCst) {
assert!(Instant::now() < deadline, "header worker did not exit");
thread::yield_now();
}
}
#[test]
fn provider_stream_idle_timeout_returns_error_with_stalled_reader() {
let _lock = stream_reader_test_lock();
let started = Instant::now();
let mut chunks = Vec::new();
let error = stream_reader_with_idle_timeout(
StalledReader,
Duration::from_millis(100),
&AgentCancellation::default(),
None,
&mut |chunk| {
chunks.push(chunk.to_string());
Ok(())
},
)
.unwrap_err()
.to_string();
assert!(started.elapsed() < Duration::from_secs(1));
assert!(error.contains("idle timeout"), "{error}");
assert!(chunks.is_empty());
}
#[test]
fn provider_stream_callback_error_cleans_up_and_releases_reader() {
let _lock = stream_reader_test_lock();
let baseline = super::PROVIDER_STREAM_READER_WORKERS.load(Ordering::SeqCst);
let released = Arc::new(AtomicBool::new(false));
let reader = CooperativeReader {
released: Arc::clone(&released),
sent_chunk: false,
};
let started = Instant::now();
let error = stream_reader_with_idle_timeout(
reader,
Duration::from_secs(5),
&AgentCancellation::default(),
None,
&mut |_| anyhow::bail!("sink failed"),
)
.unwrap_err()
.to_string();
assert!(error.contains("sink failed"), "{error}");
assert!(error.contains("cleanup incomplete"), "{error}");
assert!(started.elapsed() < Duration::from_secs(1));
released.store(true, Ordering::SeqCst);
let deadline = Instant::now() + Duration::from_secs(1);
while super::PROVIDER_STREAM_READER_WORKERS.load(Ordering::SeqCst) > baseline {
assert!(Instant::now() < deadline, "reader worker did not drain");
std::thread::sleep(Duration::from_millis(5));
}
}
#[test]
fn provider_stream_reader_callback_failures_are_hard_capped() {
let _lock = stream_reader_test_lock();
let baseline = super::PROVIDER_STREAM_READER_WORKERS.load(Ordering::SeqCst);
let mut callback_failures = 0;
for _ in 0..=(super::PROVIDER_STREAM_READER_WORKER_LIMIT * 2) {
let error = stream_reader_with_idle_timeout(
std::io::Cursor::new(b"chunk".to_vec()),
Duration::from_secs(5),
&AgentCancellation::default(),
None,
&mut |_| {
callback_failures += 1;
anyhow::bail!("sink failed")
},
)
.unwrap_err()
.to_string();
assert_eq!(error, "sink failed");
}
assert_eq!(
callback_failures,
super::PROVIDER_STREAM_READER_WORKER_LIMIT * 2 + 1
);
let deadline = Instant::now() + Duration::from_secs(1);
while super::PROVIDER_STREAM_READER_WORKERS.load(Ordering::SeqCst) > baseline {
assert!(Instant::now() < deadline, "reader workers did not drain");
std::thread::sleep(Duration::from_millis(5));
}
}
#[test]
fn provider_stream_semantic_deadline_beats_byte_idle_timeout() {
let _lock = stream_reader_test_lock();
let deadline = std::sync::atomic::AtomicU64::new(super::provider_stream_deadline_after(
Duration::from_millis(50),
));
let error = stream_reader_with_idle_timeout(
KeepaliveReader,
Duration::from_secs(1),
&AgentCancellation::default(),
Some(&deadline),
&mut |_| Ok(()),
)
.unwrap_err()
.to_string();
assert!(error.contains("no semantic progress"), "{error}");
assert!(!error.contains("idle timeout"), "{error}");
}
#[test]
fn provider_header_wait_observes_cancellation_before_idle_timeout() {
let _lock = header_wait_test_lock();
let (cancellation, handle): (AgentCancellation, AgentCancellationHandle) =
AgentCancellation::default().child_token();
let (release_tx, release_rx) = mpsc::channel();
let started = Instant::now();
let canceler = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(50));
handle.cancel();
});
let error = run_provider_header_wait_with_idle_timeout(
move || {
release_rx.recv().unwrap();
Ok(())
},
Duration::from_secs(5),
&cancellation,
)
.unwrap_err();
canceler.join().unwrap();
assert!(is_run_canceled(&error));
assert!(started.elapsed() < Duration::from_secs(1));
release_tx.send(()).unwrap();
}
#[test]
fn provider_header_completion_wins_cancellation_race() {
let (cancellation, handle) = AgentCancellation::default().child_token();
let (sender, receiver) = mpsc::sync_channel(1);
sender.send("completed").unwrap();
handle.cancel();
let result =
super::recv_timeout_cancellable(&receiver, Duration::from_secs(1), &cancellation)
.unwrap();
assert_eq!(result, "completed");
}
struct CompletionRaceReader {
started: mpsc::Sender<()>,
release: mpsc::Receiver<()>,
completed: Option<mpsc::Sender<()>>,
}
impl Read for CompletionRaceReader {
fn read(&mut self, _buffer: &mut [u8]) -> io::Result<usize> {
self.started.send(()).unwrap();
self.release.recv().unwrap();
Ok(0)
}
}
impl Drop for CompletionRaceReader {
fn drop(&mut self) {
self.completed.take().unwrap().send(()).unwrap();
}
}
#[test]
fn provider_stream_completion_wins_cancellation_race() {
let _lock = stream_reader_test_lock();
let (cancellation, handle) = AgentCancellation::default().child_token();
let (started_tx, started_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let (completed_tx, completed_rx) = mpsc::channel();
let (result_tx, result_rx) = mpsc::channel();
let worker = thread::spawn(move || {
let result = stream_reader_with_idle_timeout(
CompletionRaceReader {
started: started_tx,
release: release_rx,
completed: Some(completed_tx),
},
Duration::from_secs(1),
&cancellation,
None,
&mut |_| Ok(()),
);
result_tx.send(result).unwrap();
});
started_rx.recv().unwrap();
release_tx.send(()).unwrap();
completed_rx.recv().unwrap();
handle.cancel();
result_rx
.recv_timeout(Duration::from_secs(1))
.unwrap()
.unwrap();
worker.join().unwrap();
}
#[test]
fn provider_stream_idle_wait_observes_cancellation_before_idle_timeout() {
let _lock = stream_reader_test_lock();
let (cancellation, handle): (AgentCancellation, AgentCancellationHandle) =
AgentCancellation::default().child_token();
let started = Instant::now();
let canceler = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(50));
handle.cancel();
});
let error = stream_reader_with_idle_timeout(
StalledReader,
Duration::from_secs(5),
&cancellation,
None,
&mut |_| Ok(()),
)
.unwrap_err();
canceler.join().unwrap();
assert!(is_run_canceled(&error));
assert!(started.elapsed() < Duration::from_secs(1));
}
#[test]
fn provider_header_cancellation_does_not_stop_blocked_worker() {
let _lock = header_wait_test_lock();
let baseline = super::PROVIDER_HEADER_WAIT_WORKERS.load(Ordering::SeqCst);
let (cancellation, cancel_handle) = AgentCancellation::default().child_token();
let (started_tx, started_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let finished = Arc::new(AtomicBool::new(false));
let operation_finished = Arc::clone(&finished);
let caller = thread::spawn(move || {
run_provider_header_wait_with_idle_timeout(
move || {
started_tx.send(()).unwrap();
release_rx.recv().unwrap();
operation_finished.store(true, Ordering::SeqCst);
Ok::<_, anyhow::Error>(())
},
Duration::from_secs(5),
&cancellation,
)
});
started_rx.recv().unwrap();
cancel_handle.cancel();
let error = caller.join().unwrap().unwrap_err();
assert!(is_run_canceled(&error), "unexpected error: {error}");
assert!(!finished.load(Ordering::SeqCst));
release_tx.send(()).unwrap();
let deadline = Instant::now() + Duration::from_secs(1);
while !finished.load(Ordering::SeqCst) {
assert!(Instant::now() < deadline, "header worker did not exit");
thread::yield_now();
}
while super::PROVIDER_HEADER_WAIT_WORKERS.load(Ordering::SeqCst) > baseline {
assert!(
Instant::now() < deadline,
"header worker slot did not release"
);
thread::yield_now();
}
}
struct ControlledStreamReader {
started: mpsc::Sender<()>,
release: mpsc::Receiver<()>,
finished: Arc<AtomicBool>,
}
impl Read for ControlledStreamReader {
fn read(&mut self, _buffer: &mut [u8]) -> io::Result<usize> {
self.started.send(()).unwrap();
self.release.recv().unwrap();
self.finished.store(true, Ordering::SeqCst);
Ok(0)
}
}
#[test]
fn provider_stream_cancellation_does_not_stop_blocked_reader() {
let _lock = stream_reader_test_lock();
let baseline = super::PROVIDER_STREAM_READER_WORKERS.load(Ordering::SeqCst);
let (cancellation, cancel_handle) = AgentCancellation::default().child_token();
let (started_tx, started_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let finished = Arc::new(AtomicBool::new(false));
let reader_finished = Arc::clone(&finished);
let caller = thread::spawn(move || {
stream_reader_with_idle_timeout(
ControlledStreamReader {
started: started_tx,
release: release_rx,
finished: reader_finished,
},
Duration::from_secs(5),
&cancellation,
None,
&mut |_| Ok(()),
)
});
started_rx.recv().unwrap();
cancel_handle.cancel();
let error = caller.join().unwrap().unwrap_err();
assert!(is_run_canceled(&error), "unexpected error: {error}");
assert!(!finished.load(Ordering::SeqCst));
release_tx.send(()).unwrap();
let deadline = Instant::now() + Duration::from_secs(1);
while !finished.load(Ordering::SeqCst) {
assert!(Instant::now() < deadline, "stream reader did not exit");
thread::yield_now();
}
while super::PROVIDER_STREAM_READER_WORKERS.load(Ordering::SeqCst) > baseline {
assert!(
Instant::now() < deadline,
"stream worker slot did not release"
);
thread::yield_now();
}
}
#[test]
fn dispatch_utf8_bytes_preserves_split_sequences() {
let mut pending = Vec::new();
let mut chunks = Vec::new();
super::dispatch_utf8_bytes(&mut pending, "hello ".as_bytes(), &mut |chunk| {
chunks.push(chunk.to_string());
Ok(())
})
.unwrap();
let snowman = "☃".as_bytes();
super::dispatch_utf8_bytes(&mut pending, &snowman[..1], &mut |chunk| {
chunks.push(chunk.to_string());
Ok(())
})
.unwrap();
super::dispatch_utf8_bytes(&mut pending, &snowman[1..], &mut |chunk| {
chunks.push(chunk.to_string());
Ok(())
})
.unwrap();
assert_eq!(chunks, vec!["hello ".to_string(), "☃".to_string()]);
}
#[test]
fn provider_error_body_redacts_secret_split_at_display_limit() {
let _lock = error_body_test_lock();
let secret = format!("sk-{}", "x".repeat(24));
let body = format!("{} {secret}", "x".repeat(29));
let redacted = read_bounded_redacted_response_body_with_limit(
std::io::Cursor::new(body.into_bytes()),
33,
Duration::from_secs(1),
);
assert!(
redacted.contains("<truncated provider error body>"),
"{redacted}"
);
assert!(!redacted.contains(&secret), "{redacted}");
assert!(!redacted.contains("sk-"));
}
#[test]
fn non_success_body_read_is_bounded_and_redacted() {
let _lock = error_body_test_lock();
let secret = format!("sk-{}", "x".repeat(24));
let body = format!("{}{}", "x".repeat(32), secret);
let redacted = read_bounded_redacted_response_body_with_limit(
std::io::Cursor::new(body.into_bytes()),
32,
Duration::from_secs(1),
);
assert!(redacted.contains("<truncated provider error body>"));
assert!(!redacted.contains(&secret));
assert!(redacted.len() < 128, "{redacted}");
}
#[test]
fn provider_error_url_redacts_userinfo_query_and_fragment() {
let redacted = sanitize_provider_error_url(
"https://user:password@example.test/v1/responses?api_key=secret#token",
);
assert_eq!(redacted, "https://example.test/v1/responses");
assert!(!redacted.contains("user"));
assert!(!redacted.contains("password"));
assert!(!redacted.contains("secret"));
assert!(!redacted.contains("token"));
}
#[test]
fn redact_response_body_removes_failed_response_secrets() {
let api_key = format!("sk-{}", "x".repeat(24));
let jwt = format!(
"eyJ{}.eyJ{}.{}",
"x".repeat(8),
"x".repeat(8),
"x".repeat(8)
);
let bearer = format!("Bearer {}", "x".repeat(24));
let account_id = format!("acct-{}", "x".repeat(12));
let access_token = format!("plain-access-{}", "x".repeat(12));
let api_key_json = format!("plain-api-{}", "x".repeat(12));
let body = format!(
r#"{{"status":401,"error":"invalid credentials","key":"{api_key}","authorization":"{bearer}","token":"{jwt}","accountId":"{account_id}","access_token":"{access_token}","api_key":"{api_key_json}","input_tokens":42,"token_estimate":9,"retry":false}}"#
);
let redacted = redact_response_body(&body);
assert!(redacted.contains(r#""status":401"#));
assert!(redacted.contains(r#""error":"invalid credentials""#));
assert!(redacted.contains(r#""input_tokens":42"#));
assert!(redacted.contains(r#""token_estimate":9"#));
assert!(redacted.contains(r#""retry":false"#));
assert!(redacted.matches("<redacted>").count() >= 6);
assert!(!redacted.contains(&api_key));
assert!(!redacted.contains(&bearer));
assert!(!redacted.contains(&jwt));
assert!(!redacted.contains(&account_id));
assert!(!redacted.contains(&access_token));
assert!(!redacted.contains(&api_key_json));
assert!(!redacted.contains("sk-"));
assert!(!redacted.contains("eyJ"));
}
}