修复 400 错误响应体读取失败被误判为确定性失败

PUT 收到 400 后若响应体在解析出 OSS 错误码前超时/断流,读取错误
不再被静默吞掉:保留 status=400 并按 timeout/transport 归类重试;
已解析出确定性错误码时维持原有语义。补充单元测试与真实 HTTP
断流/挂起回归测试。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
2026-07-20 10:11:06 +00:00
parent 0932b2f4fb
commit b8e0a20ddd
2 changed files with 215 additions and 14 deletions
+1 -1
View File
@@ -17,4 +17,4 @@ tokio = { workspace = true, features = ["sync", "time"] }
tracing = { workspace = true } tracing = { workspace = true }
[dev-dependencies] [dev-dependencies]
tokio = { workspace = true, features = ["macros", "rt"] } tokio = { workspace = true, features = ["macros", "rt", "net", "io-util"] }
+214 -13
View File
@@ -1205,26 +1205,50 @@ async fn request_status_error_from_oss_put_response(mut response: reqwest::Respo
.get("x-oss-request-id") .get("x-oss-request-id")
.and_then(|value| value.to_str().ok()) .and_then(|value| value.to_str().ok())
.and_then(|value| normalize_oss_error_field(value, OSS_REQUEST_ID_MAX_BYTES)); .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 read_bounded_oss_error_body(&mut response).await
} else { } else {
Vec::new() OssErrorBodyRead {
body: Vec::new(),
read_failure: None,
}
}; };
request_status_error_from_oss_parts( request_status_error_from_oss_parts(
OssRequestOperation::Put, OssRequestOperation::Put,
status.as_u16(), status.as_u16(),
header_request_id, header_request_id,
&body, &body_read.body,
body_read.read_failure,
) )
} }
async fn read_bounded_oss_error_body(response: &mut reqwest::Response) -> Vec<u8> { /// 400 错误响应体读取失败的原因。仅在部分响应体尚未解析出确定性
/// OSS 错误码时参与重试判定,否则只体现在 message 里。
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct OssErrorBodyReadFailure {
timeout: bool,
}
struct OssErrorBodyRead {
body: Vec<u8>,
read_failure: Option<OssErrorBodyReadFailure>,
}
async fn read_bounded_oss_error_body(response: &mut reqwest::Response) -> OssErrorBodyRead {
let mut body = Vec::new(); let mut body = Vec::new();
let mut read_failure = None;
while body.len() < CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES { while body.len() < CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES {
let chunk = match response.chunk().await { let chunk = match response.chunk().await {
Ok(Some(chunk)) => chunk, 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(); let remaining = CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES - body.len();
body.extend_from_slice(&chunk[..chunk.len().min(remaining)]); 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<u8
break; break;
} }
} }
body OssErrorBodyRead { body, read_failure }
} }
fn request_status_error_from_oss_parts( fn request_status_error_from_oss_parts(
@@ -1240,6 +1264,7 @@ fn request_status_error_from_oss_parts(
status: u16, status: u16,
header_request_id: Option<String>, header_request_id: Option<String>,
body: &[u8], body: &[u8],
body_read_failure: Option<OssErrorBodyReadFailure>,
) -> OssError { ) -> OssError {
let header_request_id = header_request_id let header_request_id = header_request_id
.as_deref() .as_deref()
@@ -1249,8 +1274,17 @@ fn request_status_error_from_oss_parts(
let xml_request_id = let xml_request_id =
extract_oss_error_xml_field(bounded_body, "RequestId", OSS_REQUEST_ID_MAX_BYTES); 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 oss_request_id = header_request_id.or(xml_request_id);
let timeout = status == reqwest::StatusCode::BAD_REQUEST.as_u16() // 确定性错误码优先:已解析出 Code 时,响应体读取失败只保留在 message 里,
&& oss_code.as_deref() == Some("RequestTimeout"); // 不改变重试语义;错误码缺失时才按读取失败归类为可重试的超时/传输错误。
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}"); let mut message = format!("OSS PutObject 失败,状态码:{status}");
if let Some(oss_code) = oss_code.as_deref() { if let Some(oss_code) = oss_code.as_deref() {
message.push_str(&format!("OSS 错误码:{oss_code}")); 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() { if let Some(oss_request_id) = oss_request_id.as_deref() {
message.push_str(&format!("OSS Request ID{oss_request_id}")); message.push_str(&format!("OSS Request ID{oss_request_id}"));
} }
if body_read_failure.is_some() {
message.push_str(",错误响应体读取失败");
}
OssError::Request(OssRequestError { OssError::Request(OssRequestError {
status: Some(status), status: Some(status),
timeout, timeout,
connect: false, connect: false,
transport: false, transport,
oss_code, oss_code,
oss_request_id, oss_request_id,
operation, operation,
@@ -1380,6 +1417,9 @@ fn oss_error_is_retryable(error: &OssError) -> bool {
match (request_error.status, request_error.oss_code.as_deref()) { match (request_error.status, request_error.oss_code.as_deref()) {
(Some(400), Some("RequestTimeout")) => true, (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(408 | 429 | 500..=599), _) => true,
(Some(_), _) => false, (Some(_), _) => false,
(None, _) => request_error.timeout || request_error.connect || request_error.transport, (None, _) => request_error.timeout || request_error.connect || request_error.transport,
@@ -2230,6 +2270,7 @@ mod tests {
400, 400,
Some("header-request-id".to_string()), Some("header-request-id".to_string()),
body, body,
None,
); );
let OssError::Request(request_error) = &error else { let OssError::Request(request_error) = &error else {
panic!("OSS status failure should remain a request error"); panic!("OSS status failure should remain a request error");
@@ -2250,7 +2291,8 @@ mod tests {
let body = br#"<Error> let body = br#"<Error>
<Code>RequestTimeout</Code><RequestId>xml-request-id</RequestId> <Code>RequestTimeout</Code><RequestId>xml-request-id</RequestId>
</Error>"#; </Error>"#;
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 { let OssError::Request(request_error) = &error else {
panic!("OSS status failure should remain a request error"); panic!("OSS status failure should remain a request error");
}; };
@@ -2269,8 +2311,13 @@ mod tests {
b"<Error><Code>RequestTimeout".as_slice(), b"<Error><Code>RequestTimeout".as_slice(),
b"not xml".as_slice(), b"not xml".as_slice(),
] { ] {
let error = let error = request_status_error_from_oss_parts(
request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, body); OssRequestOperation::Put,
400,
None,
body,
None,
);
assert!(!oss_error_is_retryable(&error)); assert!(!oss_error_is_retryable(&error));
} }
@@ -2279,6 +2326,7 @@ mod tests {
403, 403,
None, None,
b"<Error><Code>RequestTimeout</Code></Error>", b"<Error><Code>RequestTimeout</Code></Error>",
None,
); );
assert!(!oss_error_is_retryable(&error)); assert!(!oss_error_is_retryable(&error));
} }
@@ -2289,7 +2337,8 @@ mod tests {
body.extend_from_slice( body.extend_from_slice(
b"<Error><Code>RequestTimeout</Code><RequestId>late</RequestId></Error>", b"<Error><Code>RequestTimeout</Code><RequestId>late</RequestId></Error>",
); );
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 { let OssError::Request(request_error) = &error else {
panic!("OSS status failure should remain a request error"); panic!("OSS status failure should remain a request error");
}; };
@@ -2300,6 +2349,158 @@ mod tests {
assert!(!oss_error_is_retryable(&error)); 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"<Error><Code>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"<Error><Code>InvalidArgument</Code><RequestId>partial",
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"<Error><Code>RequestTimeout</Code>",
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\
<Error><Code>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] #[tokio::test]
async fn reqwest_builder_error_is_not_retryable_transport() { async fn reqwest_builder_error_is_not_retryable_transport() {
let error = reqwest::Client::new() let error = reqwest::Client::new()