产物下载改为按「无数据推进」判超时
- artifact reqwest client 不再设包含响应体的总超时,改成 connect_timeout + read_timeout(reqwest 0.12 的「每次读操作」语义,读完自动重置) - 响应头之前没有数据推进时用显式 tokio::time::timeout 兜住,避免连上却不返回响应头时无限挂住 - 新增 TripoProviderClient::artifact_stall_timeout 统一承载三处预算,取值沿用 TripoSettings::request_timeout - 补 3 个 mock TCP 用例:响应头一直不来、传输中途断流按无数据推进失败,慢速但持续的下载必须完整收完(后者能拦住「改回总超时」的回归) - 同步技术方案与后端架构文档里的下载超时口径
This commit is contained in:
@@ -137,7 +137,7 @@ worker 配置:
|
||||
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS`:任务 lease 时长,默认 `600`;worker 会按约三分之一 lease、最长 30 秒的间隔续租。该值应覆盖一次心跳网络抖动窗口,不需要大于完整外部生成链路耗时。
|
||||
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS`:普通外部生成 job 的执行预算,默认 `900`。超过预算后当前 worker 停止续租并释放 worker 槽位,但不取消已启动的业务 future,也不主动写入失败 / 重试状态;在途执行交由 lease fencing 仲裁:写回在租约有效期内到达则照常完成,否则被拒绝,租约过期后任务可被重新认领,attempt 耗尽时由认领事务原子标记失败并结算退款。
|
||||
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS`:VectorEngine 图片生成 / 编辑、图标 spritesheet 生成、UI 素材提取以及角色动作、视频等长耗时 job 的执行预算,默认 `1800`。其中四类 VectorEngine 图片 job 固定为 `editor_image_generation`、`editor_image_edit`、`editor_icon_spritesheet_generation` 和 `editor_ui_design_asset_extraction`;手动去背景等不直接调用 VectorEngine 的 job 继续使用普通预算。Tripo 3D 的 `model3d_text_to_model` 与 `model3d_image_to_model` 同样使用该预算:它们在单次 attempt 内把 `get_task` 轮询到 provider 终态再下载产物,因此会持续占用一个 worker 并发位,接入真实流量前必须确认并发数与超时预算。
|
||||
- `TRIPO_BASE_URL` / `TRIPO_API_KEY` / `TRIPO_REQUEST_TIMEOUT_MS` / `TRIPO_RETRIES`:Tripo 3D provider 的网关、凭据、单请求超时与重试次数,由 HTTP 角色与 worker 共享同一份 API env。默认网关 `https://openapi.tripo3d.com/v3`;缺 API Key 时 3D 提交返回 503 且不扣费、不入队。worker 只使用 submit、单次 `get_task` 与产物下载三项能力,任务 ID 只作为 checkpoint 存在服务端。
|
||||
- `TRIPO_BASE_URL` / `TRIPO_API_KEY` / `TRIPO_REQUEST_TIMEOUT_MS` / `TRIPO_RETRIES`:Tripo 3D provider 的网关、凭据、API 单请求超时(同时也是产物下载的连接与「无数据推进」预算)与重试次数,由 HTTP 角色与 worker 共享同一份 API env。默认网关 `https://openapi.tripo3d.com/v3`;缺 API Key 时 3D 提交返回 503 且不扣费、不入队。worker 只使用 submit、单次 `get_task` 与产物下载三项能力,任务 ID 只作为 checkpoint 存在服务端。
|
||||
|
||||
worker 在单次 job 开始执行时从同一个单调时钟起点计算绝对 `job deadline` 和更早的 `provider deadline`:常规情况下为终态审计、OSS 持久化及 `complete/fail` 回写保留 `60` 秒;当整个 job 预算小于 `120` 秒时,保留其一半,避免 provider 预算被全部吃掉。该 deadline 只通过进程内 `RequestContext` 传给 VectorEngine 图片调用,不写入 HTTP DTO、队列 payload 或 SpacetimeDB;普通 HTTP / `inline` 上下文没有 deadline,保持原有行为。
|
||||
|
||||
|
||||
@@ -42,7 +42,7 @@ task 尚未完成时 `output` 为空;完成后 `output` 必须是与 task type
|
||||
|
||||
`get_task` 按 `TripoTaskHandle` 做单次查询,返回通用 `TripoTaskSnapshot`,由 `task_type` 决定 `output` 的具体 variant;adapter 不做轮询、不阻塞等待。
|
||||
|
||||
`download_model` 接受 `TripoTaskSnapshot`,先确认快照确实处于完成态,再从严格 endpoint 结果中取得已校验的 `model_url`,由 provider 自己的无鉴权 reqwest client 打开签名 URL 并返回 `TripoDownloadedArtifact` 流包装。包装只公开 `url`、`content_type`、`content_length`、`filename(name)` 和 `next_chunk()`;不提供完整 `Vec<u8>`,调用方必须逐块消费。`filename(name)` 的扩展名来自远端地址,只接受短的 ASCII 字母数字,其余退回 `glb`,避免远端地址里的 `%2F` 解码后拼出跨目录路径。provider 在流结束时校验实际接收字节数与 `Content-Length`,响应体中断或长度不一致按结构化请求错误失败;SDK 保持第三方原样,不承担产物下载。smoke example 使用异步文件写入逐块落盘;未来接入 OSS 时应把同一数据流直接送入 OSS 分片上传,不经过完整内存缓冲。
|
||||
`download_model` 接受 `TripoTaskSnapshot`,先确认快照确实处于完成态,再从严格 endpoint 结果中取得已校验的 `model_url`,由 provider 自己的无鉴权 reqwest client 打开签名 URL 并返回 `TripoDownloadedArtifact` 流包装。包装只公开 `url`、`content_type`、`content_length`、`filename(name)` 和 `next_chunk()`;不提供完整 `Vec<u8>`,调用方必须逐块消费。`filename(name)` 的扩展名来自远端地址,只接受短的 ASCII 字母数字,其余退回 `glb`,避免远端地址里的 `%2F` 解码后拼出跨目录路径。provider 在流结束时校验实际接收字节数与 `Content-Length`,响应体中断或长度不一致按结构化请求错误失败。产物下载不设总超时:几十 MB 的流只要还在出数据就不该被判失败、更不该重试后从零重传,因此下载链路只按「无数据推进」判超时,预算取 `TripoSettings::request_timeout`(连接用 `connect_timeout`,响应体用 `read_timeout`,响应头之前由显式的 `tokio::time::timeout` 兜住),超过预算没有数据推进即按传输失败返回;SDK 保持第三方原样,不承担产物下载。smoke example 使用异步文件写入逐块落盘;未来接入 OSS 时应把同一数据流直接送入 OSS 分片上传,不经过完整内存缓冲。
|
||||
|
||||
TODO:等待上游 `tripo-rust-sdk` 提供原生 artifact stream API 后,删除 provider-side reqwest 下载器,改由 SDK stream 直接承接。
|
||||
|
||||
@@ -53,4 +53,5 @@ TODO:等待上游 `tripo-rust-sdk` 提供原生 artifact stream API 后,删
|
||||
- 三个入口的输入校验可拒绝空白 `prompt` / `input`、视图不足两张和空白 `taskId`;组合校验在调用 SDK 前完成。
|
||||
- provider 错误统一为 adapter 错误类型。
|
||||
- shared contracts 的 ts-rs binding 无 feature 开关且始终可生成。
|
||||
- 产物下载只按「无数据推进」判超时:连接、响应头与每个数据块都必须有进展,慢速但持续的下载可完整收完,中途断流按请求错误失败。
|
||||
- 真实 Provider smoke 已覆盖三个入口;example 在任务未完成时跳过下载并继续跑后续入口,只打印脱敏后的结果结构和下载产物信息。
|
||||
|
||||
@@ -15,4 +15,4 @@ tokio = { workspace = true, features = ["time"] }
|
||||
url = { workspace = true }
|
||||
|
||||
[dev-dependencies]
|
||||
tokio = { workspace = true, features = ["fs", "io-util", "macros", "rt-multi-thread"] }
|
||||
tokio = { workspace = true, features = ["fs", "io-util", "macros", "net", "rt-multi-thread"] }
|
||||
|
||||
@@ -12,13 +12,19 @@ pub struct TripoProviderClient {
|
||||
pub(crate) client: TripoClient,
|
||||
artifact_client: reqwest::Client,
|
||||
artifact_retries: u32,
|
||||
/// 产物下载的「无数据推进」预算:连接、首字节和之后每个数据块都必须在它之内有进展,
|
||||
/// 下载总时长不设限。
|
||||
artifact_stall_timeout: Duration,
|
||||
}
|
||||
|
||||
impl TripoProviderClient {
|
||||
pub fn new(settings: TripoSettings) -> Result<Self, TripoError> {
|
||||
let artifact_client = reqwest::Client::builder()
|
||||
.user_agent(settings.user_agent.clone())
|
||||
.timeout(settings.request_timeout)
|
||||
// 产物是几十 MB 的流,设总超时会把「还在正常下载」判成失败,也会让重试从头重传;
|
||||
// 这里只限制「多久没有进展」:连接用 connect_timeout,传输过程用 read_timeout。
|
||||
.connect_timeout(settings.request_timeout)
|
||||
.read_timeout(settings.request_timeout)
|
||||
.build()
|
||||
.map_err(|error| TripoError::Request {
|
||||
message: format!("failed to build artifact download client: {error}"),
|
||||
@@ -28,6 +34,7 @@ impl TripoProviderClient {
|
||||
client: TripoClient::new(settings.client_options()).map_err(TripoError::from)?,
|
||||
artifact_client,
|
||||
artifact_retries: settings.retries,
|
||||
artifact_stall_timeout: settings.request_timeout,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -93,18 +100,21 @@ impl TripoProviderClient {
|
||||
task_id: &str,
|
||||
url: &TripoUrl,
|
||||
) -> Result<TripoDownloadedArtifact, TripoError> {
|
||||
let response = self
|
||||
.artifact_client
|
||||
.get(url.as_str())
|
||||
.send()
|
||||
.await
|
||||
.map_err(|error| TripoError::Request {
|
||||
message: format!(
|
||||
"artifact download transport failure for task {task_id}: {}",
|
||||
transport_error_kind(&error)
|
||||
),
|
||||
status: None,
|
||||
})?;
|
||||
// read_timeout 只覆盖响应体,响应头之前没有进展同样要按「无数据推进」判超时,
|
||||
// 否则连上却不给响应头的服务端会让下载无限挂住。
|
||||
let response = tokio::time::timeout(
|
||||
self.artifact_stall_timeout,
|
||||
self.artifact_client.get(url.as_str()).send(),
|
||||
)
|
||||
.await
|
||||
.map_err(|_elapsed| artifact_stall_error(task_id, self.artifact_stall_timeout))?
|
||||
.map_err(|error| TripoError::Request {
|
||||
message: format!(
|
||||
"artifact download transport failure for task {task_id}: {}",
|
||||
transport_error_kind(&error)
|
||||
),
|
||||
status: None,
|
||||
})?;
|
||||
let status = response.status();
|
||||
if !status.is_success() {
|
||||
return Err(TripoError::Request {
|
||||
@@ -132,6 +142,16 @@ fn download_backoff(attempt: u32) -> Duration {
|
||||
Duration::from_millis(250u64.saturating_mul(2u64.saturating_pow(attempt.min(6))))
|
||||
}
|
||||
|
||||
/// 产物下载长时间没有数据推进:连接、响应头或响应体任一段卡住都归一成这个错误。
|
||||
fn artifact_stall_error(task_id: &str, stall_timeout: Duration) -> TripoError {
|
||||
TripoError::Request {
|
||||
message: format!(
|
||||
"artifact download stalled for task {task_id}: no data received within {stall_timeout:?}"
|
||||
),
|
||||
status: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// 取完成态产物:下载入口只接受已完成且带输出的 task 快照。
|
||||
///
|
||||
/// `TripoTaskSnapshot` 字段全部公开,「完成态才有 output」只是 `map_task` 的约定,
|
||||
@@ -162,3 +182,164 @@ fn transport_error_kind(error: &reqwest::Error) -> &'static str {
|
||||
"request error"
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::net::SocketAddr;
|
||||
use std::time::Duration;
|
||||
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
use tokio::net::TcpListener;
|
||||
|
||||
use super::*;
|
||||
|
||||
/// 测试用的「无数据推进」预算:够短让用例跑得快,又给 CI 抖动留出余量。
|
||||
const STALL_TIMEOUT: Duration = Duration::from_millis(300);
|
||||
/// 正常推进的间隔:明显小于预算,保证用例只在「真的没数据」时才失败。
|
||||
const PROGRESS_INTERVAL: Duration = Duration::from_millis(100);
|
||||
|
||||
fn test_client() -> TripoProviderClient {
|
||||
let settings = TripoSettings::new(
|
||||
"test-key".to_string(),
|
||||
"http://127.0.0.1:1".to_string(),
|
||||
STALL_TIMEOUT,
|
||||
0,
|
||||
"genarrative-test-tripo/1".to_string(),
|
||||
);
|
||||
TripoProviderClient::new(settings).expect("测试配置必须能建出 client")
|
||||
}
|
||||
|
||||
fn artifact_url(addr: SocketAddr) -> TripoUrl {
|
||||
TripoUrl::parse(&format!("http://{addr}/model.glb")).expect("测试地址必须是合法产物 URL")
|
||||
}
|
||||
|
||||
async fn spawn_mock_server<F, Fut>(serve: F) -> SocketAddr
|
||||
where
|
||||
F: FnOnce(tokio::net::TcpStream) -> Fut + Send + 'static,
|
||||
Fut: std::future::Future<Output = ()> + Send,
|
||||
{
|
||||
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 Ok((socket, _)) = listener.accept().await else {
|
||||
return;
|
||||
};
|
||||
serve(socket).await;
|
||||
});
|
||||
addr
|
||||
}
|
||||
|
||||
async fn read_request(socket: &mut tokio::net::TcpStream) {
|
||||
let mut buffer = [0u8; 1024];
|
||||
let _ = socket.read(&mut buffer).await;
|
||||
}
|
||||
|
||||
/// 接受连接后一直不返回响应头,直到远超预算才断开。
|
||||
async fn spawn_silent_server() -> SocketAddr {
|
||||
spawn_mock_server(|_socket| async {
|
||||
tokio::time::sleep(Duration::from_secs(30)).await;
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// 响应头与第一块数据都正常,之后彻底断流。
|
||||
async fn spawn_stalling_body_server() -> SocketAddr {
|
||||
spawn_mock_server(|mut socket| async move {
|
||||
read_request(&mut socket).await;
|
||||
let head = "HTTP/1.1 200 OK\r\n\
|
||||
Content-Type: model/gltf-binary\r\n\
|
||||
Content-Length: 64\r\n\
|
||||
\r\n";
|
||||
if socket.write_all(head.as_bytes()).await.is_err() {
|
||||
return;
|
||||
}
|
||||
if socket.write_all(b"glTF").await.is_err() {
|
||||
return;
|
||||
}
|
||||
let _ = socket.flush().await;
|
||||
tokio::time::sleep(Duration::from_secs(30)).await;
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// 总量固定但推得很慢:总时长远超预算,每次却都有数据推进。
|
||||
async fn spawn_slow_progress_server() -> SocketAddr {
|
||||
spawn_mock_server(|mut socket| async move {
|
||||
let _ = socket.set_nodelay(true);
|
||||
read_request(&mut socket).await;
|
||||
let head = "HTTP/1.1 200 OK\r\n\
|
||||
Content-Type: model/gltf-binary\r\n\
|
||||
Content-Length: 10\r\n\
|
||||
\r\n";
|
||||
if socket.write_all(head.as_bytes()).await.is_err() {
|
||||
return;
|
||||
}
|
||||
for _ in 0..5 {
|
||||
tokio::time::sleep(PROGRESS_INTERVAL).await;
|
||||
if socket.write_all(b"gl").await.is_err() {
|
||||
return;
|
||||
}
|
||||
let _ = socket.flush().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn artifact_download_times_out_when_response_headers_never_arrive() {
|
||||
let addr = spawn_silent_server().await;
|
||||
let error = match test_client()
|
||||
.download_artifact("task-1", &artifact_url(addr))
|
||||
.await
|
||||
{
|
||||
Ok(_) => panic!("响应头一直不来必须按无数据推进失败"),
|
||||
Err(error) => error,
|
||||
};
|
||||
assert!(matches!(error, TripoError::Request { .. }), "{error:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn artifact_download_times_out_when_body_stops_progressing() {
|
||||
let addr = spawn_stalling_body_server().await;
|
||||
let mut artifact = test_client()
|
||||
.download_artifact_once("task-1", &artifact_url(addr))
|
||||
.await
|
||||
.expect("响应头与第一块数据必须正常返回");
|
||||
|
||||
let first = artifact
|
||||
.next_chunk()
|
||||
.await
|
||||
.expect("第一块数据必须正常返回")
|
||||
.expect("第一块数据必须存在");
|
||||
assert_eq!(first.len(), 4);
|
||||
|
||||
let error = artifact
|
||||
.next_chunk()
|
||||
.await
|
||||
.expect_err("传输中途断流必须按无数据推进失败");
|
||||
assert!(matches!(error, TripoError::Request { .. }), "{error:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn artifact_download_survives_slow_but_progressing_body() {
|
||||
let addr = spawn_slow_progress_server().await;
|
||||
let mut artifact = test_client()
|
||||
.download_artifact_once("task-1", &artifact_url(addr))
|
||||
.await
|
||||
.expect("响应头必须正常返回");
|
||||
|
||||
let mut received = Vec::new();
|
||||
while let Some(chunk) = artifact
|
||||
.next_chunk()
|
||||
.await
|
||||
.expect("只要还在出数据就不能判超时")
|
||||
{
|
||||
received.extend_from_slice(&chunk);
|
||||
}
|
||||
assert_eq!(received.len(), 10, "慢但持续的下载必须完整收完");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user