use crate::{
cancellation::AgentCancellation, output::redact_sensitive_text, providers::error::ProviderError,
};
use serde_json::Value;
use std::{
collections::BTreeMap,
io::Read,
sync::{
Arc, 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_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: Arc<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" | "x-codex-turn-state"
) || name.contains("token")
|| name.contains("key")
}
#[derive(Clone, Default)]
pub struct HttpResponseMetadata {
pub request_id: Option<String>,
pub codex_turn_state: Option<String>,
}
impl std::fmt::Debug for HttpResponseMetadata {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("HttpResponseMetadata")
.field("request_id", &self.request_id)
.field("codex_turn_state_present", &self.codex_turn_state.is_some())
.finish()
}
}
pub trait HttpTransport {
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(HttpResponseMetadata),
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_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(HttpResponseMetadata),
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);
let codex_turn_state = response
.headers()
.get(super::codex_session::TURN_STATE_HEADER)
.and_then(|value| value.to_str().ok())
.filter(|value| super::codex_session::valid_turn_state(value))
.map(str::to_owned);
on_response_id(HttpResponseMetadata {
request_id,
codex_turn_state,
});
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: Arc<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.as_ref()).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 {
crate::thread_join::join_with_timeout(handle, timeout).is_ok()
}
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())
}
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()
}