修复 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 }
[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")
.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<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 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<u8
break;
}
}
body
OssErrorBodyRead { body, read_failure }
}
fn request_status_error_from_oss_parts(
@@ -1240,6 +1264,7 @@ fn request_status_error_from_oss_parts(
status: u16,
header_request_id: Option<String>,
body: &[u8],
body_read_failure: Option<OssErrorBodyReadFailure>,
) -> 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#"<Error>
<Code>RequestTimeout</Code><RequestId>xml-request-id</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 {
panic!("OSS status failure should remain a request error");
};
@@ -2269,8 +2311,13 @@ mod tests {
b"<Error><Code>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"<Error><Code>RequestTimeout</Code></Error>",
None,
);
assert!(!oss_error_is_retryable(&error));
}
@@ -2289,7 +2337,8 @@ mod tests {
body.extend_from_slice(
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 {
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"<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]
async fn reqwest_builder_error_is_not_retryable_transport() {
let error = reqwest::Client::new()