From b8e0a20ddd2bece812e6312b95c6de211e95f58c Mon Sep 17 00:00:00 2001 From: Linghong Date: Mon, 20 Jul 2026 10:11:06 +0000 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=20400=20=E9=94=99=E8=AF=AF?= =?UTF-8?q?=E5=93=8D=E5=BA=94=E4=BD=93=E8=AF=BB=E5=8F=96=E5=A4=B1=E8=B4=A5?= =?UTF-8?q?=E8=A2=AB=E8=AF=AF=E5=88=A4=E4=B8=BA=E7=A1=AE=E5=AE=9A=E6=80=A7?= =?UTF-8?q?=E5=A4=B1=E8=B4=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PUT 收到 400 后若响应体在解析出 OSS 错误码前超时/断流,读取错误 不再被静默吞掉:保留 status=400 并按 timeout/transport 归类重试; 已解析出确定性错误码时维持原有语义。补充单元测试与真实 HTTP 断流/挂起回归测试。 Co-Authored-By: Claude Fable 5 --- server-rs/crates/platform-oss/Cargo.toml | 2 +- server-rs/crates/platform-oss/src/lib.rs | 227 +++++++++++++++++++++-- 2 files changed, 215 insertions(+), 14 deletions(-) diff --git a/server-rs/crates/platform-oss/Cargo.toml b/server-rs/crates/platform-oss/Cargo.toml index e06e28dc4..e514611cd 100644 --- a/server-rs/crates/platform-oss/Cargo.toml +++ b/server-rs/crates/platform-oss/Cargo.toml @@ -17,4 +17,4 @@ tokio = { workspace = true, features = ["sync", "time"] } tracing = { workspace = true } [dev-dependencies] -tokio = { workspace = true, features = ["macros", "rt"] } +tokio = { workspace = true, features = ["macros", "rt", "net", "io-util"] } diff --git a/server-rs/crates/platform-oss/src/lib.rs b/server-rs/crates/platform-oss/src/lib.rs index 38170861c..c563ee158 100644 --- a/server-rs/crates/platform-oss/src/lib.rs +++ b/server-rs/crates/platform-oss/src/lib.rs @@ -1205,26 +1205,50 @@ async fn request_status_error_from_oss_put_response(mut response: reqwest::Respo .get("x-oss-request-id") .and_then(|value| value.to_str().ok()) .and_then(|value| normalize_oss_error_field(value, OSS_REQUEST_ID_MAX_BYTES)); - let body = if status == reqwest::StatusCode::BAD_REQUEST { + let body_read = if status == reqwest::StatusCode::BAD_REQUEST { read_bounded_oss_error_body(&mut response).await } else { - Vec::new() + OssErrorBodyRead { + body: Vec::new(), + read_failure: None, + } }; request_status_error_from_oss_parts( OssRequestOperation::Put, status.as_u16(), header_request_id, - &body, + &body_read.body, + body_read.read_failure, ) } -async fn read_bounded_oss_error_body(response: &mut reqwest::Response) -> Vec { +/// 400 错误响应体读取失败的原因。仅在部分响应体尚未解析出确定性 +/// OSS 错误码时参与重试判定,否则只体现在 message 里。 +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct OssErrorBodyReadFailure { + timeout: bool, +} + +struct OssErrorBodyRead { + body: Vec, + read_failure: Option, +} + +async fn read_bounded_oss_error_body(response: &mut reqwest::Response) -> OssErrorBodyRead { let mut body = Vec::new(); + let mut read_failure = None; while body.len() < CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES { let chunk = match response.chunk().await { Ok(Some(chunk)) => chunk, - Ok(None) | Err(_) => break, + Ok(None) => break, + Err(error) => { + // 断流/超时不丢弃已读字节:部分响应体可能已含确定性错误码。 + read_failure = Some(OssErrorBodyReadFailure { + timeout: error.is_timeout(), + }); + break; + } }; let remaining = CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES - body.len(); body.extend_from_slice(&chunk[..chunk.len().min(remaining)]); @@ -1232,7 +1256,7 @@ async fn read_bounded_oss_error_body(response: &mut reqwest::Response) -> Vec, body: &[u8], + body_read_failure: Option, ) -> OssError { let header_request_id = header_request_id .as_deref() @@ -1249,8 +1274,17 @@ fn request_status_error_from_oss_parts( let xml_request_id = extract_oss_error_xml_field(bounded_body, "RequestId", OSS_REQUEST_ID_MAX_BYTES); let oss_request_id = header_request_id.or(xml_request_id); - let timeout = status == reqwest::StatusCode::BAD_REQUEST.as_u16() - && oss_code.as_deref() == Some("RequestTimeout"); + // 确定性错误码优先:已解析出 Code 时,响应体读取失败只保留在 message 里, + // 不改变重试语义;错误码缺失时才按读取失败归类为可重试的超时/传输错误。 + let unclassified_read_failure = if oss_code.is_some() { + None + } else { + body_read_failure + }; + let timeout = (status == reqwest::StatusCode::BAD_REQUEST.as_u16() + && oss_code.as_deref() == Some("RequestTimeout")) + || unclassified_read_failure.is_some_and(|failure| failure.timeout); + let transport = unclassified_read_failure.is_some_and(|failure| !failure.timeout); let mut message = format!("OSS PutObject 失败,状态码:{status}"); if let Some(oss_code) = oss_code.as_deref() { message.push_str(&format!(",OSS 错误码:{oss_code}")); @@ -1258,12 +1292,15 @@ fn request_status_error_from_oss_parts( if let Some(oss_request_id) = oss_request_id.as_deref() { message.push_str(&format!(",OSS Request ID:{oss_request_id}")); } + if body_read_failure.is_some() { + message.push_str(",错误响应体读取失败"); + } OssError::Request(OssRequestError { status: Some(status), timeout, connect: false, - transport: false, + transport, oss_code, oss_request_id, operation, @@ -1380,6 +1417,9 @@ fn oss_error_is_retryable(error: &OssError) -> bool { match (request_error.status, request_error.oss_code.as_deref()) { (Some(400), Some("RequestTimeout")) => true, + // 400 是唯一会读取错误响应体的状态码:错误码缺失且响应体读取 + // 超时/断流时,无法证明是确定性 400,按传输错误重试。 + (Some(400), None) if request_error.timeout || request_error.transport => true, (Some(408 | 429 | 500..=599), _) => true, (Some(_), _) => false, (None, _) => request_error.timeout || request_error.connect || request_error.transport, @@ -2230,6 +2270,7 @@ mod tests { 400, Some("header-request-id".to_string()), body, + None, ); let OssError::Request(request_error) = &error else { panic!("OSS status failure should remain a request error"); @@ -2250,7 +2291,8 @@ mod tests { let body = br#" RequestTimeoutxml-request-id "#; - let error = request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, body); + let error = + request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, body, None); let OssError::Request(request_error) = &error else { panic!("OSS status failure should remain a request error"); }; @@ -2269,8 +2311,13 @@ mod tests { b"RequestTimeout".as_slice(), b"not xml".as_slice(), ] { - let error = - request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, body); + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + body, + None, + ); assert!(!oss_error_is_retryable(&error)); } @@ -2279,6 +2326,7 @@ mod tests { 403, None, b"RequestTimeout", + None, ); assert!(!oss_error_is_retryable(&error)); } @@ -2289,7 +2337,8 @@ mod tests { body.extend_from_slice( b"RequestTimeoutlate", ); - let error = request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, &body); + let error = + request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, &body, None); let OssError::Request(request_error) = &error else { panic!("OSS status failure should remain a request error"); }; @@ -2300,6 +2349,158 @@ mod tests { assert!(!oss_error_is_retryable(&error)); } + #[test] + fn oss_400_without_code_and_broken_body_read_is_retryable() { + for (read_timeout, expect_transport) in [(true, false), (false, true)] { + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + b"Request", + Some(OssErrorBodyReadFailure { + timeout: read_timeout, + }), + ); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert_eq!(request_error.oss_code, None); + assert_eq!(request_error.timeout, read_timeout); + assert_eq!(request_error.transport, expect_transport); + assert!(request_error.message.contains("错误响应体读取失败")); + assert!(oss_error_is_retryable(&error)); + } + } + + #[test] + fn oss_400_with_parsed_code_keeps_deterministic_semantics_on_read_failure() { + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + b"InvalidArgumentpartial", + Some(OssErrorBodyReadFailure { timeout: true }), + ); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.oss_code.as_deref(), Some("InvalidArgument")); + assert!(!request_error.timeout); + assert!(!request_error.transport); + assert!(request_error.message.contains("错误响应体读取失败")); + assert!(!oss_error_is_retryable(&error)); + + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + b"RequestTimeout", + Some(OssErrorBodyReadFailure { timeout: false }), + ); + assert!(oss_error_is_retryable(&error)); + } + + const MOCK_PUT_BODY: &[u8] = b"animation-frame-bytes"; + + /// 极简 HTTP/1.1 mock:读完整个 PUT 请求后返回 400 与部分 XML 响应体 + /// (Content-Length 大于实际发送字节),`stall_before_close` 决定挂住 + /// 连接触发客户端读超时,还是直接断开触发传输错误。 + async fn spawn_broken_error_body_server(stall_before_close: bool) -> std::net::SocketAddr { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("mock server should bind"); + let addr = listener + .local_addr() + .expect("mock server should expose its addr"); + tokio::spawn(async move { + let Ok((mut socket, _)) = listener.accept().await else { + return; + }; + let mut received = Vec::new(); + let mut buffer = [0u8; 4096]; + while !received.ends_with(MOCK_PUT_BODY) { + match socket.read(&mut buffer).await { + Ok(0) | Err(_) => return, + Ok(read) => received.extend_from_slice(&buffer[..read]), + } + } + let response = "HTTP/1.1 400 Bad Request\r\n\ + x-oss-request-id: mock-request-id\r\n\ + Content-Type: application/xml\r\n\ + Content-Length: 4096\r\n\ + \r\n\ + Request"; + if socket.write_all(response.as_bytes()).await.is_err() { + return; + } + let _ = socket.flush().await; + if stall_before_close { + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + } + }); + addr + } + + #[tokio::test] + async fn oss_400_with_broken_error_body_stream_is_retryable_transport() { + let addr = spawn_broken_error_body_server(false).await; + let response = reqwest::Client::new() + .put(format!("http://{addr}/generated-animations/frame01.png")) + .body(MOCK_PUT_BODY.to_vec()) + .send() + .await + .expect("response headers should arrive before the body breaks"); + assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); + + let error = request_status_error_from_oss_put_response(response).await; + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert_eq!(request_error.oss_code, None); + assert!(!request_error.timeout); + assert!(request_error.transport); + assert_eq!( + request_error.oss_request_id.as_deref(), + Some("mock-request-id") + ); + assert!(request_error.message.contains("错误响应体读取失败")); + assert!(oss_error_is_retryable(&error)); + } + + #[tokio::test] + async fn oss_400_with_stalled_error_body_stream_is_retryable_timeout() { + let addr = spawn_broken_error_body_server(true).await; + let client = reqwest::Client::builder() + .timeout(std::time::Duration::from_millis(300)) + .build() + .expect("test client should build"); + let response = client + .put(format!("http://{addr}/generated-animations/frame01.png")) + .body(MOCK_PUT_BODY.to_vec()) + .send() + .await + .expect("response headers should arrive before the body stalls"); + assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); + + let error = request_status_error_from_oss_put_response(response).await; + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert_eq!(request_error.oss_code, None); + assert!(request_error.timeout); + assert!(!request_error.transport); + assert!(oss_error_is_retryable(&error)); + } + #[tokio::test] async fn reqwest_builder_error_is_not_retryable_transport() { let error = reqwest::Client::new()