任务状态查询本身带上退避重试,瞬时故障不再作废一次计费任务
- platform-tripo:`get_task` 内部按 `TripoSettings::retries` 对瞬时故障(连接抖动 / 超时、408、429、5xx)退避重试,确定性失败(任务 id 非法、响应结构不符、404)立即返回;重试字段与退避函数改名,不再只服务产物下载 - api-server worker:查询在重试后仍是瞬时故障时继续按轮询间隔查到预算用尽,只有确定性失败才立刻收口 - 新增用例:一次 500 后成功(重试生效)、404 只打一次请求(不重试),并给测试客户端加了可指定 base_url 的构造 - docs/technical:写明查询带退避重试与「确定性失败立即返回」的边界
This commit is contained in:
@@ -44,7 +44,7 @@ TODO:图生输入当前由 api-server 自己读站内对象的字节再上传
|
||||
|
||||
## 任务查询与下载
|
||||
|
||||
`get_task` 按 `TripoTaskHandle` 做单次查询,返回通用 `TripoTaskSnapshot`,由 `task_type` 决定 `output` 的具体 variant;adapter 不做轮询、不阻塞等待。
|
||||
`get_task` 按 `TripoTaskHandle` 做单次查询,返回通用 `TripoTaskSnapshot`,由 `task_type` 决定 `output` 的具体 variant;adapter 不做轮询、不阻塞等待。单次查询内部带退避重试:生成本来要跑几分钟,一次连接抖动不该作废一笔已经提交、不能重来的计费任务,所以瞬时故障(连接抖动 / 超时、408、429、5xx)按 `TripoSettings::retries` 重试,确定性失败(任务 id 非法、响应结构不符、404)立即返回。
|
||||
|
||||
`download_model`(及预览图的 `download_rendered_image`)接受 `TripoTaskSnapshot` 与调用方给的体积上限 `max_bytes`,先确认快照确实处于完成态,再从严格 endpoint 结果中取得已校验的产物地址,由 provider 自己的无鉴权 reqwest client 打开签名 URL,读完整个 body 后返回 `TripoArtifactBytes`(`url`、`content_type`、`content_length`、`bytes`、`filename(name)`)。流式句柄 `TripoDownloadedArtifact`(`next_chunk()`)保留为内部实现:**重试必须包住整个 body** —— 只重试「拿到响应头」这一步的时候,几十 MB 的字节其实是在之后才传输的,一次 CDN 抖动就会毁掉一次已经扣费、provider 任务也跑完的生成。因此每次尝试都重新取响应头并整体重下(不做断点续传,签名地址会过期),body 读取失败与长度不一致都归一成可重试的传输错误,重试上限与退避沿用 `TripoSettings::retries`。**可重试性由产生错误的一方显式判定,不再由「`status` 是不是 `None`」推断**:传输层的超时 / 连接失败、「无数据推进」超时、body 提前结束都可重试;客户端构造失败、响应体格式错误(空 body / 缺 `data` 字段)、DNS 与 TLS 之外的永久失败不可重试 —— 只看 `status.is_none()` 会把这两类混在一起,把故障拖到重试耗尽才暴露。`filename(name)` 的扩展名来自远端地址,只接受短的 ASCII 字母数字,其余退回 `glb`,避免远端地址里的 `%2F` 解码后拼出跨目录路径。每次读取都校验实际接收字节数与 `Content-Length`,读到 `max_bytes` 之上直接按输出违约失败(不重试):宁可失败退款也不要把 api-server 内存打满。产物下载不设总超时:几十 MB 的流只要还在出数据就不该被判失败,重下只发生在传输真的断了的时候,因此下载链路只按「无数据推进」判超时,预算取 `TripoSettings::request_timeout`(连接用 `connect_timeout`,响应体用 `read_timeout`,响应头之前由显式的 `tokio::time::timeout` 兜住),超过预算没有数据推进即按传输失败返回;SDK 保持第三方原样,不承担产物下载。smoke example 一次写入完整字节;未来接入 OSS 时应把同一数据流直接送入 OSS 分片上传,不经过完整内存缓冲。
|
||||
|
||||
|
||||
@@ -222,7 +222,17 @@ async fn poll_until_terminal(
|
||||
if Instant::now() >= provider_deadline {
|
||||
return Err(provider_deadline_error());
|
||||
}
|
||||
let snapshot = client.get_task(handle).await.map_err(map_provider_error)?;
|
||||
let snapshot = match client.get_task(handle).await {
|
||||
Ok(snapshot) => snapshot,
|
||||
// 查询本身已经带过一轮退避重试;仍是瞬时故障时继续按轮询间隔查到预算用尽,
|
||||
// 不因为一次 provider 抖动就作废这笔已经提交、不能重来的任务。
|
||||
Err(error) if error.is_retryable() => {
|
||||
let remaining = provider_deadline.saturating_duration_since(Instant::now());
|
||||
tokio::time::sleep(MODEL3D_POLL_INTERVAL.min(remaining)).await;
|
||||
continue;
|
||||
}
|
||||
Err(error) => return Err(map_provider_error(error)),
|
||||
};
|
||||
match snapshot.status {
|
||||
Model3dTaskStatus::Completed => return Ok(snapshot),
|
||||
Model3dTaskStatus::Queued | Model3dTaskStatus::Running => {}
|
||||
|
||||
@@ -12,7 +12,8 @@ use super::{
|
||||
pub struct TripoProviderClient {
|
||||
pub(crate) client: TripoClient,
|
||||
artifact_client: reqwest::Client,
|
||||
artifact_retries: u32,
|
||||
/// 瞬时故障的重试次数:产物下载与任务状态查询共用同一份配置。
|
||||
retries: u32,
|
||||
/// 产物下载的「无数据推进」预算:连接、首字节和之后每个数据块都必须在它之内有进展,
|
||||
/// 下载总时长不设限。
|
||||
artifact_stall_timeout: Duration,
|
||||
@@ -42,22 +43,38 @@ impl TripoProviderClient {
|
||||
Ok(Self {
|
||||
client: TripoClient::new(settings.client_options()).map_err(TripoError::from)?,
|
||||
artifact_client,
|
||||
artifact_retries: settings.retries,
|
||||
retries: settings.retries,
|
||||
artifact_stall_timeout: settings.request_timeout,
|
||||
})
|
||||
}
|
||||
|
||||
/// 查询任务状态。
|
||||
///
|
||||
/// 单次查询失败不直接判死:生成本来就要跑几分钟,一次连接抖动不该作废一笔不可重来的
|
||||
/// 计费任务。瞬时故障(连接抖动 / 超时、408、429、5xx)按 [`retry_backoff`] 退避重试,
|
||||
/// 与产物下载同一口径;确定性失败(任务 id 非法、响应结构不符)立刻返回,等在这里也是
|
||||
/// 同一个结果。
|
||||
pub async fn get_task(
|
||||
&self,
|
||||
handle: &TripoTaskHandle,
|
||||
) -> Result<TripoTaskSnapshot, TripoError> {
|
||||
validate_task_id(&handle.task_id)?;
|
||||
let task = self
|
||||
.client
|
||||
.get_task(&handle.task_id)
|
||||
.await
|
||||
.map_err(TripoError::from)?;
|
||||
map_task(task)
|
||||
let total_attempts = self.retries.saturating_add(1);
|
||||
let mut attempt = 1;
|
||||
loop {
|
||||
let result = match self.client.get_task(&handle.task_id).await {
|
||||
Ok(task) => map_task(task),
|
||||
Err(error) => Err(TripoError::from(error)),
|
||||
};
|
||||
match result {
|
||||
Ok(snapshot) => return Ok(snapshot),
|
||||
Err(error) if attempt < total_attempts && error.is_retryable() => {
|
||||
tokio::time::sleep(retry_backoff(attempt)).await;
|
||||
attempt += 1;
|
||||
}
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 下载完整的模型产物(响应头 + 整个 body);`max_bytes` 是单次读取的硬上限。
|
||||
@@ -98,7 +115,7 @@ impl TripoProviderClient {
|
||||
url: &TripoUrl,
|
||||
max_bytes: u64,
|
||||
) -> Result<TripoArtifactBytes, TripoError> {
|
||||
let total_attempts = self.artifact_retries.saturating_add(1);
|
||||
let total_attempts = self.retries.saturating_add(1);
|
||||
let mut attempt = 1;
|
||||
|
||||
loop {
|
||||
@@ -111,7 +128,7 @@ impl TripoProviderClient {
|
||||
match result {
|
||||
Ok(artifact) => return Ok(artifact),
|
||||
Err(error) if attempt < total_attempts && error.is_retryable() => {
|
||||
tokio::time::sleep(download_backoff(attempt)).await;
|
||||
tokio::time::sleep(retry_backoff(attempt)).await;
|
||||
attempt += 1;
|
||||
}
|
||||
Err(error) => return Err(error),
|
||||
@@ -174,7 +191,7 @@ impl TripoProviderClient {
|
||||
}
|
||||
}
|
||||
|
||||
fn download_backoff(attempt: u32) -> Duration {
|
||||
fn retry_backoff(attempt: u32) -> Duration {
|
||||
Duration::from_millis(250u64.saturating_mul(2u64.saturating_pow(attempt.min(6))))
|
||||
}
|
||||
|
||||
@@ -231,6 +248,8 @@ fn transport_error_source(error: &reqwest::Error) -> String {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU32, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
@@ -249,11 +268,16 @@ mod tests {
|
||||
test_client_with_retries(0)
|
||||
}
|
||||
|
||||
/// 带产物重试次数的测试客户端:验证「整体重下」需要至少一次重试。
|
||||
/// 带重试次数的测试客户端:验证「整体重下」需要至少一次重试。
|
||||
/// 默认 base_url 指向一个连不上的端口,避免用例意外打到真实 provider。
|
||||
fn test_client_with_retries(retries: u32) -> TripoProviderClient {
|
||||
test_client_with_base_url("http://127.0.0.1:1", retries)
|
||||
}
|
||||
|
||||
fn test_client_with_base_url(base_url: &str, retries: u32) -> TripoProviderClient {
|
||||
let settings = TripoSettings::new(
|
||||
"test-key".to_string(),
|
||||
"http://127.0.0.1:1".to_string(),
|
||||
base_url.to_string(),
|
||||
STALL_TIMEOUT,
|
||||
retries,
|
||||
"genarrative-test-tripo/1".to_string(),
|
||||
@@ -448,6 +472,101 @@ mod tests {
|
||||
assert_eq!(artifact.content_length, Some(COMPLETE_BODY.len() as u64));
|
||||
}
|
||||
|
||||
/// 第一次查询返回 500,第二次给出正常状态:瞬时故障必须重试后成功。
|
||||
#[tokio::test]
|
||||
async fn task_status_query_retries_transient_failures() {
|
||||
let listener = TcpListener::bind("127.0.0.1:0")
|
||||
.await
|
||||
.expect("mock server should bind");
|
||||
let addr = listener
|
||||
.local_addr()
|
||||
.expect("mock server should expose addr");
|
||||
tokio::spawn(async move {
|
||||
let mut first_attempt = true;
|
||||
loop {
|
||||
let Ok((mut socket, _)) = listener.accept().await else {
|
||||
return;
|
||||
};
|
||||
read_request(&mut socket).await;
|
||||
let response = if first_attempt {
|
||||
first_attempt = false;
|
||||
"HTTP/1.1 500 Internal Server Error\r\nContent-Length: 0\r\n\r\n".to_string()
|
||||
} else {
|
||||
let body = serde_json::json!({
|
||||
"code": 0,
|
||||
"data": {
|
||||
"task_id": "task-1",
|
||||
"type": "text_to_model",
|
||||
"status": "running",
|
||||
"progress": 10
|
||||
}
|
||||
})
|
||||
.to_string();
|
||||
format!(
|
||||
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\
|
||||
Content-Length: {}\r\n\r\n{body}",
|
||||
body.len()
|
||||
)
|
||||
};
|
||||
let _ = socket.write_all(response.as_bytes()).await;
|
||||
let _ = socket.flush().await;
|
||||
}
|
||||
});
|
||||
|
||||
let snapshot = test_client_with_base_url(&format!("http://{addr}"), 1)
|
||||
.get_task(&TripoTaskHandle {
|
||||
task_id: "task-1".to_string(),
|
||||
})
|
||||
.await
|
||||
.expect("瞬时 500 必须重试后拿到状态");
|
||||
assert_eq!(snapshot.status, Model3dTaskStatus::Running);
|
||||
}
|
||||
|
||||
/// 404 是确定性失败:只打一次请求就返回,不按退避重试。
|
||||
#[tokio::test]
|
||||
async fn task_status_query_does_not_retry_deterministic_failures() {
|
||||
let connections = Arc::new(AtomicU32::new(0));
|
||||
let observed = Arc::clone(&connections);
|
||||
let listener = TcpListener::bind("127.0.0.1:0")
|
||||
.await
|
||||
.expect("mock server should bind");
|
||||
let addr = listener
|
||||
.local_addr()
|
||||
.expect("mock server should expose addr");
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
let Ok((mut socket, _)) = listener.accept().await else {
|
||||
return;
|
||||
};
|
||||
observed.fetch_add(1, Ordering::SeqCst);
|
||||
read_request(&mut socket).await;
|
||||
let body =
|
||||
serde_json::json!({ "code": 1004, "message": "task not found" }).to_string();
|
||||
let response = format!(
|
||||
"HTTP/1.1 404 Not Found\r\nContent-Type: application/json\r\n\
|
||||
Content-Length: {}\r\n\r\n{body}",
|
||||
body.len()
|
||||
);
|
||||
let _ = socket.write_all(response.as_bytes()).await;
|
||||
let _ = socket.flush().await;
|
||||
}
|
||||
});
|
||||
|
||||
let error = test_client_with_base_url(&format!("http://{addr}"), 2)
|
||||
.get_task(&TripoTaskHandle {
|
||||
task_id: "task-1".to_string(),
|
||||
})
|
||||
.await
|
||||
.expect_err("404 是确定性失败,不重试也拿不到状态");
|
||||
assert!(!error.is_retryable(), "404 不该被判成可重试:{error:?}");
|
||||
tokio::time::sleep(Duration::from_millis(700)).await;
|
||||
assert_eq!(
|
||||
connections.load(Ordering::SeqCst),
|
||||
1,
|
||||
"确定性失败只允许打一次请求"
|
||||
);
|
||||
}
|
||||
|
||||
/// 产物超过调用方给的体积上限:按输出违约失败,且不重试。
|
||||
#[tokio::test]
|
||||
async fn artifact_download_rejects_bodies_over_the_caller_limit() {
|
||||
|
||||
Reference in New Issue
Block a user