e20ff2e89e
- platform-tripo:`TaskStatus::Unknown` 不再归一成 `OutputSchema` 终态失败,改成新的 `TripoError::TaskStatusUnknown`,并把它的可重试性判为真(SDK 自己也没把 unknown 算进终态) - platform-tripo:补用例——provider 回 `"status":"unknown"` 时必须返回可重试的 `TaskStatusUnknown`,不能当成正常快照 - api-server:错误映射补 `TaskStatusUnknown` 分支(502 + 固定文案,轮询侧早已按可重试继续查到截止时间) - api-server worker:状态查询遇到可重试错误时补一条 warn 日志,便于区分「一直是 unknown」和「一次网络抖动」 - docs:Provider 集成方案写明 unknown 按可重试处理及其理由
693 lines
29 KiB
Rust
693 lines
29 KiB
Rust
//! Tripo 3D 生成的 worker 执行链路。
|
||
//!
|
||
//! 与其它生成任务的关键差异是 **at-most-once submit**:provider 不接受幂等重放,
|
||
//! 所以只有“没有 checkpoint 的 job”才 submit,submit 成功后必须先落 checkpoint 再轮询;
|
||
//! 单次尝试(`max_attempts = 1`)失败即终态收口。
|
||
//!
|
||
//! 注意这不是「崩溃后接着查询」的续跑能力:3D job 的 `max_attempts = 1`,租约耗尽的 job 会被
|
||
//! 直接判终态失败、不会被重新 claim,所以不存在「重新 claim 后继续查询」的真实路径。
|
||
//! [`existing_checkpoint`] 的分支是一条**防御性硬约束**(任何带 checkpoint 的执行都绝不会
|
||
//! 再 submit 一次),checkpoint 的日常用途是人工对账。
|
||
|
||
use std::time::{Duration, Instant};
|
||
|
||
use axum::http::StatusCode;
|
||
use platform_tripo::{TripoProviderClient, TripoTaskHandle, TripoTaskSnapshot};
|
||
use serde_json::json;
|
||
use shared_contracts::editor_generation::{
|
||
editor_generation_stable_asset_id, editor_generation_stable_resource_id,
|
||
};
|
||
use shared_contracts::model3d::common::{Model3dGenerationTargetRef, Model3dTaskStatus};
|
||
use spacetime_client::{
|
||
EditorAssetCreateRecordInput, EditorProjectResourceCreateRecordInput,
|
||
ExternalGenerationJobProviderCheckpointRecordInput, ExternalGenerationJobRecord,
|
||
editor_project::EditorGenerationResultPersistItemRecordInput,
|
||
};
|
||
|
||
use crate::{
|
||
asset_billing::execute_billable_asset_operation_with_cost,
|
||
editor_project::{
|
||
EditorCanvasLayoutPlan, EditorGenerationCaller, build_editor_canvas_placement_item,
|
||
current_utc_micros, editor_generation_job_result_payload_json,
|
||
editor_media_src_from_object_key, generated_canvas_layer_id,
|
||
persist_editor_generation_result_atomically,
|
||
},
|
||
http_error::AppError,
|
||
state::AppState,
|
||
};
|
||
|
||
use super::{
|
||
artifacts::{MODEL3D_MAX_ARTIFACT_BYTES, read_artifact},
|
||
errors::map_provider_error,
|
||
image_source::resolve_image_input,
|
||
job::{MODEL3D_PROVIDER_KIND, Model3dJobRequest, parse_model3d_job_request},
|
||
provider::{TRIPO_PROVIDER, tripo_provider_client},
|
||
result::{build_result, job_result_payload},
|
||
storage::{
|
||
MODEL3D_ASSET_KIND, Model3dArtifactSlot, StoredModel3dArtifact, discard_model3d_artifact,
|
||
store_model3d_artifact,
|
||
},
|
||
};
|
||
|
||
/// 轮询间隔:provider 侧单次生成通常要几分钟,秒级轮询只会白烧查询配额。
|
||
const MODEL3D_POLL_INTERVAL: Duration = Duration::from_secs(5);
|
||
|
||
const MODEL3D_MODEL_SLOT: &str = "model";
|
||
const MODEL3D_PREVIEW_SLOT: &str = "preview";
|
||
|
||
pub(crate) async fn process_model3d_job(
|
||
state: &AppState,
|
||
caller: &EditorGenerationCaller,
|
||
job: &ExternalGenerationJobRecord,
|
||
provider_deadline: Instant,
|
||
) -> Result<(), String> {
|
||
let request =
|
||
parse_model3d_job_request(job.job_kind.as_str(), job.request_payload_json.as_str())?;
|
||
execute_billable_asset_operation_with_cost(
|
||
state,
|
||
caller.owner_user_id.as_str(),
|
||
MODEL3D_ASSET_KIND,
|
||
job.job_id.as_str(),
|
||
job.price_mud_points,
|
||
run_model3d_job(state, caller, job, request, provider_deadline),
|
||
)
|
||
.await
|
||
// 这里返回的字符串会成为 job 的 `last_error_message`,并原样回给前端(见
|
||
// `editor_generation_queue_state` 的 `error` 字段),所以刻意只取 provider 错误体里
|
||
// 人类可读的 message:用户可见文案不带内部 reason / taskStatus 枚举。
|
||
// 评审 #59 提过把结构化 details 一并写进去,按决定不改 —— 那会改掉已落库与前端展示的
|
||
// 文案口径;排障要定位 provider 侧任务时看 job 的 provider task id(checkpoint)。
|
||
.map_err(|error| error.body_text())
|
||
}
|
||
|
||
async fn run_model3d_job(
|
||
state: &AppState,
|
||
caller: &EditorGenerationCaller,
|
||
job: &ExternalGenerationJobRecord,
|
||
request: Model3dJobRequest,
|
||
provider_deadline: Instant,
|
||
) -> Result<(), AppError> {
|
||
let client = tripo_provider_client(&state.config)?;
|
||
let handle = match existing_checkpoint(job) {
|
||
Some(task_id) => TripoTaskHandle { task_id },
|
||
None => {
|
||
// submit 是不可逆且计费的调用:预算可能在别处就先耗光了(图生 3D 解析与上传
|
||
// 大图、任务排在其他 job 之后),这时提交上去也只会立刻被 deadline 判失败,
|
||
// 白花一次 provider 调用。只有还轮询得起才提交。
|
||
ensure_submit_budget(provider_deadline)?;
|
||
submit_and_checkpoint(state, caller, job, &client, &request, provider_deadline).await?
|
||
}
|
||
};
|
||
let snapshot = poll_until_terminal(&client, &handle, provider_deadline).await?;
|
||
let model = read_artifact(
|
||
client
|
||
.download_model(&snapshot, MODEL3D_MAX_ARTIFACT_BYTES)
|
||
.await
|
||
.map_err(map_provider_error)?,
|
||
)?;
|
||
let preview = read_artifact(
|
||
client
|
||
.download_rendered_image(&snapshot, MODEL3D_MAX_ARTIFACT_BYTES)
|
||
.await
|
||
.map_err(map_provider_error)?,
|
||
)?;
|
||
let (preview_width, preview_height) = preview_dimensions(preview.bytes.as_slice())?;
|
||
let stored_model = store_model3d_artifact(
|
||
state,
|
||
caller,
|
||
Model3dArtifactSlot::Model,
|
||
model.bytes,
|
||
model.content_type.as_str(),
|
||
)
|
||
.await?;
|
||
let stored_preview = match store_model3d_artifact(
|
||
state,
|
||
caller,
|
||
Model3dArtifactSlot::Preview,
|
||
preview.bytes,
|
||
preview.content_type.as_str(),
|
||
)
|
||
.await
|
||
{
|
||
Ok(stored) => stored,
|
||
Err(error) => {
|
||
discard_model3d_artifact(state, &stored_model).await;
|
||
return Err(error);
|
||
}
|
||
};
|
||
// 产物已经进了 OSS,但还没登记成正式资源:这一步失败就把刚传上去的对象删掉,
|
||
// 不让按 job 命名的私有残留对象一直堆在 bucket 里。
|
||
if let Err(error) = persist_result(
|
||
state,
|
||
caller,
|
||
job,
|
||
&request,
|
||
&stored_model,
|
||
&stored_preview,
|
||
preview_width,
|
||
preview_height,
|
||
)
|
||
.await
|
||
{
|
||
discard_model3d_artifact(state, &stored_model).await;
|
||
discard_model3d_artifact(state, &stored_preview).await;
|
||
return Err(error);
|
||
}
|
||
Ok(())
|
||
}
|
||
|
||
/// checkpoint 是 provider 侧任务的唯一凭据,也是 at-most-once 的开关:有值就只能查询,
|
||
/// 绝不允许再次 submit。
|
||
///
|
||
/// 当前 `max_attempts = 1` 下 job 不会被重新 claim(租约耗尽直接终态失败),所以这里守住的
|
||
/// 是一条防御性路径;它存在的意义是把「不许二次 submit」写成代码里的硬约束,而不是靠配置保证。
|
||
fn existing_checkpoint(job: &ExternalGenerationJobRecord) -> Option<String> {
|
||
job.provider_task_id
|
||
.as_deref()
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
.map(ToOwned::to_owned)
|
||
}
|
||
|
||
/// 提交前的预算检查:本次尝试已经没有轮询预算时,不发起不可逆且计费的 provider submit。
|
||
///
|
||
/// 调用点有两处:进入提交流程之前,以及图生 3D 在下载 / 上传源图之后、真正 submit 之前
|
||
/// (那两步都可能耗时到把预算用光)。
|
||
fn ensure_submit_budget(provider_deadline: Instant) -> Result<(), AppError> {
|
||
if Instant::now() >= provider_deadline {
|
||
return Err(provider_deadline_error());
|
||
}
|
||
Ok(())
|
||
}
|
||
|
||
async fn submit_and_checkpoint(
|
||
state: &AppState,
|
||
caller: &EditorGenerationCaller,
|
||
job: &ExternalGenerationJobRecord,
|
||
client: &TripoProviderClient,
|
||
request: &Model3dJobRequest,
|
||
provider_deadline: Instant,
|
||
) -> Result<TripoTaskHandle, AppError> {
|
||
let handle = match request {
|
||
Model3dJobRequest::TextToModel(request) => {
|
||
ensure_submit_budget(provider_deadline)?;
|
||
client.submit_text_to_model(&request.generation).await
|
||
}
|
||
Model3dJobRequest::ImageToModel(request) => {
|
||
// 解析源图要先下载站内图片、再上传给 provider,这一步可能把剩下的预算用光。
|
||
// 提交不可撤销且计费,所以真正发起之前再确认一次:宁可本次尝试判失败,
|
||
// 也不要提交上去紧接着被 deadline 判失败,白花一次外部调用。
|
||
let input = resolve_image_input(
|
||
state,
|
||
client,
|
||
caller.owner_user_id.as_str(),
|
||
&request.source,
|
||
)
|
||
.await?;
|
||
ensure_submit_budget(provider_deadline)?;
|
||
client
|
||
.submit_image_to_model(&input, &request.generation)
|
||
.await
|
||
}
|
||
}
|
||
.map_err(map_provider_error)?;
|
||
// submit 已经消耗了 provider 额度。checkpoint 写不进去时不能再退回可重试队列,
|
||
// 否则下一次 claim 会再次 submit;这里直接失败并保留“provider 已消耗”的对账信息。
|
||
state
|
||
.spacetime_client()
|
||
.set_external_generation_job_provider_checkpoint(
|
||
ExternalGenerationJobProviderCheckpointRecordInput {
|
||
job_id: job.job_id.clone(),
|
||
worker_id: worker_id(caller)?.to_string(),
|
||
lease_token: lease_token(job)?.to_string(),
|
||
provider_kind: MODEL3D_PROVIDER_KIND.to_string(),
|
||
provider_task_id: handle.task_id.clone(),
|
||
},
|
||
)
|
||
.await
|
||
.map_err(|error| {
|
||
AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"reason": "model3d-checkpoint-write-failed",
|
||
"message": format!(
|
||
"3D 生成任务已提交但 checkpoint 写入失败,本次尝试终态失败,等待人工对账:{error}"
|
||
),
|
||
}))
|
||
})?;
|
||
caller.report_processing_phase(state).await?;
|
||
Ok(handle)
|
||
}
|
||
|
||
async fn poll_until_terminal(
|
||
client: &TripoProviderClient,
|
||
handle: &TripoTaskHandle,
|
||
provider_deadline: Instant,
|
||
) -> Result<TripoTaskSnapshot, AppError> {
|
||
loop {
|
||
if Instant::now() >= provider_deadline {
|
||
return Err(provider_deadline_error());
|
||
}
|
||
let snapshot = match client.get_task(handle).await {
|
||
Ok(snapshot) => snapshot,
|
||
// 查询本身已经带过一轮退避重试;仍是瞬时故障时继续按轮询间隔查到预算用尽,
|
||
// 不因为一次 provider 抖动就作废这笔已经提交、不能重来的任务。
|
||
Err(error) if error.is_retryable() => {
|
||
// 包括 provider 报的 `unknown` 状态:继续查到截止时间,日志留痕便于区分
|
||
// 「一直是 unknown」和「一次网络抖动」。
|
||
tracing::warn!(error = %error, "3D 生成任务状态查询遇到可重试错误,继续轮询");
|
||
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 => {}
|
||
Model3dTaskStatus::Failed
|
||
| Model3dTaskStatus::Cancelled
|
||
| Model3dTaskStatus::Expired => return Err(task_not_completed(&snapshot)),
|
||
}
|
||
let remaining = provider_deadline.saturating_duration_since(Instant::now());
|
||
tokio::time::sleep(MODEL3D_POLL_INTERVAL.min(remaining)).await;
|
||
}
|
||
}
|
||
|
||
#[allow(clippy::too_many_arguments)]
|
||
async fn persist_result(
|
||
state: &AppState,
|
||
caller: &EditorGenerationCaller,
|
||
job: &ExternalGenerationJobRecord,
|
||
request: &Model3dJobRequest,
|
||
// 取引用而不是拿走所有权:登记失败时调用方还要按对象键把刚传上去的产物删掉。
|
||
model: &StoredModel3dArtifact,
|
||
preview: &StoredModel3dArtifact,
|
||
preview_width: u32,
|
||
preview_height: u32,
|
||
) -> Result<(), AppError> {
|
||
let operation = caller.operation.as_ref().ok_or_else(|| {
|
||
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"message": "落库 3D 生成结果时缺少稳定 operation 上下文。",
|
||
}))
|
||
})?;
|
||
let owner_user_id = caller.owner_user_id.as_str();
|
||
// 3D 资源在画布与素材库里按预览图渲染,模型本体留在 objectKey,因此两者互不覆盖。
|
||
let preview_src = editor_media_src_from_object_key(preview.object_key.as_str());
|
||
let completed_at_micros = current_utc_micros();
|
||
let audit_prompt = request.audit_prompt();
|
||
// 落库值必须是契约 wire 取值;取不到就按服务端问题失败,不写入 Debug 兜底名。
|
||
let model_version = request.model_version().map_err(|error| {
|
||
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"message": format!("落库 3D 生成结果时:{error}。"),
|
||
}))
|
||
})?;
|
||
// 落点与其它生成接口同形的平坦字段:校验层保证「至少一个」,两个都给时两条落点都写。
|
||
let target = request.target();
|
||
let project_id = target
|
||
.project_id
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
.map(str::to_string);
|
||
let folder_id = target
|
||
.asset_folder_id
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
.map(str::to_string);
|
||
let label = target
|
||
.asset_label
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
.map(str::to_string)
|
||
.unwrap_or_else(|| MODEL3D_DEFAULT_ASSET_LABEL.to_string());
|
||
let canvas_completion = target.canvas_completion.cloned();
|
||
let canvas_project_id = project_id.clone();
|
||
let resource_id = editor_generation_stable_resource_id(
|
||
owner_user_id,
|
||
operation.operation_kind.as_str(),
|
||
job.job_id.as_str(),
|
||
MODEL3D_MODEL_SLOT,
|
||
);
|
||
let asset_id = editor_generation_stable_asset_id(
|
||
owner_user_id,
|
||
operation.operation_kind.as_str(),
|
||
job.job_id.as_str(),
|
||
MODEL3D_MODEL_SLOT,
|
||
);
|
||
let has_resource = project_id.is_some();
|
||
let has_asset = folder_id.is_some();
|
||
// 派生来源只来自请求:资源 / 素材记录的 `source_resource_id` 记的是「从哪个画布资源
|
||
// 派生出来的」,与本次新生成的输出资源 id 无关(图生 3D 的 source,文生 3D 没有)。
|
||
let source_resource_id = request.source_resource_id();
|
||
let project_resource = project_id.map(|project_id| EditorProjectResourceCreateRecordInput {
|
||
resource_id: resource_id.clone(),
|
||
project_id,
|
||
owner_user_id: owner_user_id.to_string(),
|
||
asset_object_id: Some(model.asset_object.asset_object_id.clone()),
|
||
image_src: preview_src.clone(),
|
||
object_key: Some(model.object_key.clone()),
|
||
width: preview_width,
|
||
height: preview_height,
|
||
source_type: "generated".to_string(),
|
||
prompt: Some(audit_prompt.clone()),
|
||
actual_prompt: None,
|
||
model: Some(model_version.to_string()),
|
||
provider: Some(MODEL3D_PROVIDER_KIND.to_string()),
|
||
// 平台自己的 operation ID;provider task ID 只留在 checkpoint,不落业务行。
|
||
task_id: Some(job.job_id.clone()),
|
||
source_resource_id: source_resource_id.clone(),
|
||
asset_kind: Some(MODEL3D_ASSET_KIND.to_string()),
|
||
generation_inputs_json: None,
|
||
updated_at_micros: completed_at_micros,
|
||
image_sequence_frames_json: None,
|
||
image_sequence_duration_ms: None,
|
||
});
|
||
let asset_id_for_result = asset_id.clone();
|
||
let asset = folder_id.map(|folder_id| EditorAssetCreateRecordInput {
|
||
asset_id,
|
||
owner_user_id: owner_user_id.to_string(),
|
||
folder_id,
|
||
label,
|
||
asset_object_id: Some(model.asset_object.asset_object_id.clone()),
|
||
image_src: preview_src.clone(),
|
||
object_key: Some(model.object_key.clone()),
|
||
width: preview_width,
|
||
height: preview_height,
|
||
source_type: "generated".to_string(),
|
||
prompt: Some(audit_prompt),
|
||
actual_prompt: None,
|
||
model: Some(model_version.to_string()),
|
||
provider: Some(MODEL3D_PROVIDER_KIND.to_string()),
|
||
task_id: Some(job.job_id.clone()),
|
||
asset_kind: Some(MODEL3D_ASSET_KIND.to_string()),
|
||
generation_inputs_json: None,
|
||
// 素材行与画布资源行是同一份产物的两种落点:有本项目的资源行时指向它(与其它
|
||
// 生成工具一致),只落素材库时退回请求里的来源资源,不把溯源信息丢掉。
|
||
source_resource_id: has_resource
|
||
.then(|| resource_id.clone())
|
||
.or_else(|| source_resource_id.clone()),
|
||
generation_cost_mud_points: job.price_mud_points,
|
||
now_micros: completed_at_micros,
|
||
thumbnail_src: Some(preview_src.clone()),
|
||
group_task_id: None,
|
||
group_task_expected_asset_count: None,
|
||
image_sequence_frames_json: None,
|
||
image_sequence_duration_ms: None,
|
||
});
|
||
// 结果先构造:它引用的是对象元数据,与随后要搬进 item 的 asset_object 是两份值。
|
||
let result = build_result(
|
||
request.job_kind(),
|
||
target_ref(
|
||
has_resource.then(|| resource_id.clone()),
|
||
has_asset.then_some(asset_id_for_result),
|
||
),
|
||
model,
|
||
preview,
|
||
);
|
||
// 预览图单独登记成 asset_object:它没有自己的资源行,但必须属于当前账号,
|
||
// 否则后续按 objectKey 解析归属时会认为它未登记。
|
||
let items = vec![
|
||
EditorGenerationResultPersistItemRecordInput {
|
||
slot: MODEL3D_MODEL_SLOT.to_string(),
|
||
asset_object: Some(model.asset_object.clone()),
|
||
project_resource,
|
||
asset,
|
||
binding: None,
|
||
},
|
||
EditorGenerationResultPersistItemRecordInput {
|
||
slot: MODEL3D_PREVIEW_SLOT.to_string(),
|
||
asset_object: Some(preview.asset_object.clone()),
|
||
project_resource: None,
|
||
asset: None,
|
||
binding: None,
|
||
},
|
||
];
|
||
let canvas_layout_plan = match canvas_completion.as_ref() {
|
||
Some(completion) if has_resource => {
|
||
let layer_id = generated_canvas_layer_id(resource_id.as_str());
|
||
EditorCanvasLayoutPlan::generation(
|
||
owner_user_id,
|
||
canvas_project_id.as_deref(),
|
||
Some(completion),
|
||
vec![build_editor_canvas_placement_item(
|
||
completion,
|
||
resource_id.as_str(),
|
||
preview_width,
|
||
preview_height,
|
||
layer_id.clone(),
|
||
)],
|
||
Some(layer_id),
|
||
)
|
||
}
|
||
_ => EditorCanvasLayoutPlan::None,
|
||
};
|
||
let job_result_payload_json =
|
||
editor_generation_job_result_payload_json(caller, job_result_payload(result)?);
|
||
persist_editor_generation_result_atomically(
|
||
state,
|
||
caller,
|
||
items,
|
||
canvas_layout_plan,
|
||
job_result_payload_json,
|
||
completed_at_micros,
|
||
)
|
||
.await?;
|
||
Ok(())
|
||
}
|
||
|
||
/// 结果里的落点引用:写了项目资源行就带 `resourceId`,写了素材行就带 `assetId`;
|
||
/// 两个落点都写(前端画布链路的常态)时两个都带 —— 与请求落点同形的平坦字段。
|
||
fn target_ref(resource_id: Option<String>, asset_id: Option<String>) -> Model3dGenerationTargetRef {
|
||
Model3dGenerationTargetRef {
|
||
resource_id,
|
||
asset_id,
|
||
}
|
||
}
|
||
|
||
fn preview_dimensions(bytes: &[u8]) -> Result<(u32, u32), AppError> {
|
||
let reader = image::ImageReader::new(std::io::Cursor::new(bytes))
|
||
.with_guessed_format()
|
||
.map_err(|error| output_schema_error(format!("预览图格式无法识别:{error}")))?;
|
||
let (width, height) = reader
|
||
.into_dimensions()
|
||
.map_err(|error| output_schema_error(format!("预览图尺寸无法读取:{error}")))?;
|
||
if width == 0 || height == 0 {
|
||
return Err(output_schema_error("预览图尺寸为零".to_string()));
|
||
}
|
||
Ok((width, height))
|
||
}
|
||
|
||
pub(crate) const MODEL3D_DEFAULT_ASSET_LABEL: &str = "3D 模型";
|
||
|
||
fn worker_id(caller: &EditorGenerationCaller) -> Result<&str, AppError> {
|
||
caller
|
||
.operation
|
||
.as_ref()
|
||
.and_then(|operation| operation.worker_id.as_deref())
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
.ok_or_else(|| missing_operation_context("worker_id"))
|
||
}
|
||
|
||
fn lease_token(job: &ExternalGenerationJobRecord) -> Result<&str, AppError> {
|
||
job.lease_token
|
||
.as_deref()
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
.ok_or_else(|| missing_operation_context("lease_token"))
|
||
}
|
||
|
||
fn missing_operation_context(field: &str) -> AppError {
|
||
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"message": format!("3D 生成任务缺少 {field},无法写入 provider checkpoint。"),
|
||
}))
|
||
}
|
||
|
||
fn provider_deadline_error() -> AppError {
|
||
AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"reason": "tripo-provider-timeout",
|
||
"message": "3D 生成在本次 worker 预算内没有完成,本次尝试终态失败。",
|
||
}))
|
||
}
|
||
|
||
fn task_not_completed(snapshot: &TripoTaskSnapshot) -> AppError {
|
||
AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"reason": "tripo-task-failed",
|
||
"message": "3D 生成任务未成功完成。",
|
||
"taskStatus": snapshot.status,
|
||
"providerCode": snapshot.failure.as_ref().and_then(|failure| failure.code),
|
||
"providerMessage": snapshot.failure.as_ref().and_then(|failure| failure.message.clone()),
|
||
}))
|
||
}
|
||
|
||
fn output_schema_error(message: String) -> AppError {
|
||
AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({
|
||
"provider": TRIPO_PROVIDER,
|
||
"reason": "tripo-output-schema",
|
||
"message": message,
|
||
}))
|
||
}
|
||
|
||
#[cfg(test)]
|
||
mod tests {
|
||
use super::*;
|
||
use crate::tripo3d::job::Model3dJobKind;
|
||
use module_assets::{AssetObjectAccessPolicy, AssetObjectUpsertInput};
|
||
|
||
/// 预算已经耗尽时不能再提交:submit 不可逆且计费,提交上去也只会立刻被 deadline 判失败。
|
||
#[test]
|
||
fn submit_is_rejected_once_the_polling_budget_is_gone() {
|
||
let expired = ensure_submit_budget(Instant::now())
|
||
.expect_err("预算已耗尽的尝试不能再发起 provider submit");
|
||
assert_eq!(expired.status_code(), StatusCode::BAD_GATEWAY);
|
||
assert!(
|
||
ensure_submit_budget(Instant::now() + Duration::from_secs(1)).is_ok(),
|
||
"还有预算时必须放行"
|
||
);
|
||
}
|
||
|
||
/// at-most-once submit 完全依赖这条判定:checkpoint 有值就只能续跑查询。
|
||
/// 空白值必须等同于“没有 checkpoint”,否则一次空写入会把 job 永久锁死。
|
||
#[test]
|
||
fn checkpoint_is_only_recognized_when_provider_task_id_is_present() {
|
||
let mut job = job_fixture();
|
||
assert_eq!(existing_checkpoint(&job), None);
|
||
|
||
job.provider_task_id = Some(" ".to_string());
|
||
assert_eq!(existing_checkpoint(&job), None);
|
||
|
||
job.provider_task_id = Some(" task-123 ".to_string());
|
||
assert_eq!(existing_checkpoint(&job).as_deref(), Some("task-123"));
|
||
}
|
||
|
||
/// 完成结果按端点严格区分,且只带正式对象元数据:provider task ID、SDK 类型与
|
||
/// 带签名的临时地址都不允许出现在返回给调用方的结果里。
|
||
#[test]
|
||
fn completed_result_is_strict_per_endpoint_and_free_of_provider_facts() {
|
||
let model = stored_artifact("editor/model3d/task-1/model.glb", "model/gltf-binary", 4096);
|
||
let preview = stored_artifact("editor/model3d/task-1/preview.webp", "image/webp", 2048);
|
||
|
||
for (kind, expected_kind) in [
|
||
(Model3dJobKind::TextToModel, "textToModel"),
|
||
(Model3dJobKind::ImageToModel, "imageToModel"),
|
||
] {
|
||
let result = build_result(
|
||
kind,
|
||
Model3dGenerationTargetRef {
|
||
resource_id: Some("resource-1".to_string()),
|
||
asset_id: None,
|
||
},
|
||
&model,
|
||
&preview,
|
||
);
|
||
let payload = job_result_payload(result).expect("完成结果应可序列化");
|
||
let result = &payload["result"];
|
||
|
||
assert_eq!(result["kind"], json!(expected_kind));
|
||
assert_eq!(result["target"]["resourceId"], json!("resource-1"));
|
||
// 只落项目资源时不得凭空带出素材 ID:结果引用与请求落点同形。
|
||
assert_eq!(result["target"].get("assetId"), None);
|
||
assert_eq!(
|
||
result["model"]["objectKey"],
|
||
json!("editor/model3d/task-1/model.glb")
|
||
);
|
||
assert_eq!(result["model"]["contentLength"], json!(4096));
|
||
assert_eq!(result["preview"]["contentLength"], json!(2048));
|
||
|
||
let serialized = payload.to_string();
|
||
for forbidden in ["provider", "taskId", "providerTaskId", "https://"] {
|
||
assert!(
|
||
!serialized.contains(forbidden),
|
||
"完成结果不得携带 {forbidden}:{serialized}"
|
||
);
|
||
}
|
||
}
|
||
}
|
||
|
||
/// 预览图尺寸是资源行的宽高来源,读不出来就宁可失败退款,也不写占位尺寸。
|
||
#[test]
|
||
fn preview_dimensions_require_a_real_decodable_image() {
|
||
let image = image::RgbaImage::from_pixel(3, 2, image::Rgba([10, 20, 30, 255]));
|
||
let mut bytes = Vec::new();
|
||
image
|
||
.write_to(
|
||
&mut std::io::Cursor::new(&mut bytes),
|
||
image::ImageFormat::Png,
|
||
)
|
||
.expect("测试用 PNG 应可编码");
|
||
|
||
assert_eq!(
|
||
preview_dimensions(bytes.as_slice()).expect("应读出尺寸"),
|
||
(3, 2)
|
||
);
|
||
assert!(preview_dimensions(b"not an image").is_err());
|
||
}
|
||
|
||
fn stored_artifact(
|
||
object_key: &str,
|
||
content_type: &str,
|
||
content_length: u64,
|
||
) -> StoredModel3dArtifact {
|
||
StoredModel3dArtifact {
|
||
object_key: object_key.to_string(),
|
||
content_type: content_type.to_string(),
|
||
content_length,
|
||
sha256: "0".repeat(64),
|
||
asset_object: AssetObjectUpsertInput {
|
||
asset_object_id: format!("asset-object-{object_key}"),
|
||
bucket: "genarrative-test".to_string(),
|
||
object_key: object_key.to_string(),
|
||
access_policy: AssetObjectAccessPolicy::Private,
|
||
content_type: Some(content_type.to_string()),
|
||
content_length,
|
||
content_hash: Some("0".repeat(64)),
|
||
version: 1,
|
||
source_job_id: Some("extgen-1".to_string()),
|
||
owner_user_id: Some("user-1".to_string()),
|
||
profile_id: None,
|
||
entity_id: Some("extgen-1".to_string()),
|
||
asset_kind: MODEL3D_ASSET_KIND.to_string(),
|
||
updated_at_micros: 1_780_444_800_000_000,
|
||
},
|
||
}
|
||
}
|
||
|
||
fn job_fixture() -> ExternalGenerationJobRecord {
|
||
ExternalGenerationJobRecord {
|
||
job_id: "extgen-1".to_string(),
|
||
dedupe_key: "model3d-generation:user-1:model3d_text_to_model:issue-1".to_string(),
|
||
job_kind: "model3d_text_to_model".to_string(),
|
||
owner_user_id: "user-1".to_string(),
|
||
source_module: "editor-canvas".to_string(),
|
||
source_entity_id: "project-1".to_string(),
|
||
request_label: "文生 3D 模型".to_string(),
|
||
request_payload_json: "{}".to_string(),
|
||
status: "running".to_string(),
|
||
attempt: 1,
|
||
max_attempts: 1,
|
||
last_error_message: None,
|
||
worker_id: Some("worker-a".to_string()),
|
||
lease_expires_at: Some("2026-09-21T00:00:00Z".to_string()),
|
||
available_at: "2026-09-21T00:00:00Z".to_string(),
|
||
result_payload_json: None,
|
||
created_at: "2026-09-21T00:00:00Z".to_string(),
|
||
started_at: Some("2026-09-21T00:00:00Z".to_string()),
|
||
completed_at: None,
|
||
updated_at: "2026-09-21T00:00:00Z".to_string(),
|
||
updated_at_micros: 1_780_444_800_000_000,
|
||
lease_token: Some("lease-1".to_string()),
|
||
price_mud_points: 20,
|
||
refund_ledger_id: None,
|
||
notification_acknowledged_at: None,
|
||
notification_acknowledged_at_micros: None,
|
||
phase: Some("generating".to_string()),
|
||
provider_kind: None,
|
||
provider_task_id: None,
|
||
}
|
||
}
|
||
}
|