mirror of
https://github.com/openai/codex.git
synced 2026-09-28 08:43:01 +08:00
Make diagnostic report uploads resilient to slow networks (#42096)
## Why Diagnostic reports can span several envelopes, and the previous 10-second shared network budget could expire before slow uploads and their attachments completed. ## What changed - Give each diagnostic report a single five-minute deadline shared by the event, attachments, retries, and retry backoff. - Stop reading or sending later attachments once the deadline or a Sentry rate limit is reached. - Limit the app server to three concurrent report uploads and return an overloaded JSON-RPC error for additional requests until a slot is released. ## Testing Add coverage for slow multi-envelope reports, deadline-aware retries, skipped attachments, rate-limit handling, and concurrency-slot release after failures. GitOrigin-RevId: bdacf6c9731df16d2763e204565986eb6d540233
This commit is contained in:
@@ -1,5 +1,6 @@
|
||||
use super::feedback_thread_index::FeedbackThreadIndex;
|
||||
use super::*;
|
||||
use crate::error_code::OVERLOADED_ERROR_CODE;
|
||||
use codex_connectors::ConnectorDirectoryCacheContext;
|
||||
use codex_connectors::ConnectorDirectoryCacheKey;
|
||||
use codex_connectors::connector_runtime_cache_path;
|
||||
@@ -11,6 +12,7 @@ use codex_feedback::guardian_review_failures;
|
||||
use codex_rollout::RolloutRecorder;
|
||||
use sha2::Digest;
|
||||
use sha2::Sha256;
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct FeedbackRequestProcessor {
|
||||
@@ -20,6 +22,7 @@ pub(crate) struct FeedbackRequestProcessor {
|
||||
feedback: CodexFeedback,
|
||||
log_db: Option<LogDbLayer>,
|
||||
state_db: Option<StateDbHandle>,
|
||||
uploads: Arc<Semaphore>,
|
||||
}
|
||||
|
||||
impl FeedbackRequestProcessor {
|
||||
@@ -38,6 +41,7 @@ impl FeedbackRequestProcessor {
|
||||
feedback,
|
||||
log_db,
|
||||
state_db,
|
||||
uploads: Arc::new(Semaphore::new(/*permits*/ 3)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -59,6 +63,17 @@ impl FeedbackRequestProcessor {
|
||||
"sending feedback is disabled by configuration",
|
||||
));
|
||||
}
|
||||
let permit = self
|
||||
.uploads
|
||||
.clone()
|
||||
.try_acquire_owned()
|
||||
.map_err(|_| JSONRPCErrorError {
|
||||
code: OVERLOADED_ERROR_CODE,
|
||||
message:
|
||||
"Three feedback uploads are already in progress; try again after one finishes"
|
||||
.to_string(),
|
||||
data: None,
|
||||
})?;
|
||||
|
||||
let FeedbackUploadParams {
|
||||
classification,
|
||||
@@ -269,6 +284,8 @@ impl FeedbackRequestProcessor {
|
||||
let runtime_handle = tokio::runtime::Handle::current();
|
||||
|
||||
let upload_result = tokio::task::spawn_blocking(move || {
|
||||
// Cancelling the RPC waiter must not release a still-running upload's slot.
|
||||
let _permit = permit;
|
||||
let tags = (!upload_tags.is_empty()).then_some(&upload_tags);
|
||||
runtime_handle.block_on(snapshot.upload_feedback(
|
||||
FeedbackUploadOptions {
|
||||
|
||||
@@ -5,21 +5,16 @@ use app_test_support::TestAppServer;
|
||||
use codex_app_server_protocol::RequestId;
|
||||
use pretty_assertions::assert_eq;
|
||||
use serde_json::json;
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use tokio::io::BufReader;
|
||||
use tokio::net::TcpListener;
|
||||
use tokio::time::timeout;
|
||||
use wiremock::Mock;
|
||||
use wiremock::MockServer;
|
||||
use wiremock::ResponseTemplate;
|
||||
use wiremock::matchers::method;
|
||||
|
||||
#[tokio::test]
|
||||
async fn feedback_upload_reports_transport_failure_as_json_rpc_error() -> Result<()> {
|
||||
let proxy = MockServer::start().await;
|
||||
Mock::given(method("CONNECT"))
|
||||
.respond_with(ResponseTemplate::new(503))
|
||||
.expect(1)
|
||||
.mount(&proxy)
|
||||
.await;
|
||||
let proxy_uri = proxy.uri();
|
||||
async fn feedback_upload_limits_concurrency_and_releases_failed_uploads() -> Result<()> {
|
||||
let proxy = TcpListener::bind("127.0.0.1:0").await?;
|
||||
let proxy_uri = format!("http://{}", proxy.local_addr()?);
|
||||
let mut app_server = TestAppServer::builder()
|
||||
.with_env_overrides(&[
|
||||
("HTTPS_PROXY", Some(proxy_uri.as_str())),
|
||||
@@ -30,19 +25,68 @@ async fn feedback_upload_reports_transport_failure_as_json_rpc_error() -> Result
|
||||
.build_initialized()
|
||||
.await?;
|
||||
|
||||
let request_id = app_server
|
||||
let mut pending = Vec::new();
|
||||
for _ in 0..3 {
|
||||
let request_id = app_server
|
||||
.send_raw_request(
|
||||
"feedback/upload",
|
||||
Some(json!({ "classification": "bug", "includeLogs": false })),
|
||||
)
|
||||
.await?;
|
||||
let (stream, _) = timeout(Duration::from_secs(/*secs*/ 15), proxy.accept()).await??;
|
||||
let mut stream = BufReader::new(stream);
|
||||
let mut request = String::new();
|
||||
timeout(
|
||||
Duration::from_secs(/*secs*/ 15),
|
||||
stream.read_line(&mut request),
|
||||
)
|
||||
.await??;
|
||||
assert!(request.starts_with("CONNECT "));
|
||||
pending.push((request_id, stream.into_inner()));
|
||||
}
|
||||
|
||||
let excess_id = app_server
|
||||
.send_raw_request(
|
||||
"feedback/upload",
|
||||
Some(json!({ "classification": "bug", "includeLogs": false })),
|
||||
)
|
||||
.await?;
|
||||
let error = timeout(
|
||||
Duration::from_secs(15),
|
||||
app_server.read_stream_until_error_message(RequestId::Integer(request_id)),
|
||||
Duration::from_secs(/*secs*/ 15),
|
||||
app_server.read_stream_until_error_message(RequestId::Integer(excess_id)),
|
||||
)
|
||||
.await??;
|
||||
|
||||
assert_eq!(error.error.code, -32001);
|
||||
|
||||
let unavailable =
|
||||
b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n";
|
||||
for (request_id, mut stream) in pending {
|
||||
stream.write_all(unavailable).await?;
|
||||
stream.shutdown().await?;
|
||||
let error = timeout(
|
||||
Duration::from_secs(/*secs*/ 15),
|
||||
app_server.read_stream_until_error_message(RequestId::Integer(request_id)),
|
||||
)
|
||||
.await??;
|
||||
assert_eq!(error.error.code, -32603);
|
||||
assert!(error.error.message.contains("failed to upload feedback"));
|
||||
}
|
||||
|
||||
let request_id = app_server
|
||||
.send_raw_request(
|
||||
"feedback/upload",
|
||||
Some(json!({ "classification": "bug", "includeLogs": false })),
|
||||
)
|
||||
.await?;
|
||||
let (mut stream, _) = timeout(Duration::from_secs(/*secs*/ 15), proxy.accept()).await??;
|
||||
stream.write_all(unavailable).await?;
|
||||
stream.shutdown().await?;
|
||||
let error = timeout(
|
||||
Duration::from_secs(/*secs*/ 15),
|
||||
app_server.read_stream_until_error_message(RequestId::Integer(request_id)),
|
||||
)
|
||||
.await??;
|
||||
assert_eq!(error.error.code, -32603);
|
||||
assert!(error.error.message.contains("failed to upload feedback"));
|
||||
Ok(())
|
||||
}
|
||||
|
||||
+122
-11
@@ -56,7 +56,7 @@ pub const WINDOWS_SANDBOX_LOG_ATTACHMENT_FILENAME: &str = "windows-sandbox.log";
|
||||
const DEFAULT_MAX_BYTES: usize = 4 * 1024 * 1024; // 4 MiB
|
||||
const SENTRY_DSN: &str =
|
||||
"https://ae32ed50620d7a7792c1ce5df38b3e3e@o33249.ingest.us.sentry.io/4510195390611458";
|
||||
const UPLOAD_TIMEOUT_SECS: u64 = 10;
|
||||
const UPLOAD_TIMEOUT: Duration = Duration::from_secs(/*secs*/ 300);
|
||||
// Raw collection budgets used by the report API, not the interactive upload.
|
||||
pub const MAX_ATTACHMENT_BYTES: usize = 64 * 1024 * 1024;
|
||||
pub const MAX_ATTACHMENTS_BYTES: usize = 126 * 1024 * 1024;
|
||||
@@ -553,8 +553,13 @@ impl FeedbackSnapshot {
|
||||
options: FeedbackUploadOptions<'_>,
|
||||
http_client_factory: &HttpClientFactory,
|
||||
) -> Result<()> {
|
||||
self.upload_feedback_with_dsn(options, http_client_factory, SENTRY_DSN)
|
||||
.await
|
||||
self.upload_feedback_with_dsn(
|
||||
options,
|
||||
http_client_factory,
|
||||
SENTRY_DSN,
|
||||
Instant::now() + UPLOAD_TIMEOUT,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn upload_feedback_with_dsn(
|
||||
@@ -562,6 +567,7 @@ impl FeedbackSnapshot {
|
||||
options: FeedbackUploadOptions<'_>,
|
||||
http_client_factory: &HttpClientFactory,
|
||||
dsn: &str,
|
||||
deadline: Instant,
|
||||
) -> Result<()> {
|
||||
use std::str::FromStr;
|
||||
|
||||
@@ -592,15 +598,15 @@ impl FeedbackSnapshot {
|
||||
http_client_factory.clone(),
|
||||
ClientRouteClass::Other,
|
||||
);
|
||||
// Acknowledge the report before reading diagnostic files. Keep one network wait
|
||||
// budget for the whole submission; file reads and compression do not consume it.
|
||||
let mut remaining_upload_time = Duration::from_secs(UPLOAD_TIMEOUT_SECS);
|
||||
// Accept the report before reading diagnostics; all envelopes share one deadline.
|
||||
let mut rate_limited = false;
|
||||
let status = upload::send_gzip_envelope(
|
||||
&client_pool,
|
||||
&dsn,
|
||||
event_body,
|
||||
upload::EnvelopeKind::Event,
|
||||
&mut remaining_upload_time,
|
||||
deadline,
|
||||
&mut rate_limited,
|
||||
)
|
||||
.await?;
|
||||
anyhow::ensure!(
|
||||
@@ -608,7 +614,7 @@ impl FeedbackSnapshot {
|
||||
"Sentry rejected feedback upload with HTTP status {status}"
|
||||
);
|
||||
|
||||
let attachments = self.feedback_attachments(
|
||||
let mut attachments = self.feedback_attachments(
|
||||
options.include_logs,
|
||||
options.extra_attachments,
|
||||
options.extra_attachment_paths,
|
||||
@@ -617,8 +623,16 @@ impl FeedbackSnapshot {
|
||||
let mut uploaded_attachments = 0;
|
||||
let mut attachments_failed = false;
|
||||
// Keep attachments linked to the accepted event, without replaying its contents.
|
||||
for attachment in attachments {
|
||||
if remaining_upload_time.is_zero() {
|
||||
// Inspect the upper bound without reading the next file after the deadline.
|
||||
while attachments.size_hint().1 != Some(0) {
|
||||
if rate_limited || Instant::now() >= deadline {
|
||||
attachments_failed = true;
|
||||
break;
|
||||
}
|
||||
let Some(attachment) = attachments.next() else {
|
||||
break;
|
||||
};
|
||||
if Instant::now() >= deadline {
|
||||
attachments_failed = true;
|
||||
break;
|
||||
}
|
||||
@@ -630,7 +644,8 @@ impl FeedbackSnapshot {
|
||||
&dsn,
|
||||
body,
|
||||
upload::EnvelopeKind::Attachment,
|
||||
&mut remaining_upload_time,
|
||||
deadline,
|
||||
&mut rate_limited,
|
||||
)
|
||||
.await?;
|
||||
status = Some(response_status.as_u16());
|
||||
@@ -855,6 +870,7 @@ mod tests {
|
||||
use std::fs;
|
||||
use std::sync::atomic::AtomicUsize;
|
||||
use std::sync::atomic::Ordering;
|
||||
use std::time::Duration;
|
||||
|
||||
use super::*;
|
||||
use crate::FeedbackDiagnostic;
|
||||
@@ -967,10 +983,94 @@ mod tests {
|
||||
},
|
||||
&HttpClientFactory::new(OutboundProxyPolicy::ReqwestDefault),
|
||||
dsn,
|
||||
Instant::now() + UPLOAD_TIMEOUT,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn feedback_upload_allows_slow_reports() {
|
||||
let server = MockServer::start().await;
|
||||
Mock::given(method("POST"))
|
||||
.and(path("/api/42/envelope/"))
|
||||
.respond_with(
|
||||
ResponseTemplate::new(StatusCode::OK)
|
||||
.set_delay(Duration::from_secs(/*secs*/ 4)),
|
||||
)
|
||||
.expect(/*r*/ 3)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
let dsn = format!("http://public@{}/42", server.address());
|
||||
let attachment = FeedbackAttachment {
|
||||
filename: "later.txt".to_string(),
|
||||
content_type: None,
|
||||
buffer: b"later diagnostic".to_vec(),
|
||||
};
|
||||
upload_test_feedback(&CodexFeedback::new(), &dsn, &[attachment])
|
||||
.await
|
||||
.expect("all three envelopes should finish across twelve seconds of network waits");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn feedback_upload_deadline_stops_retries_and_later_attachments() {
|
||||
let server = MockServer::start().await;
|
||||
let attempt = AtomicUsize::default();
|
||||
Mock::given(method("POST"))
|
||||
.and(path("/api/42/envelope/"))
|
||||
.respond_with(move |_: &wiremock::Request| {
|
||||
match attempt.fetch_add(/*val*/ 1, Ordering::SeqCst) {
|
||||
0 | 2 => ResponseTemplate::new(StatusCode::OK)
|
||||
.set_delay(Duration::from_secs(/*secs*/ 1)),
|
||||
1 => ResponseTemplate::new(StatusCode::SERVICE_UNAVAILABLE)
|
||||
.set_delay(Duration::from_secs(/*secs*/ 1)),
|
||||
_ => ResponseTemplate::new(StatusCode::OK)
|
||||
.set_delay(Duration::from_secs(/*secs*/ 20)),
|
||||
}
|
||||
})
|
||||
.expect(/*r*/ 4)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
let dsn = format!("http://public@{}/42", server.address());
|
||||
let attachments = ["stalled.txt", "later.txt"].map(|filename| FeedbackAttachment {
|
||||
filename: filename.to_string(),
|
||||
content_type: None,
|
||||
buffer: filename.as_bytes().to_vec(),
|
||||
});
|
||||
let snapshot = CodexFeedback::new()
|
||||
.snapshot(/*session_id*/ None)
|
||||
.with_feedback_diagnostics(FeedbackDiagnostics::default());
|
||||
let error = tokio::time::timeout(
|
||||
Duration::from_secs(/*secs*/ 10),
|
||||
snapshot.upload_feedback_with_dsn(
|
||||
FeedbackUploadOptions {
|
||||
classification: "bug",
|
||||
reason: None,
|
||||
tags: None,
|
||||
include_logs: true,
|
||||
extra_attachments: &attachments,
|
||||
extra_attachment_paths: &[],
|
||||
session_source: Some(SessionSource::Cli),
|
||||
logs_override: Some(b"log contents".to_vec()),
|
||||
},
|
||||
&HttpClientFactory::new(OutboundProxyPolicy::ReqwestDefault),
|
||||
&dsn,
|
||||
Instant::now() + Duration::from_secs(/*secs*/ 8),
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("each request must use only the time left in the report deadline")
|
||||
.expect_err("the accepted report must report incomplete attachments");
|
||||
assert_eq!(
|
||||
error.to_string(),
|
||||
"feedback report was accepted, but some attachments failed to upload"
|
||||
);
|
||||
let requests = server.received_requests().await.unwrap();
|
||||
assert_eq!(requests.len(), 4);
|
||||
assert_eq!(requests[1].body, requests[2].body);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn feedback_upload_retries_diagnostics_without_replaying_core() {
|
||||
let server = MockServer::start().await;
|
||||
@@ -1114,6 +1214,7 @@ mod tests {
|
||||
},
|
||||
&HttpClientFactory::new(OutboundProxyPolicy::ReqwestDefault),
|
||||
&format!("http://public@{}/42", server.address()),
|
||||
Instant::now() + UPLOAD_TIMEOUT,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -1198,6 +1299,16 @@ mod tests {
|
||||
0,
|
||||
ResponseTemplate::new(StatusCode::OK).insert_header("Retry-After", "60"),
|
||||
),
|
||||
(
|
||||
1,
|
||||
ResponseTemplate::new(StatusCode::SERVICE_UNAVAILABLE)
|
||||
.insert_header("Retry-After", "300"),
|
||||
),
|
||||
(
|
||||
1,
|
||||
ResponseTemplate::new(StatusCode::SERVICE_UNAVAILABLE)
|
||||
.insert_header("Retry-After", "invalid"),
|
||||
),
|
||||
] {
|
||||
let server = MockServer::start().await;
|
||||
let attempt = AtomicUsize::default();
|
||||
|
||||
@@ -46,24 +46,21 @@ pub(super) fn gzip_envelope_request(
|
||||
.timeout(timeout)
|
||||
}
|
||||
|
||||
/// Send an already-serialized envelope, charging requests and backoff to the shared budget.
|
||||
/// Send an already-serialized envelope within the report's shared deadline.
|
||||
pub(super) async fn send_gzip_envelope(
|
||||
client_pool: &RouteAwareClientPool,
|
||||
dsn: &Dsn,
|
||||
body: Vec<u8>,
|
||||
kind: EnvelopeKind,
|
||||
remaining_upload_time: &mut Duration,
|
||||
deadline: Instant,
|
||||
rate_limited: &mut bool,
|
||||
) -> Result<StatusCode> {
|
||||
let request_started_at = Instant::now();
|
||||
let body = Bytes::from(body);
|
||||
// Retain the exact gzip bytes for at most two diagnostic retries. Never replay the core.
|
||||
let mut retry_delays = [250, 500].map(Duration::from_millis).into_iter();
|
||||
let response = loop {
|
||||
let timeout = remaining_upload_time.saturating_sub(request_started_at.elapsed());
|
||||
if timeout.is_zero() {
|
||||
*remaining_upload_time = Duration::ZERO;
|
||||
anyhow::bail!("feedback upload network wait budget exhausted");
|
||||
}
|
||||
let timeout = deadline.saturating_duration_since(Instant::now());
|
||||
anyhow::ensure!(!timeout.is_zero(), "feedback upload deadline exceeded");
|
||||
let response = gzip_envelope_request(client_pool, dsn, body.clone(), timeout)
|
||||
.send()
|
||||
.await;
|
||||
@@ -77,7 +74,7 @@ pub(super) async fn send_gzip_envelope(
|
||||
}) {
|
||||
// Even accepted requests can impose a quota on events or attachments.
|
||||
// https://develop.sentry.dev/sdk/foundations/transport/rate-limiting/
|
||||
*remaining_upload_time = Duration::ZERO;
|
||||
*rate_limited = true;
|
||||
break response;
|
||||
}
|
||||
let retryable = match &response {
|
||||
@@ -96,9 +93,8 @@ pub(super) async fn send_gzip_envelope(
|
||||
response.status() == StatusCode::TOO_MANY_REQUESTS
|
||||
&& (retry_delay.is_none() || !response.headers().contains_key("Retry-After"))
|
||||
}) {
|
||||
// Bare 429s default to a cooldown longer than this upload budget.
|
||||
// Never let later files bypass an exhausted rate limit.
|
||||
*remaining_upload_time = Duration::ZERO;
|
||||
*rate_limited = true;
|
||||
break response;
|
||||
}
|
||||
let Some(retry_delay) = retry_delay else {
|
||||
@@ -111,20 +107,16 @@ pub(super) async fn send_gzip_envelope(
|
||||
.map(|value| {
|
||||
// An invalid cooldown is not permission to send immediately.
|
||||
parse_retry_after(value.to_str().unwrap_or_default())
|
||||
.unwrap_or(*remaining_upload_time)
|
||||
.unwrap_or_else(|| deadline.saturating_duration_since(Instant::now()))
|
||||
});
|
||||
let delay = retry_delay.max(retry_after.unwrap_or_default());
|
||||
if delay >= remaining_upload_time.saturating_sub(request_started_at.elapsed()) {
|
||||
// The cooldown applies to later files too, not just this retry.
|
||||
*remaining_upload_time = Duration::ZERO;
|
||||
if delay >= deadline.saturating_duration_since(Instant::now()) {
|
||||
// Leave long cooldowns for a later submission, including its remaining files.
|
||||
*rate_limited = true;
|
||||
break response;
|
||||
}
|
||||
tokio::time::sleep(delay).await;
|
||||
if request_started_at.elapsed() >= *remaining_upload_time {
|
||||
break response;
|
||||
}
|
||||
};
|
||||
*remaining_upload_time = remaining_upload_time.saturating_sub(request_started_at.elapsed());
|
||||
response
|
||||
.map(|response| response.status())
|
||||
.context("failed to upload feedback to Sentry")
|
||||
|
||||
Reference in New Issue
Block a user