修复 400 错误响应体读取失败被误判为确定性失败
PUT 收到 400 后若响应体在解析出 OSS 错误码前超时/断流,读取错误 不再被静默吞掉:保留 status=400 并按 timeout/transport 归类重试; 已解析出确定性错误码时维持原有语义。补充单元测试与真实 HTTP 断流/挂起回归测试。 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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"] }
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user