use std::{ future::Future, io, pin::Pin, time::{Duration, Instant}, }; use axum::Json; use serde_json::{Value, json}; use shared_kernel::offset_datetime_to_unix_micros; use spacetime_client::{ ExternalGenerationJobClaimRecordInput, ExternalGenerationJobCompleteRecordInput, ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord, ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription, }; use tokio::{ sync::{OwnedSemaphorePermit, Semaphore}, task::{JoinHandle, JoinSet}, time::sleep, }; use tracing::{error, info, warn}; const MAX_EDITOR_GENERATION_WARNING_CHARS: usize = 2_048; // provider 必须先结束,给失败审计、计费结算和队列终态写回保留有效 lease 内的收尾窗口。 const EXTERNAL_GENERATION_WORKER_TERMINAL_WRITE_RESERVE: Duration = Duration::from_secs(60); const EDITOR_GENERATION_SLICE_WARNING_PREFIX: &str = "图集已生成,但自动拆分未完成:"; const EDITOR_GENERATION_WARNING_REDACTED_MESSAGE: &str = "自动拆分未完成(告警详情含内联媒体引用,已省略)"; use crate::{ asset_billing::with_external_generation_billing_attempt_context, character_animation_assets::{ generate_editor_character_animation_for_owner, generate_editor_video_for_owner, }, config::AppConfig, editor_generation_queue::{ EDITOR_BACKGROUND_MUSIC_GENERATION_JOB_KIND, EDITOR_BACKGROUND_REMOVAL_JOB_KIND, EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND, EDITOR_ICON_SPEC_GENERATION_JOB_KIND, EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND, EDITOR_IMAGE_EDIT_JOB_KIND, EDITOR_IMAGE_GENERATION_JOB_KIND, EDITOR_SOUND_EFFECT_GENERATION_JOB_KIND, EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND, EDITOR_VIDEO_GENERATION_JOB_KIND, GAME_CREATOR_CLIENT_GENERATION_DEDUPE_PREFIX, }, editor_project::{ EDITOR_GENERATION_MULTIPLE_WARNINGS_CODE, EDITOR_IMAGE_EDIT_QUEUE_PAYLOAD_VERSION, EditorBackgroundRemovalRequest, EditorGenerationCaller, EditorGenerationOperationContext, EditorGenerationPhaseReporter, EditorGenerationQueueResultContext, EditorImageEditQueuePayload, EditorImageEditRequest, EditorImageEditResolvedSource, EditorImageGenerationRequest, EditorUiDesignAssetExtractionRequest, compact_editor_generation_result, edit_editor_image_for_owner_with_source_snapshot, extract_editor_ui_design_assets_for_owner, generate_editor_image_for_owner, remove_editor_image_background_for_owner, }, editor_project_icon::{ EditorIconSpecGenerationRequest, EditorIconSpritesheetGenerationRequest, generate_editor_icon_spritesheet_for_owner, generate_icon_spec_for_owner, }, request_context::RequestContext, state::AppState, vector_engine_audio_generation::{ generate_editor_background_music_for_owner, generate_editor_sound_effect_for_owner, }, }; #[cfg(any())] use crate::{ puzzle::{ ExternalGenerationWriteLeaseGuard, PuzzleCompileDraftWorkerPayload, PuzzleGenerateImagesWorkerPayload, PuzzleGenerateUiBackgroundWorkerPayload, execute_puzzle_compile_draft_worker_job, execute_puzzle_generate_images_worker_job, execute_puzzle_generate_ui_background_worker_job, release_puzzle_compile_background_claim, }, puzzle_clear::{ PUZZLE_CLEAR_COMPILE_DRAFT_JOB_KIND, PuzzleClearCompileDraftWorkerPayload, execute_puzzle_clear_compile_draft_worker_job, }, state::PuzzleApiState, wooden_fish::{ WOODEN_FISH_GENERATE_IMAGE_ASSETS_JOB_KIND, WoodenFishGenerateImageAssetsWorkerPayload, execute_wooden_fish_generate_image_assets_worker_job, }, }; #[cfg(any())] pub(crate) const PUZZLE_COMPILE_DRAFT_JOB_KIND: &str = "puzzle_compile_draft"; #[cfg(any())] pub(crate) const PUZZLE_GENERATE_IMAGES_JOB_KIND: &str = "puzzle_generate_images"; #[cfg(any())] pub(crate) const PUZZLE_GENERATE_UI_BACKGROUND_JOB_KIND: &str = "puzzle_generate_ui_background"; pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<(), io::Error> { let worker_id = state.config.external_generation_worker_id.clone(); let concurrency = state.config.external_generation_worker_concurrency.max(1); let poll_interval = state.config.external_generation_worker_poll_interval; let lease = state.config.external_generation_worker_lease; // 超时任务不能立即取消(在途 procedure 仍可能写回),因此执行容量必须同时 // 约束 active 与 detached work;否则每次超时都会释放 tasks 槽位,实际内存占用 // 会超过配置并发。 let work_slots = std::sync::Arc::new(Semaphore::new(concurrency)); let mut tasks = JoinSet::new(); let mut shutdown = external_generation_worker_shutdown_signal(); let mut queue_wake = None; info!( worker_id, concurrency, poll_interval_ms = poll_interval.as_millis(), lease_seconds = lease.as_secs(), job_timeout_seconds = state .config .external_generation_worker_job_timeout .as_secs(), long_job_timeout_seconds = state .config .external_generation_worker_long_job_timeout .as_secs(), "external generation worker 已启动" ); loop { // 持续有队列任务时不会进入等待分支,因此必须在每轮主动回收已完成的 // JoinHandle;否则 permit 虽已归还,JoinSet 仍会保留每个历史任务的句柄。 reap_finished_external_generation_worker_tasks(&mut tasks); ensure_external_generation_queue_wake_subscription(&state, &mut queue_wake).await; while work_slots.available_permits() == 0 { tokio::select! { _ = shutdown.as_mut() => { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); } permit = work_slots.clone().acquire_owned() => { if permit.is_err() { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); } // 只用 acquire 作为容量变化唤醒信号,许可立即归还;真正领取任务 // 时在下方按返回的 job 数量逐个 try_acquire。 } } } let available = work_slots.available_permits().max(1); let now_micros = current_utc_micros(); let lease_expires_at_micros = now_micros.saturating_add(duration_micros_i64(lease)); let claim_jobs = state.spacetime_client().claim_external_generation_jobs( ExternalGenerationJobClaimRecordInput { worker_id: worker_id.clone(), limit: available.min(u32::MAX as usize) as u32, lease_expires_at_micros, claimed_at_micros: now_micros, }, ); tokio::pin!(claim_jobs); let jobs = match tokio::select! { _ = shutdown.as_mut() => { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); } result = &mut claim_jobs => result } { Ok(jobs) => jobs, Err(error) => { error!(error = %error, "领取外部生成任务失败,等待下一轮重试"); if await_one_task_or_queue_wake_or_sleep_or_shutdown( &mut tasks, &mut queue_wake, sleep(poll_interval), &mut shutdown, ) .await { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); } continue; } }; if jobs.is_empty() { if await_one_task_or_queue_wake_or_sleep_or_shutdown( &mut tasks, &mut queue_wake, sleep(poll_interval), &mut shutdown, ) .await { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); } continue; } for job in jobs { let state = state.clone(); let worker_id = worker_id.clone(); let permit = work_slots .clone() .try_acquire_owned() .expect("claimed job must have an execution capacity permit"); tasks.spawn(async move { if let Err(error) = process_external_generation_job(state, worker_id, lease, job, permit).await { error!(error = %error, "external generation worker 执行任务失败"); } }); } } } async fn ensure_external_generation_queue_wake_subscription( state: &AppState, queue_wake: &mut Option, ) { if queue_wake.is_some() { return; } match state .spacetime_client() .subscribe_external_generation_queue_wake() .await { Ok(subscription) => { *queue_wake = Some(subscription); info!("external generation worker 已订阅队列变更唤醒"); } Err(error) => { warn!(error = %error, "external generation worker 订阅队列变更失败,暂时使用间隔兜底"); } } } type ExternalGenerationShutdownSignal = Pin + Send>>; fn external_generation_worker_shutdown_signal() -> ExternalGenerationShutdownSignal { Box::pin(async { wait_for_external_generation_worker_shutdown_signal().await; }) } #[cfg(unix)] async fn wait_for_external_generation_worker_shutdown_signal() { use tokio::signal::unix::{SignalKind, signal}; let mut sigterm = signal(SignalKind::terminate()).ok(); tokio::select! { result = tokio::signal::ctrl_c() => { if let Err(error) = result { warn!(error = %error, "external generation worker 监听 SIGINT 失败"); } } _ = async { if let Some(sigterm) = sigterm.as_mut() { sigterm.recv().await; } else { std::future::pending::<()>().await; } } => {} } } #[cfg(not(unix))] async fn wait_for_external_generation_worker_shutdown_signal() { if let Err(error) = tokio::signal::ctrl_c().await { warn!(error = %error, "external generation worker 监听 Ctrl-C 失败"); } } async fn await_worker_task(tasks: &mut JoinSet<()>) { if let Some(result) = tasks.join_next().await && let Err(error) = result { error!(error = %error, "external generation worker 子任务 panic"); } } fn reap_finished_external_generation_worker_tasks(tasks: &mut JoinSet<()>) { while let Some(result) = tasks.try_join_next() { if let Err(error) = result { error!(error = %error, "external generation worker 子任务 panic"); } } } async fn await_one_task_or_queue_wake_or_sleep_or_shutdown( tasks: &mut JoinSet<()>, queue_wake: &mut Option, sleeper: impl Future, shutdown: &mut ExternalGenerationShutdownSignal, ) -> bool { let queue_changed = await_external_generation_queue_wake(queue_wake); tokio::pin!(queue_changed); tokio::pin!(sleeper); if tasks.is_empty() { tokio::select! { _ = shutdown.as_mut() => true, _ = &mut queue_changed => false, _ = &mut sleeper => false, } } else { tokio::select! { _ = shutdown.as_mut() => true, _ = &mut queue_changed => false, _ = &mut sleeper => false, result = tasks.join_next() => { if let Some(Err(error)) = result { error!(error = %error, "external generation worker 子任务 panic"); } false } } } } async fn await_external_generation_queue_wake( queue_wake: &mut Option, ) { let result = if let Some(subscription) = queue_wake.as_mut() { subscription.changed().await } else { std::future::pending().await }; if let Err(error) = result { warn!(error = %error, "external generation worker 队列变更订阅已失效,等待兜底间隔后重连"); *queue_wake = None; std::future::pending::<()>().await; } } async fn drain_external_generation_worker_tasks(tasks: &mut JoinSet<()>) { info!( in_flight_jobs = tasks.len(), "external generation worker 收到停机信号,停止领取新任务并等待当前任务完成" ); while !tasks.is_empty() { await_worker_task(tasks).await; } info!("external generation worker 已完成优雅停机"); } async fn process_external_generation_job( state: AppState, worker_id: String, lease: Duration, job: ExternalGenerationJobRecord, permit: OwnedSemaphorePermit, ) -> Result<(), String> { let heartbeat_interval = external_generation_worker_heartbeat_interval(lease); let job_timeout = external_generation_worker_job_timeout(&state.config, job.job_kind.as_str()); let (job_deadline, provider_deadline) = external_generation_worker_deadlines(Instant::now(), job_timeout); let billing_price_mud_points = match external_generation_billing_price_mud_points(&job) { Ok(price_mud_points) => price_mud_points, Err(message) => { fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let billing_claim_attempt = match external_generation_billing_claim_attempt(&job) { Ok(claim_attempt) => claim_attempt, Err(message) => { fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; // work 单独 spawn:已发出的 SpacetimeDB procedure 无法通过 drop 客户端 future 撤销, // 若超时时 drop 在途写回再抢先写失败态,可能出现“业务写回已提交、任务却被标记失败 // 并退款”的不一致。因此超时/续租失败时不取消 work,也不在客户端写失败态,交由服务端 // lease fencing 仲裁:写回在租约有效期内到达则任务照常完成,否则被拒绝;租约到期后 // 任务被重新认领,attempt 耗尽时由认领事务原子地标记失败并结算退款。 let mut work_handle = tokio::spawn(with_external_generation_billing_attempt_context( job.job_id.clone(), billing_claim_attempt, billing_price_mud_points, process_external_generation_job_once( state.clone(), worker_id.clone(), job.clone(), provider_deadline, ), )); let heartbeat = maintain_external_generation_job_lease(&state, &worker_id, &job, lease, heartbeat_interval); match await_external_generation_job_execution( join_external_generation_work(&mut work_handle), heartbeat, job_deadline, ) .await { ExternalGenerationJobExecutionOutcome::Finished(result) => result, ExternalGenerationJobExecutionOutcome::TimedOut => { let message = external_generation_worker_timeout_message(&job, job_timeout); warn!( job_id = %job.job_id, job_kind = %job.job_kind, timeout_seconds = job_timeout.as_secs(), "external generation worker 任务超过执行预算,停止续租并保留 worker 槽位,在途执行交由租约仲裁" ); detach_external_generation_work_until_lease_expiry( work_handle, &job, lease, "任务超过执行预算", Some(permit), ); Err(message) } ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error) => { detach_external_generation_work_until_lease_expiry( work_handle, &job, lease, "任务租约续期失败", Some(permit), ); Err(error) } } } async fn join_external_generation_work( work_handle: &mut JoinHandle>, ) -> Result<(), String> { match work_handle.await { Ok(result) => result, Err(join_error) => Err(format!("external generation work 任务 panic:{join_error}")), } } /// 停止等待 work 后不能直接取消它:在途 procedure 可能已在服务端提交。让它继续跑完当前 /// 尝试,写回由 lease fencing 裁决是否生效。等到租约必然过期(2 倍租约时长,此后一切 /// 围栏写回都会被拒绝)仍未结束的任务才安全取消,取消触发的计费补偿退款按 attempt /// 账本幂等,与队列侧耗尽结算不会重复。 fn detach_external_generation_work_until_lease_expiry( mut work_handle: JoinHandle>, job: &ExternalGenerationJobRecord, lease: Duration, reason: &'static str, permit: Option, ) { let job_id = job.job_id.clone(); let job_kind = job.job_kind.clone(); let grace = lease.saturating_mul(2); tokio::spawn(async move { match tokio::time::timeout(grace, &mut work_handle).await { Ok(Ok(Ok(()))) => info!( job_id = %job_id, job_kind = %job_kind, reason, "external generation worker 脱管任务已在租约仲裁窗口内完成写回" ), Ok(Ok(Err(error))) => info!( job_id = %job_id, job_kind = %job_kind, reason, error = %error, "external generation worker 脱管任务在租约仲裁窗口内以失败结束" ), Ok(Err(join_error)) => error!( job_id = %job_id, job_kind = %job_kind, reason, error = %join_error, "external generation worker 脱管任务 panic" ), Err(_) => { work_handle.abort(); // 仅调用 abort 不会从 JoinHandle/JoinSet 中消费完成结果;等待被取消 // 的 handle,确保 permit 与任务句柄在同一生命周期内一起释放。 let _ = work_handle.await; warn!( job_id = %job_id, job_kind = %job_kind, reason, grace_ms = grace.as_millis() as u64, "external generation worker 脱管任务超过租约仲裁窗口仍未结束,已取消;其后续写回将被 lease fencing 拒绝" ); } } // 保持执行许可直到 work 真正结束或被取消,避免超时任务脱管后继续 // 累积图片/音频响应占用。 drop(permit); }); } #[derive(Debug, PartialEq, Eq)] enum ExternalGenerationJobExecutionOutcome { Finished(Result<(), String>), TimedOut, LeaseRenewalFailed(String), } async fn await_external_generation_job_execution( work: W, heartbeat: H, job_deadline: Instant, ) -> ExternalGenerationJobExecutionOutcome where W: Future>, H: Future>, { tokio::pin!(work); tokio::pin!(heartbeat); let job_deadline_sleep = tokio::time::sleep_until(job_deadline.into()); tokio::pin!(job_deadline_sleep); tokio::select! { biased; result = &mut work => ExternalGenerationJobExecutionOutcome::Finished(result), _ = &mut job_deadline_sleep => ExternalGenerationJobExecutionOutcome::TimedOut, result = &mut heartbeat => match result { Ok(()) => unreachable!("external generation heartbeat monitor should not finish"), Err(error) => ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error), }, } } async fn maintain_external_generation_job_lease( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, lease: Duration, heartbeat_interval: Duration, ) -> Result<(), String> { loop { sleep(heartbeat_interval).await; renew_job_lease(state, worker_id, job, lease).await?; } } async fn process_external_generation_job_once( state: AppState, worker_id: String, job: ExternalGenerationJobRecord, provider_deadline: Instant, ) -> Result<(), String> { match job.job_kind.as_str() { #[cfg(any())] PUZZLE_COMPILE_DRAFT_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("拼图生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ); let puzzle_state = PuzzleApiState::from_ref(&state); let write_guard = build_external_generation_write_lease_guard(&worker_id, &job)?; match execute_puzzle_compile_draft_worker_job( &puzzle_state, &request_context, payload.clone(), write_guard, ) .await { Ok(session) => { let result = complete_job( &state, &worker_id, &job, Some( json!({ "sessionId": session.session_id, "progressPercent": session.progress_percent, }) .to_string(), ), ) .await; if result.is_ok() { release_puzzle_compile_background_claim(&puzzle_state, &payload); } result } Err(error) => { let message = error.body_text(); let should_release_claim = error.should_fail_queue_job(); let result = fail_queue_job_after_worker_error( &state, &worker_id, &job, &error, &message, ) .await; if result.is_ok() && should_release_claim { release_puzzle_compile_background_claim(&puzzle_state, &payload); } result?; Err(message) } } } #[cfg(any())] PUZZLE_GENERATE_IMAGES_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("拼图关卡图片生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ); let puzzle_state = PuzzleApiState::from_ref(&state); let write_guard = build_external_generation_write_lease_guard(&worker_id, &job)?; match execute_puzzle_generate_images_worker_job( &puzzle_state, &request_context, payload, write_guard, ) .await { Ok(session) => { complete_job( &state, &worker_id, &job, Some( json!({ "sessionId": session.session_id, "progressPercent": session.progress_percent, }) .to_string(), ), ) .await } Err(error) => { let message = error.body_text(); fail_queue_job_after_worker_error(&state, &worker_id, &job, &error, &message) .await?; Err(message) } } } #[cfg(any())] PUZZLE_GENERATE_UI_BACKGROUND_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("拼图 UI 背景图生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ); let puzzle_state = PuzzleApiState::from_ref(&state); let write_guard = build_external_generation_write_lease_guard(&worker_id, &job)?; match execute_puzzle_generate_ui_background_worker_job( &puzzle_state, &request_context, payload, write_guard, ) .await { Ok(session) => { complete_job( &state, &worker_id, &job, Some( json!({ "sessionId": session.session_id, "progressPercent": session.progress_percent, }) .to_string(), ), ) .await } Err(error) => { let message = error.body_text(); fail_queue_job_after_worker_error(&state, &worker_id, &job, &error, &message) .await?; Err(message) } } } #[cfg(any())] PUZZLE_CLEAR_COMPILE_DRAFT_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("拼消消生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ); match execute_puzzle_clear_compile_draft_worker_job(&state, &request_context, payload) .await { Ok(session) => { complete_job( &state, &worker_id, &job, Some( json!({ "sessionId": session.session_id, "status": session.status, }) .to_string(), ), ) .await } Err(response) => { let message = response_error_message(response).await; fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } #[cfg(any())] WOODEN_FISH_GENERATE_IMAGE_ASSETS_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("敲木鱼图片生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ); match execute_wooden_fish_generate_image_assets_worker_job( &state, &request_context, payload, ) .await { Ok(session) => { complete_job( &state, &worker_id, &job, Some( json!({ "sessionId": session.session_id, "status": session.status, }) .to_string(), ), ) .await } Err(response) => { let message = response_error_message(response).await; fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_IMAGE_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布生图任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_editor_image_for_owner( &state, &request_context, editor_generation_worker_caller(&worker_id, &job)?, payload, ) .await { Ok(_) => Ok(()), Err(error) => { let message = error.body_text(); fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_IMAGE_EDIT_JOB_KIND => { let payload = match parse_editor_image_edit_worker_payload(&job) { Ok(payload) => payload, Err(message) => { fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match edit_editor_image_for_owner_with_source_snapshot( &state, &request_context, editor_generation_worker_caller(&worker_id, &job)?, payload.request, payload.source.as_ref(), ) .await { Ok(_) => Ok(()), Err(error) => { let message = error.body_text(); fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_BACKGROUND_REMOVAL_JOB_KIND => { let mut payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布去除背景任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; payload.task_id.get_or_insert_with(|| job.job_id.clone()); let request_context = worker_request_context(&job, provider_deadline); match remove_editor_image_background_for_owner( &state, &request_context, editor_generation_worker_caller(&worker_id, &job)?, payload, ) .await { Ok(_) => Ok(()), Err(error) => { let message = error.body_text(); fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布图标素材任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_editor_icon_spritesheet_for_owner( &state, &request_context, editor_generation_worker_caller(&worker_id, &job)?, payload, ) .await { Ok(_) => Ok(()), Err(error) => { let message = error.body_text(); fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_ICON_SPEC_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布图标规范任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_icon_spec_for_owner( &state, &request_context, editor_generation_worker_caller(&worker_id, &job)?, payload, ) .await { Ok(_) => Ok(()), Err(error) => { let message = error.body_text(); fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND => { let payload = match serde_json::from_str::( job.request_payload_json.as_str(), ) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布 UI 素材提取任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match extract_editor_ui_design_assets_for_owner( &state, &request_context, editor_generation_worker_caller(&worker_id, &job)?, payload, ) .await { Ok(_) => Ok(()), Err(error) => { let message = error.body_text(); fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::< shared_contracts::assets::EditorCharacterAnimationGenerateRequest, >(job.request_payload_json.as_str()) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布角色动作任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_editor_character_animation_for_owner( state.clone(), request_context, editor_generation_worker_caller(&worker_id, &job)?, Ok(Json(payload)), ) .await { Ok(_) => Ok(()), Err(response) => { let message = response_error_message(response).await; fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_VIDEO_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::< shared_contracts::assets::EditorVideoGenerateRequest, >(job.request_payload_json.as_str()) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布视频生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_editor_video_for_owner( state.clone(), request_context, editor_generation_worker_caller(&worker_id, &job)?, Ok(Json(payload)), ) .await { Ok(_) => Ok(()), Err(response) => { let message = response_error_message(response).await; fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_SOUND_EFFECT_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::< shared_contracts::assets::EditorSoundEffectGenerateRequest, >(job.request_payload_json.as_str()) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布音效生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_editor_sound_effect_for_owner( state.clone(), request_context, editor_generation_worker_caller(&worker_id, &job)?, Ok(Json(payload)), ) .await { Ok(_) => Ok(()), Err(response) => { let message = response_error_message(response).await; fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } EDITOR_BACKGROUND_MUSIC_GENERATION_JOB_KIND => { let payload = match serde_json::from_str::< shared_contracts::assets::EditorBackgroundMusicGenerateRequest, >(job.request_payload_json.as_str()) { Ok(payload) => payload, Err(error) => { let message = format!("图片画布背景音乐生成任务参数解析失败:{error}"); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } }; let request_context = worker_request_context(&job, provider_deadline); match generate_editor_background_music_for_owner( state.clone(), request_context, editor_generation_worker_caller(&worker_id, &job)?, Ok(Json(payload)), ) .await { Ok(_) => Ok(()), Err(response) => { let message = response_error_message(response).await; fail_job(&state, &worker_id, &job, message.clone()).await?; Err(message) } } } unknown => { warn!( job_id = job.job_id, job_kind = unknown, "external generation worker 收到暂不支持的任务类型" ); fail_job( &state, &worker_id, &job, format!("暂不支持的外部生成任务类型:{unknown}"), ) .await } } } async fn response_error_message(response: axum::response::Response) -> String { use axum::body::to_bytes; let status = response.status(); let body_bytes = match to_bytes(response.into_body(), 64 * 1024).await { Ok(bytes) => bytes, Err(error) => { return format!("外部生成任务失败:{status},响应读取失败:{error}"); } }; let body_text = String::from_utf8_lossy(&body_bytes).trim().to_string(); if body_text.is_empty() { return format!("外部生成任务失败:{status}"); } if let Ok(body_json) = serde_json::from_str::(&body_text) && let Some(message) = extract_api_error_display_message(&body_json) { return message; } body_text } fn extract_api_error_display_message(body_json: &serde_json::Value) -> Option { let error = body_json.get("error")?; error .get("details") .and_then(|details| { read_trimmed_json_string(details.get("reason")) .or_else(|| read_trimmed_json_string(details.get("message"))) }) .or_else(|| read_trimmed_json_string(error.get("message"))) } fn read_trimmed_json_string(value: Option<&serde_json::Value>) -> Option { value .and_then(serde_json::Value::as_str) .map(str::trim) .filter(|message| !message.is_empty()) .map(ToOwned::to_owned) } #[cfg(any())] async fn fail_queue_job_after_worker_error( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, error: &crate::puzzle::PuzzleExternalGenerationWorkerError, message: &str, ) -> Result<(), String> { if error.should_fail_queue_job() { return fail_job(state, worker_id, job, message.to_string()).await; } warn!( job_id = job.job_id, job_kind = job.job_kind, "external generation worker 业务失败态尚未写回,保留任务租约等待后续重试" ); Ok(()) } fn worker_request_context( job: &ExternalGenerationJobRecord, provider_deadline: Instant, ) -> RequestContext { RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ) .with_external_call_deadline(provider_deadline) } struct ParsedEditorImageEditWorkerPayload { request: EditorImageEditRequest, source: Option, } #[derive(serde::Deserialize)] #[serde(rename_all = "camelCase")] struct LegacyEditorImageEditWorkerPayload { prompt: String, source_image_src: String, size: Option, model: Option, aspect_ratio: Option, image_size: Option, reference_image_srcs: Option>, project_id: Option, generation_inputs: Option, asset_folder_id: Option, asset_label: Option, source_resource_id: Option, target_layer_id: Option, canvas_completion: Option, } fn parse_editor_image_edit_worker_payload( job: &ExternalGenerationJobRecord, ) -> Result { if let Ok(payload) = serde_json::from_str::(job.request_payload_json.as_str()) { if payload.version != EDITOR_IMAGE_EDIT_QUEUE_PAYLOAD_VERSION { return Err(format!( "图片画布改图任务载荷版本不受支持:{}", payload.version )); } return Ok(ParsedEditorImageEditWorkerPayload { request: payload.request, source: Some(payload.source), }); } let legacy: LegacyEditorImageEditWorkerPayload = serde_json::from_str(job.request_payload_json.as_str()) .map_err(|error| format!("图片画布改图任务参数解析失败:{error}"))?; let source_reference_id = legacy .source_resource_id .as_deref() .map(str::trim) .filter(|value| !value.is_empty()) .unwrap_or_else(|| legacy.source_image_src.trim()) .to_string(); if source_reference_id.is_empty() { return Err("历史图片画布改图任务缺少可迁移的业务引用 ID".to_string()); } Ok(ParsedEditorImageEditWorkerPayload { request: EditorImageEditRequest { prompt: legacy.prompt, source_reference_id, size: legacy.size, model: legacy.model, aspect_ratio: legacy.aspect_ratio, image_size: legacy.image_size, reference_image_srcs: legacy.reference_image_srcs, project_id: legacy.project_id, generation_inputs: legacy.generation_inputs, asset_folder_id: legacy.asset_folder_id, asset_label: legacy.asset_label, target_layer_id: legacy.target_layer_id, canvas_completion: legacy.canvas_completion, }, source: None, }) } fn editor_generation_worker_caller( worker_id: &str, job: &ExternalGenerationJobRecord, ) -> Result { Ok(EditorGenerationCaller { owner_user_id: job.owner_user_id.clone(), audit_subject_user_id: Some(job.owner_user_id.clone()), audit_project_id: Some(job.source_entity_id.clone()), phase_reporter: Some(editor_generation_phase_reporter(worker_id, job)?), operation: Some(EditorGenerationOperationContext { operation_kind: job.job_kind.clone(), operation_id: job.job_id.clone(), operation_fingerprint: shared_contracts::editor_generation::editor_generation_request_fingerprint( job.job_kind.as_str(), job.request_payload_json.as_str(), ), worker_id: Some(worker_id.to_string()), lease_token: Some(require_job_lease_token(job)?), queue_result_context: Some(EditorGenerationQueueResultContext::from_job(job)), }), }) } fn editor_generation_phase_reporter( worker_id: &str, job: &ExternalGenerationJobRecord, ) -> Result { Ok(EditorGenerationPhaseReporter::new( job.job_id.clone(), worker_id.to_string(), require_job_lease_token(job)?, )) } async fn complete_editor_generation_job( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, response: Value, ) -> Result<(), String> { complete_job( state, worker_id, job, Some(editor_generation_result_payload_json(job, &response)), ) .await } fn editor_generation_result_payload_json( job: &ExternalGenerationJobRecord, response: &Value, ) -> String { let mut payload = json!({ "sourceModule": job.source_module.clone(), "sourceEntityId": job.source_entity_id.clone(), }); if is_editor_agent_generation_job(job) && let Some(object) = payload.as_object_mut() { // The Agent needs this compact result to restore its tool-call card. Other jobs keep // master's metadata-only completion payload to avoid turning the queue into an asset API. object.insert( // TODO extract const "editor-agent-tool-call-result".to_string(), compact_editor_generation_result(response.clone()), ); } if is_result_recovery_generation_job(job) && let Some(object) = payload.as_object_mut() { object.insert( "result".to_string(), compact_external_api_generation_result(response.clone()), ); } if let Some(warning) = extract_editor_generation_warning(response) && let Some(object) = payload.as_object_mut() { object.insert("warning".to_string(), warning); } payload.to_string() } fn is_result_recovery_generation_job(job: &ExternalGenerationJobRecord) -> bool { let dedupe_key = job.dedupe_key.trim(); dedupe_key.starts_with("external-api-generation:") || dedupe_key.starts_with(&format!("{GAME_CREATOR_CLIENT_GENERATION_DEDUPE_PREFIX}:")) } fn is_editor_agent_generation_job(job: &ExternalGenerationJobRecord) -> bool { job.dedupe_key.trim().starts_with("editor-agent:") } fn compact_external_api_generation_result(result: Value) -> Value { let mut result = result.get("data").cloned().unwrap_or(result); let Some(object) = result.as_object_mut() else { return Value::Null; }; object.retain(|key, _| { matches!( key.as_str(), "ok" | "imageSrc" | "videoSrc" | "audioSrc" | "previewVideoPath" | "thumbnailSrc" | "objectKey" | "assetObjectId" | "width" | "height" | "sourceType" | "model" | "taskId" | "durationSeconds" | "loop" | "resolution" | "priceMudPoints" | "audioKind" | "spritesheetImageSrc" | "spritesheetWidth" | "spritesheetHeight" | "iconImageSrcs" | "sliceLayout" | "frames" | "frameCount" | "frameWidth" | "frameHeight" | "fps" | "resource" | "asset" | "spritesheetResource" | "spritesheetAsset" | "warning" | "sliceWarning" ) }); for field in ["resource", "spritesheetResource"] { if let Some(resource) = object.get_mut(field).and_then(Value::as_object_mut) { compact_external_generation_resource(resource); } } for field in ["asset", "spritesheetAsset"] { if let Some(asset) = object.get_mut(field).and_then(Value::as_object_mut) { compact_external_generation_asset(asset); } } if let Some(icons) = object .get_mut("iconImageSrcs") .and_then(Value::as_array_mut) { for icon in icons { let Some(icon) = icon.as_object_mut() else { continue; }; icon.retain(|key, _| { matches!( key.as_str(), "name" | "imageSrc" | "objectKey" | "width" | "height" | "resource" | "asset" ) }); if let Some(resource) = icon.get_mut("resource").and_then(Value::as_object_mut) { compact_external_generation_resource(resource); } if let Some(asset) = icon.get_mut("asset").and_then(Value::as_object_mut) { compact_external_generation_asset(asset); } remove_unstable_external_generation_media_fields(icon); } } if let Some(frames) = object.get_mut("frames").and_then(Value::as_array_mut) { for frame in frames { let Some(frame) = frame.as_object_mut() else { continue; }; frame.retain(|key, _| { matches!( key.as_str(), "frameIndex" | "imageSrc" | "objectKey" | "assetObjectId" | "width" | "height" ) }); remove_unstable_external_generation_media_fields(frame); } } for field in ["warning", "sliceWarning"] { if let Some(warning) = object.get_mut(field).and_then(Value::as_object_mut) && let Some(reason) = warning.get_mut("reason") && let Some(value) = reason.as_str() { *reason = Value::String(normalize_editor_generation_warning_reason(value)); } } remove_unstable_external_generation_media_fields(object); result } fn compact_external_generation_resource(resource: &mut serde_json::Map) { resource.retain(|key, _| { matches!( key.as_str(), "resourceId" | "projectId" | "objectKey" | "assetObjectId" | "imageSrc" | "width" | "height" | "sourceType" | "assetKind" | "taskId" | "sourceResourceId" ) }); remove_unstable_external_generation_media_fields(resource); } fn compact_external_generation_asset(asset: &mut serde_json::Map) { asset.retain(|key, _| { matches!( key.as_str(), "assetId" | "folderId" | "objectKey" | "assetObjectId" | "imageSrc" | "thumbnailSrc" | "width" | "height" | "sourceType" | "assetKind" | "taskId" ) }); remove_unstable_external_generation_media_fields(asset); } fn remove_unstable_external_generation_media_fields(object: &mut serde_json::Map) { object.retain(|key, value| { if !matches!( key.as_str(), "imageSrc" | "videoSrc" | "audioSrc" | "previewVideoPath" | "thumbnailSrc" | "spritesheetImageSrc" ) { return true; } value .as_str() .is_some_and(is_stable_external_generation_media_reference) }); } fn is_stable_external_generation_media_reference(value: &str) -> bool { let value = value.trim(); !value.is_empty() && value.starts_with('/') && !value.starts_with("//") && !value.contains('?') && !value.contains('#') && !value.to_ascii_lowercase().starts_with("data:") && !value.to_ascii_lowercase().starts_with("blob:") && !value.to_ascii_lowercase().starts_with("http://") && !value.to_ascii_lowercase().starts_with("https://") } fn is_editor_internal_processing_model(model: &str) -> bool { matches!( model.trim().to_ascii_lowercase().as_str(), "anime-seg" | "bgfilter complex" | "birefnet" | "connected-components" | "screen-color-keying" | "segment-common-image" ) } fn extract_editor_generation_warning_fields( warning: Option<&Value>, is_slice_warning: bool, ) -> Option<(String, String)> { let warning = warning?; let code = warning.get("code")?.as_str()?.trim(); let reason = warning.get("reason")?.as_str()?.trim(); if code.is_empty() || reason.is_empty() { return None; } let reason = if is_slice_warning { format!("{EDITOR_GENERATION_SLICE_WARNING_PREFIX}{reason}") } else { reason.to_string() }; Some((code.to_string(), reason)) } fn extract_editor_generation_warning(response: &Value) -> Option { let data = response.get("data").unwrap_or(response); // 中文注释:风格归一化和像素规整产生的通用 warning 可以与 sliceWarning 并存。 // 队列结果只有一个有界 warning 字段,因此按与 inline 响应相同的策略归一: // code 不同时收敛为 multiple-generation-warnings,reason 按“通用在前、拆分在后” // 顺序拼接,再交给既有上界收敛,不允许其中任何一条被静默丢弃。 let common = extract_editor_generation_warning_fields(data.get("warning"), false); let slice = extract_editor_generation_warning_fields(data.get("sliceWarning"), true); let (code, reason) = match (common, slice) { (None, None) => return None, (Some(warning), None) | (None, Some(warning)) => warning, (Some((common_code, common_reason)), Some((slice_code, slice_reason))) => { let code = if common_code == slice_code { common_code } else { EDITOR_GENERATION_MULTIPLE_WARNINGS_CODE.to_string() }; (code, format!("{common_reason} {slice_reason}")) } }; let reason = normalize_editor_generation_warning_reason(reason.as_str()); Some(json!({ "code": code, "reason": reason, })) } fn normalize_editor_generation_warning_reason(reason: &str) -> String { let normalized = reason.to_ascii_lowercase(); if normalized.contains("data:") || normalized.contains("blob:") || normalized.contains("http://") || normalized.contains("https://") || normalized.contains("x-amz-") || normalized.contains("signature=") { return EDITOR_GENERATION_WARNING_REDACTED_MESSAGE.to_string(); } let mut chars = reason.chars(); let mut bounded = chars .by_ref() .take(MAX_EDITOR_GENERATION_WARNING_CHARS) .collect::(); if chars.next().is_some() { bounded.push('…'); } bounded } async fn complete_job( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, result_payload_json: Option, ) -> Result<(), String> { state .spacetime_client() .complete_external_generation_job(ExternalGenerationJobCompleteRecordInput { job_id: job.job_id.clone(), worker_id: worker_id.to_string(), lease_token: require_job_lease_token(job)?, result_payload_json, completed_at_micros: current_utc_micros(), }) .await .map(|_| ()) .map_err(|error| error.to_string()) } async fn fail_job( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, error_message: String, ) -> Result<(), String> { let now_micros = current_utc_micros(); state .spacetime_client() .fail_external_generation_job(ExternalGenerationJobFailRecordInput { job_id: job.job_id.clone(), worker_id: worker_id.to_string(), lease_token: require_job_lease_token(job)?, error_message, retry_after_micros: now_micros.saturating_add(60_000_000), failed_at_micros: now_micros, // 退款由 SpacetimeDB 在验证 lease 后与失败终态同事务结算。 refund_ledger_id: None, }) .await .map(|_| ()) .map_err(|error| error.to_string()) } async fn renew_job_lease( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, lease: Duration, ) -> Result<(), String> { let now_micros = current_utc_micros(); state .spacetime_client() .renew_external_generation_job_lease(ExternalGenerationJobRenewLeaseRecordInput { job_id: job.job_id.clone(), worker_id: worker_id.to_string(), lease_token: require_job_lease_token(job)?, lease_expires_at_micros: now_micros.saturating_add(duration_micros_i64(lease)), renewed_at_micros: now_micros, }) .await .map(|_| ()) .map_err(|error| error.to_string()) } fn require_job_lease_token(job: &ExternalGenerationJobRecord) -> Result { job.lease_token .as_deref() .map(str::trim) .filter(|value| !value.is_empty()) .map(ToOwned::to_owned) .ok_or_else(|| format!("external_generation_job {} 缺少 lease token", job.job_id)) } #[cfg(any())] fn build_external_generation_write_lease_guard( worker_id: &str, job: &ExternalGenerationJobRecord, ) -> Result { Ok(ExternalGenerationWriteLeaseGuard::from_claimed_job( job.job_id.clone(), worker_id.to_string(), require_job_lease_token(job)?, )) } fn duration_micros_i64(duration: Duration) -> i64 { duration.as_micros().min(i64::MAX as u128) as i64 } fn external_generation_billing_price_mud_points( job: &ExternalGenerationJobRecord, ) -> Result { u32::try_from(job.price_mud_points).map_err(|_| { format!( "外部生成任务 {} 的入队价格 price_mud_points={} 超出 worker 支持范围 0..={}", job.job_id, job.price_mud_points, u32::MAX ) }) } fn external_generation_billing_claim_attempt( job: &ExternalGenerationJobRecord, ) -> Result { if job.attempt == 0 { return Err(format!( "外部生成任务 {} 缺少有效的 claim attempt", job.job_id )); } Ok(job.attempt) } fn external_generation_worker_heartbeat_interval(lease: Duration) -> Duration { let heartbeat_millis = (lease.as_millis() / 3).clamp(250, 30_000) as u64; Duration::from_millis(heartbeat_millis) } fn external_generation_worker_deadlines( started_at: Instant, job_timeout: Duration, ) -> (Instant, Instant) { let terminal_write_reserve = EXTERNAL_GENERATION_WORKER_TERMINAL_WRITE_RESERVE.min(job_timeout / 2); let Some(job_deadline) = started_at.checked_add(job_timeout) else { return (started_at, started_at); }; let provider_deadline = started_at .checked_add(job_timeout.saturating_sub(terminal_write_reserve)) .unwrap_or(started_at); (job_deadline, provider_deadline) } fn external_generation_worker_job_timeout(config: &AppConfig, job_kind: &str) -> Duration { match job_kind { EDITOR_IMAGE_GENERATION_JOB_KIND | EDITOR_IMAGE_EDIT_JOB_KIND | EDITOR_ICON_SPEC_GENERATION_JOB_KIND | EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND | EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND | EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND | EDITOR_VIDEO_GENERATION_JOB_KIND => config.external_generation_worker_long_job_timeout, _ => config.external_generation_worker_job_timeout, } } fn external_generation_worker_timeout_message( job: &ExternalGenerationJobRecord, timeout: Duration, ) -> String { format!( "外部生成任务 {}({})超过 worker 执行预算 {} 秒,已停止当前尝试", job.job_id, job.job_kind, timeout.as_secs() ) } fn current_utc_micros() -> i64 { offset_datetime_to_unix_micros(time::OffsetDateTime::now_utc()) } #[cfg(test)] mod tests { use super::*; #[cfg(any())] #[test] fn worker_write_guard_uses_claimed_job_lease_token() { let job = external_generation_job_record_fixture(Some("lease-1")); let guard = build_external_generation_write_lease_guard("worker-a", &job) .expect("guard should build"); assert_eq!(guard.job_id.as_deref(), Some("extgen-1")); assert_eq!(guard.worker_id.as_deref(), Some("worker-a")); assert_eq!(guard.lease_token.as_deref(), Some("lease-1")); } #[test] fn editor_worker_caller_carries_phase_reporter() { let job = external_generation_job_record_fixture(Some("lease-1")); let caller = editor_generation_worker_caller("worker-a", &job) .expect("worker caller should include claimed job phase reporter"); assert!(caller.phase_reporter.is_some()); } #[test] fn editor_image_edit_worker_parses_versioned_snapshot_payload() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.job_kind = EDITOR_IMAGE_EDIT_JOB_KIND.to_string(); job.request_payload_json = json!({ "version": 1, "request": { "prompt": "改成蓝色", "sourceReferenceId": "resource-source" }, "source": { "referenceKind": "project-resource", "sourceReferenceId": "resource-source", "resourceId": "resource-source", "assetId": null, "assetObjectId": "object-source", "bucket": "bucket", "objectKey": "generated/source.png", "sourceAssetKind": "character", "effectiveAssetKind": "character", "mediaType": "image", "targetResourceId": null } }) .to_string(); let payload = parse_editor_image_edit_worker_payload(&job) .expect("versioned image edit payload should parse"); assert_eq!(payload.request.source_reference_id, "resource-source"); assert_eq!( payload.source.expect("snapshot should exist").object_key, "generated/source.png" ); } #[test] fn legacy_editor_image_edit_payload_migrates_only_business_ids() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.job_kind = EDITOR_IMAGE_EDIT_JOB_KIND.to_string(); job.request_payload_json = json!({ "prompt": "改成蓝色", "sourceImageSrc": "raw/generated/source.png", "sourceResourceId": "resource-source", "projectId": "project-1", "assetKind": "editor_agent_edit_image", "generationInputs": { "source": "editor-agent" } }) .to_string(); let legacy_agent_payload = parse_editor_image_edit_worker_payload(&job) .expect("历史 Agent 图片编辑 payload 应由 worker 解析"); assert_eq!( legacy_agent_payload.request.source_reference_id, "resource-source" ); assert!(legacy_agent_payload.source.is_none()); job.request_payload_json = json!({ "prompt": "改成蓝色", "sourceImageSrc": "raw/generated/source.png", "assetKind": "character" }) .to_string(); let raw_only_payload = parse_editor_image_edit_worker_payload(&job) .expect("raw-only legacy payload should migrate without object-key lookup"); assert_eq!( raw_only_payload.request.source_reference_id, "raw/generated/source.png" ); assert!(raw_only_payload.source.is_none()); } #[cfg(any())] #[test] fn worker_write_guard_requires_claimed_job_lease_token() { let job = external_generation_job_record_fixture(None); let error = build_external_generation_write_lease_guard("worker-a", &job) .expect_err("missing token should fail"); assert!(error.contains("缺少 lease token")); } #[test] fn worker_billing_price_rejects_values_above_u32() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.price_mud_points = u64::from(u32::MAX); assert_eq!( external_generation_billing_price_mud_points(&job), Ok(u32::MAX) ); job.price_mud_points += 1; let error = external_generation_billing_price_mud_points(&job) .expect_err("out-of-range queued price should fail"); assert!(error.contains("extgen-1")); assert!(error.contains("price_mud_points=4294967296")); assert!(error.contains("4294967295")); } #[test] fn worker_billing_requires_claim_attempt() { let mut job = external_generation_job_record_fixture(Some("lease-1")); assert_eq!(external_generation_billing_claim_attempt(&job), Ok(1)); job.attempt = 0; let error = external_generation_billing_claim_attempt(&job) .expect_err("unclaimed job should not enter billing context"); assert!(error.contains("extgen-1")); assert!(error.contains("claim attempt")); } #[tokio::test] async fn worker_response_error_message_prefers_api_error_details() { let response = axum::http::Response::builder() .status(axum::http::StatusCode::BAD_GATEWAY) .body(axum::body::Body::from( json!({ "ok": false, "data": null, "error": { "code": "UPSTREAM_ERROR", "message": "上游服务请求失败", "details": { "provider": "vector-engine", "reason": "VECTOR_ENGINE_API_KEY 未配置", "message": "提交编辑器音效任务失败:missing field sound" } }, "meta": {} }) .to_string(), )) .expect("response should build"); let message = response_error_message(response).await; assert_eq!(message, "VECTOR_ENGINE_API_KEY 未配置"); } #[tokio::test] async fn worker_heartbeat_wait_does_not_stop_work_progress() { let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1)); let work_started = std::sync::Arc::new(tokio::sync::Notify::new()); let work_connection = connection.clone(); let work_started_signal = work_started.clone(); let work = async move { let permit = work_connection .acquire_owned() .await .expect("work should acquire the only connection"); work_started_signal.notify_one(); tokio::time::sleep(Duration::from_millis(25)).await; drop(permit); Ok(()) }; let heartbeat_connection = connection.clone(); let heartbeat = async move { work_started.notified().await; let _permit = heartbeat_connection .acquire_owned() .await .expect("heartbeat should acquire after work releases the connection"); std::future::pending::>().await }; let outcome = await_external_generation_job_execution( work, heartbeat, Instant::now() + Duration::from_secs(1), ) .await; assert_eq!( outcome, ExternalGenerationJobExecutionOutcome::Finished(Ok(())) ); } #[tokio::test] async fn worker_deadline_drops_local_work_future_before_returning() { let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1)); let work_connection = connection.clone(); let work = async move { let _permit = work_connection .acquire_owned() .await .expect("work should acquire the only connection"); std::future::pending::>().await }; let heartbeat = std::future::pending::>(); let outcome = await_external_generation_job_execution( work, heartbeat, Instant::now() + Duration::from_millis(25), ) .await; assert_eq!(outcome, ExternalGenerationJobExecutionOutcome::TimedOut); assert!( connection.try_acquire().is_ok(), "deadline 返回前应先 drop 传入的 work future 并释放其持有的资源" ); } #[tokio::test] async fn worker_detached_work_keeps_running_after_deadline() { let finished = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); let finished_flag = finished.clone(); let work_handle = tokio::spawn(async move { tokio::time::sleep(Duration::from_millis(20)).await; finished_flag.store(true, std::sync::atomic::Ordering::SeqCst); Ok(()) }); let job = external_generation_job_record_fixture(Some("lease-1")); detach_external_generation_work_until_lease_expiry( work_handle, &job, Duration::from_millis(200), "任务超过执行预算", None, ); tokio::time::sleep(Duration::from_millis(100)).await; assert!( finished.load(std::sync::atomic::Ordering::SeqCst), "超时后在途 work 应继续执行完成,而不是被取消" ); } #[tokio::test] async fn worker_detached_work_keeps_execution_slot_until_finished() { let slots = std::sync::Arc::new(tokio::sync::Semaphore::new(1)); let permit = slots .clone() .acquire_owned() .await .expect("the only execution slot should be available"); let work_handle = tokio::spawn(async { tokio::time::sleep(Duration::from_millis(20)).await; Ok(()) }); let job = external_generation_job_record_fixture(Some("lease-1")); detach_external_generation_work_until_lease_expiry( work_handle, &job, Duration::from_millis(200), "任务超过执行预算", Some(permit), ); assert!( slots.try_acquire().is_err(), "脱管 work 完成前不得重新领取执行容量" ); tokio::time::sleep(Duration::from_millis(100)).await; assert!( slots.try_acquire().is_ok(), "脱管 work 完成后应归还执行容量" ); } #[tokio::test] async fn worker_reaps_completed_tasks_while_queue_remains_busy() { let mut tasks = JoinSet::new(); let completed = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); const TASK_COUNT: usize = 128; for _ in 0..TASK_COUNT { let completed = completed.clone(); tasks.spawn(async move { completed.fetch_add(1, std::sync::atomic::Ordering::SeqCst); }); } while completed.load(std::sync::atomic::Ordering::SeqCst) < TASK_COUNT { tokio::task::yield_now().await; } assert_eq!(tasks.len(), TASK_COUNT); reap_finished_external_generation_worker_tasks(&mut tasks); assert!(tasks.is_empty(), "已完成任务的 JoinHandle 应在每轮被回收"); } #[tokio::test] async fn worker_detached_work_is_aborted_after_lease_arbitration_window() { let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1)); let work_connection = connection.clone(); let work_handle = tokio::spawn(async move { let _permit = work_connection .acquire_owned() .await .expect("work should acquire the only connection"); std::future::pending::>().await }); let job = external_generation_job_record_fixture(Some("lease-1")); detach_external_generation_work_until_lease_expiry( work_handle, &job, Duration::from_millis(10), "任务超过执行预算", None, ); let reacquired = tokio::time::timeout(Duration::from_secs(2), connection.acquire_owned()).await; assert!( reacquired.is_ok(), "超过租约仲裁窗口后应取消仍未结束的 work 并释放其持有的资源" ); } #[test] fn editor_generation_result_payload_keeps_only_lightweight_slice_warning() { let job = external_generation_job_record_fixture(Some("lease-1")); let response = json!({ "spritesheetImageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST", "iconImageSrcs": [{"imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST"}], "sliceWarning": { "code": "insufficient-connected-components", "reason": "连通域数量不足" } }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker 结果应是合法 JSON"); assert_eq!(payload["sourceModule"], json!("editor")); assert_eq!(payload["sourceEntityId"], json!("project-1")); assert_eq!( payload["warning"], json!({ "code": "insufficient-connected-components", "reason": "图集已生成,但自动拆分未完成:连通域数量不足" }) ); assert!(payload.get("spritesheetImageSrc").is_none()); assert!(payload.get("iconImageSrcs").is_none()); assert!(payload.get("editor-agent-tool-call-result").is_none()); } #[test] fn editor_agent_result_payload_keeps_compact_response() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.dedupe_key = "editor-agent:conversation-1:7:generate-image".to_string(); job.request_payload_json = json!({ "generationInputs": { "source": "editor-agent" }, }) .to_string(); let response = json!({ "imageSrc": "/api/assets/object/generated.png", "objectKey": "users/user-1/generated.png", "assetObjectId": "asset-object-1", "width": 1024, "height": 1024, "sourceType": "generated", "prompt": "castle", "actualPrompt": null, "model": "gpt-image-2", "provider": "VectorEngine", "generationInputs": { "fields": { "prompt": "castle" }, "characterAnimation": { "frameCount": 8 }, "screenColorHex": "#CFEFFF", "mattingProvider": "BgFilter", "mattingModel": "birefnet", }, "taskId": "provider-task-1", "resource": { "resourceId": "resource-1", "objectKey": "users/user-1/generated.png", "assetObjectId": "asset-object-1", "imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST", }, "asset": { "assetId": "asset-1" }, "project": { "projectId": "project-1" }, }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker result should be valid JSON"); assert_eq!( payload["editor-agent-tool-call-result"]["imageSrc"], json!("/api/assets/object/generated.png") ); assert_eq!( payload["editor-agent-tool-call-result"]["resource"], json!({ "resourceId": "resource-1", "objectKey": "users/user-1/generated.png", "assetObjectId": "asset-object-1", }) ); assert!( payload["editor-agent-tool-call-result"] .get("asset") .is_none() ); assert!( payload["editor-agent-tool-call-result"] .get("project") .is_none() ); assert_eq!( payload["editor-agent-tool-call-result"]["model"], json!("gpt-image-2") ); assert!( payload["editor-agent-tool-call-result"] .get("provider") .is_none() ); assert_eq!( payload["editor-agent-tool-call-result"]["generationInputs"], json!({ "fields": { "prompt": "castle" }, "characterAnimation": { "frameCount": 8 }, }) ); assert!(!payload.to_string().contains("data:image")); } #[test] fn compact_editor_generation_result_removes_internal_models_case_insensitively() { for model in [ "anime-seg", " BgFilter Complex ", "BIREFNET", "connected-components", "screen-color-keying", "segment-common-image", ] { let compact = compact_editor_generation_result(json!({ "model": model, "provider": "internal-provider", })); assert!(compact.get("model").is_none(), "model={model}"); assert!(compact.get("provider").is_none(), "model={model}"); } } #[test] fn editor_agent_audio_compact_results_remain_reconcileable() { use crate::editor_agent::reconcile_completed_editor_agent_tool_call_for_test; use platform_editor_agent::{ agent::tools::{ generate_background_music::GenerateBackgroundMusicTool, generate_sound_effect::GenerateSoundEffectTool, }, framework::tool::Tool, }; use shared_contracts::editor_agent::{EditorAgentMessage, EditorAgentToolCallStatus}; for (tool_name, prompt, actual_prompt, audio_kind, model, provider) in [ ( GenerateSoundEffectTool::NAME, "按钮点击声", "A short button click", "sound-effect", "eleven_text_to_sound_v2", "elevenlabs", ), ( GenerateBackgroundMusicTool::NAME, "森林背景音乐", "森林背景音乐", "background-music", "chirp-v5", "vectorengine", ), ] { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.dedupe_key = format!("editor-agent:conversation-1:7:{tool_name}"); job.request_payload_json = json!({ "generationInputs": { "source": "editor-agent" }, }) .to_string(); let response = json!({ "ok": true, "audioSrc": "/generated/audio.mp3", "width": 420, "height": 120, "sourceType": "generated", "prompt": prompt, "actualPrompt": actual_prompt, "model": model, "provider": provider, "taskId": "task-1", "priceMudPoints": 5, "audioKind": audio_kind, "durationSeconds": if audio_kind == "sound-effect" { json!(5.25) } else { Value::Null }, "loop": if audio_kind == "sound-effect" { json!(false) } else { Value::Null }, }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker compact payload should serialize"); assert!(payload.get("editor-agent-tool-call-result").is_some()); assert!( payload["editor-agent-tool-call-result"] .get("provider") .is_none() ); let mut message: EditorAgentMessage = serde_json::from_value(json!({ "id": 1, "role": "system", "text": "waiting", "attachments": [], "toolCall": { "toolName": tool_name, "status": "not_completed", "args": { "prompt": prompt }, "displayArgs": { "stringArgs": [], "imageArgs": [], "extras": { "priceMudPoints": 5 } }, "externalJobId": "job-1", "images": [], "audios": [] }, "createdAt": "2026-08-06T00:00:00Z" })) .expect("pending audio Agent message should deserialize"); let payload_json = payload.to_string(); reconcile_completed_editor_agent_tool_call_for_test( &mut message, Some(payload_json.as_str()), ) .expect("worker compact audio result should reconcile"); let tool_call = message.tool_call.expect("tool call should remain present"); assert_eq!(tool_call.status, EditorAgentToolCallStatus::Completed); assert_eq!(tool_call.audios[0].audio_src, "/generated/audio.mp3"); } } #[test] fn editor_agent_spritesheet_result_keeps_all_persisted_slices() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.dedupe_key = "editor-agent:conversation-1:7:generate-icon-spritesheet".to_string(); job.request_payload_json = json!({ "generationInputs": { "source": "editor-agent" }, }) .to_string(); let response = json!({ "spritesheetImageSrc": "/api/assets/object/sheet.png", "spritesheetWidth": 512, "spritesheetHeight": 512, "taskId": "provider-task-1", "spritesheetResource": { "resourceId": "sheet-resource", "objectKey": "users/user-1/sheet.png", "assetObjectId": "sheet-object", }, "iconImageSrcs": [ { "name": "backpack", "imageSrc": "/api/assets/object/backpack.png", "width": 64, "height": 64, "resource": { "resourceId": "icon-resource-1", "objectKey": "users/user-1/backpack.png", "assetObjectId": "icon-object-1", "sourceResourceId": "sheet-resource", "imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST", }, "asset": { "assetId": "icon-asset-1" }, }, { "name": "map", "imageSrc": "/api/assets/object/map.png", "width": 64, "height": 64, "resource": { "resourceId": "icon-resource-2", "objectKey": "users/user-1/map.png", "assetObjectId": "icon-object-2", "sourceResourceId": "sheet-resource", }, }, ], }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker result should be valid JSON"); assert_eq!( payload["editor-agent-tool-call-result"]["iconImageSrcs"] .as_array() .map(Vec::len), Some(2) ); assert_eq!( payload["editor-agent-tool-call-result"]["iconImageSrcs"][0]["resource"], json!({ "resourceId": "icon-resource-1", "objectKey": "users/user-1/backpack.png", "assetObjectId": "icon-object-1", "sourceResourceId": "sheet-resource", }) ); assert!( payload["editor-agent-tool-call-result"]["iconImageSrcs"][0] .get("asset") .is_none() ); assert!(!payload.to_string().contains("data:image")); } #[test] fn editor_generation_result_payload_merges_common_and_slice_warnings() { let job = external_generation_job_record_fixture(Some("lease-1")); let response = json!({ "data": { "imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST", "warning": { "code": "unsupported-image-style", "reason": "不支持的图片风格,已按无风格继续生成。" }, "sliceWarning": { "code": "insufficient-connected-components", "reason": "有效连通域不足" } } }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker 结果应是合法 JSON"); // 中文注释:风格归一化告警与拆分告警可以并存,队列只有一个 warning 字段, // 必须拼接后收敛 code,不能让其中任何一条消失。 assert_eq!( payload["warning"], json!({ "code": "multiple-generation-warnings", "reason": "不支持的图片风格,已按无风格继续生成。 图集已生成,但自动拆分未完成:有效连通域不足" }) ); assert!(payload.get("imageSrc").is_none()); } #[test] fn external_generation_result_payload_keeps_fixed_spritesheet_layout() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.dedupe_key = "external-api-generation:conversation-1:7:icon-spritesheet".to_string(); let response = json!({ "sliceLayout": "grid-2x2", "iconImageSrcs": [ { "name": "素材 1", "imageSrc": "/api/assets/object/one.png" }, { "name": "素材 2", "imageSrc": "/api/assets/object/two.png" }, { "name": "素材 3", "imageSrc": "/api/assets/object/three.png" }, { "name": "素材 4", "imageSrc": "/api/assets/object/four.png" } ] }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker result should be valid JSON"); assert_eq!(payload["result"]["sliceLayout"], json!("grid-2x2")); assert_eq!( payload["result"]["iconImageSrcs"].as_array().map(Vec::len), Some(4) ); } #[test] fn editor_generation_result_payload_keeps_single_warning_untouched() { let job = external_generation_job_record_fixture(Some("lease-1")); let response = json!({ "data": { "warning": { "code": "postprocess-failed-source-preserved", "reason": "生成任务成功,后处理失败。" } } }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker 结果应是合法 JSON"); // 中文注释:透明背景最终失败不会进入拆分,此时仍是单条告警,原样保留。 assert_eq!( payload["warning"], json!({ "code": "postprocess-failed-source-preserved", "reason": "生成任务成功,后处理失败。" }) ); } #[test] fn editor_image_job_completion_is_committed_by_atomic_persistence() { let source = include_str!("external_generation_worker.rs"); let start = source .find("EDITOR_IMAGE_GENERATION_JOB_KIND => {") .expect("editor image worker branch should exist"); let branch_tail = &source[start..]; let end = branch_tail .find("EDITOR_IMAGE_EDIT_JOB_KIND => {") .expect("editor image worker branch end marker should exist"); let branch = &branch_tail[..end]; assert!(branch.contains("Ok(_) => Ok(())")); assert!( !branch.contains("complete_editor_generation_job"), "editor image completion must be part of the atomic result persistence" ); } #[test] fn icon_spec_job_completion_is_committed_by_atomic_persistence() { let source = include_str!("external_generation_worker.rs"); let start = source .find("EDITOR_ICON_SPEC_GENERATION_JOB_KIND => {") .expect("icon spec worker branch should exist"); let branch_tail = &source[start..]; let end = branch_tail .find("EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND => {") .expect("icon spec worker branch end marker should exist"); let branch = &branch_tail[..end]; assert!(branch.contains("Ok(_) => Ok(())")); assert!( !branch.contains("complete_editor_generation_job"), "icon spec completion must be part of the atomic result persistence" ); } #[test] fn editor_generation_result_payload_accepts_envelope_and_redacts_inline_media() { let job = external_generation_job_record_fixture(Some("lease-1")); let response = json!({ "data": { "sliceWarning": { "code": "slice-persistence-failed", "reason": "provider returned data:image/png;base64,AAAA" } } }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("worker 结果应是合法 JSON"); assert_eq!( payload["warning"]["reason"], json!(EDITOR_GENERATION_WARNING_REDACTED_MESSAGE) ); assert!(!payload.to_string().contains("data:image")); } #[test] fn external_api_result_keeps_stable_artifacts_and_removes_unstable_media() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.dedupe_key = "external-api-generation:editor_image_generation:fingerprint".to_string(); let response = json!({ "data": { "imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST", "videoSrc": "blob:https://example.test/video", "audioSrc": "https://cdn.example.test/audio.mp3?X-Amz-Signature=secret", "previewVideoPath": "https://cdn.example.test/stable-looking-but-external.mp4", "thumbnailSrc": "/api/assets/object/thumbnail.png?expires=1&signature=secret", "objectKey": "users/user-1/generated/main.png", "assetObjectId": "asset-object-main", "width": 1024, "height": 1024, "durationSeconds": 7.42, "loop": true, "provider": "internal-provider-must-not-persist", "resource": { "resourceId": "resource-main", "projectId": "project-1", "objectKey": "users/user-1/generated/main.png", "assetObjectId": "asset-object-main", "sourceResourceId": "source-resource-main", "imageSrc": "https://cdn.example.test/main.png?signature=secret", "width": 1024, "height": 1024, "prompt": "不应复制完整资源元数据" }, "asset": { "assetId": "asset-main", "folderId": "folder-1", "objectKey": "users/user-1/generated/main.png", "assetObjectId": "asset-object-main", "imageSrc": "/api/assets/object/main.png", "thumbnailSrc": "https://cdn.example.test/thumb.png?signature=secret", "width": 1024, "height": 1024, "generationInputs": {"private": true} }, "frames": [ { "imageSrc": "/generated/action/frame-01.png", "objectKey": "generated/action/frame-01.png", "assetObjectId": "asset-object-frame-01", "width": 192, "height": 256, "provider": "internal-provider-must-not-persist" }, { "imageSrc": "/generated/action/frame-02.png", "objectKey": "generated/action/frame-02.png", "assetObjectId": "asset-object-frame-02", "width": 192, "height": 256 } ], "project": { "projectId": "project-1", "canvas": {"layers": ["large-layout-must-not-persist"]} }, "warning": { "code": "dimension-restore-fallback", "reason": "已保留 provider 实际输出尺寸。" } }, "meta": { "requestId": "worker-envelope-must-not-persist" } }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("外部生成结果应是合法 JSON"); let result = &payload["result"]; assert!(result.get("project").is_none()); assert!(result.get("provider").is_none()); for unstable_field in [ "imageSrc", "videoSrc", "audioSrc", "previewVideoPath", "thumbnailSrc", ] { assert!( result.get(unstable_field).is_none(), "不稳定媒体字段 {unstable_field} 不得持久化" ); } assert_eq!( result["objectKey"], json!("users/user-1/generated/main.png") ); assert_eq!(result["assetObjectId"], json!("asset-object-main")); assert_eq!(result["durationSeconds"], json!(7.42)); assert_eq!(result["loop"], json!(true)); assert_eq!(result["resource"]["resourceId"], json!("resource-main")); assert_eq!( result["resource"]["sourceResourceId"], json!("source-resource-main") ); assert_eq!( result["resource"]["objectKey"], json!("users/user-1/generated/main.png") ); assert!(result["resource"].get("imageSrc").is_none()); assert!(result["resource"].get("prompt").is_none()); assert_eq!(result["asset"]["assetId"], json!("asset-main")); assert_eq!( result["asset"]["imageSrc"], json!("/api/assets/object/main.png") ); assert!(result["asset"].get("thumbnailSrc").is_none()); assert!(result["asset"].get("generationInputs").is_none()); assert_eq!( result["frames"][0], json!({ "imageSrc": "/generated/action/frame-01.png", "objectKey": "generated/action/frame-01.png", "assetObjectId": "asset-object-frame-01", "width": 192, "height": 256 }) ); assert_eq!( result["warning"], json!({ "code": "dimension-restore-fallback", "reason": "已保留 provider 实际输出尺寸。" }) ); assert_eq!(payload["warning"], result["warning"]); assert!(result.get("prompt").is_none()); assert!(result.get("actualPrompt").is_none()); let serialized = payload.to_string().to_ascii_lowercase(); for forbidden in [ "data:", "blob:", "x-amz-signature", "?signature=", "large-layout", ] { assert!( !serialized.contains(forbidden), "compact result 不应包含 {forbidden}" ); } } #[test] fn game_creator_client_result_keeps_completed_grid_spritesheet_for_recovery() { let mut job = external_generation_job_record_fixture(Some("lease-1")); job.dedupe_key = format!( "{GAME_CREATOR_CLIENT_GENERATION_DEDUPE_PREFIX}:editor_icon_spritesheet_generation:fingerprint" ); job.job_kind = EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND.to_string(); let response = json!({ "data": { "spritesheetImageSrc": "/api/assets/object/core-sheet.png", "spritesheetWidth": 1024, "spritesheetHeight": 1024, "sliceLayout": "grid-2x2", "spritesheetResource": { "resourceId": "sheet-resource-1", "objectKey": "users/user-1/core-sheet.png", "imageSrc": "/api/assets/object/core-sheet.png" }, "spritesheetAsset": { "assetId": "sheet-asset-1", "objectKey": "users/user-1/core-sheet.png", "imageSrc": "/api/assets/object/core-sheet.png" }, "iconImageSrcs": [ {"name": "玩家", "objectKey": "users/user-1/player.png", "imageSrc": "/api/assets/object/player.png"}, {"name": "目标", "objectKey": "users/user-1/targets.png", "imageSrc": "/api/assets/object/targets.png"}, {"name": "场景", "objectKey": "users/user-1/scene.png", "imageSrc": "/api/assets/object/scene.png"}, {"name": "反馈", "objectKey": "users/user-1/feedback.png", "imageSrc": "/api/assets/object/feedback.png"} ] } }); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(&job, &response)) .expect("游戏创作客户端完成结果应持久化为合法 JSON"); assert_eq!(payload["result"]["sliceLayout"], json!("grid-2x2")); assert_eq!( payload["result"]["iconImageSrcs"].as_array().map(Vec::len), Some(4) ); assert_eq!( payload["result"]["spritesheetResource"]["resourceId"], json!("sheet-resource-1") ); assert_eq!( payload["result"]["spritesheetAsset"]["assetId"], json!("sheet-asset-1") ); assert_eq!( payload["result"]["iconImageSrcs"][0]["objectKey"], json!("users/user-1/player.png") ); assert!(!payload.to_string().contains("prompt")); } #[test] fn non_external_job_does_not_publish_query_result() { let job = external_generation_job_record_fixture(Some("lease-1")); let payload: Value = serde_json::from_str(&editor_generation_result_payload_json( &job, &json!({ "objectKey": "users/user-1/generated/main.png", "resource": {"resourceId": "resource-main"} }), )) .expect("普通编辑器任务结果应为合法 JSON"); assert!(payload.get("result").is_none()); } #[test] fn worker_job_timeout_uses_long_budget_for_image_and_video_jobs() { let config = AppConfig { external_generation_worker_job_timeout: Duration::from_secs(60), external_generation_worker_long_job_timeout: Duration::from_secs(600), ..AppConfig::default() }; assert_eq!( external_generation_worker_job_timeout(&config, EDITOR_IMAGE_GENERATION_JOB_KIND), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout(&config, EDITOR_IMAGE_EDIT_JOB_KIND), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout(&config, EDITOR_ICON_SPEC_GENERATION_JOB_KIND), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout( &config, EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND ), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout( &config, EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND ), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout( &config, EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND ), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout(&config, EDITOR_VIDEO_GENERATION_JOB_KIND), Duration::from_secs(600) ); assert_eq!( external_generation_worker_job_timeout(&config, EDITOR_BACKGROUND_REMOVAL_JOB_KIND), Duration::from_secs(60) ); assert_eq!( external_generation_worker_job_timeout( &config, EDITOR_SOUND_EFFECT_GENERATION_JOB_KIND ), Duration::from_secs(60) ); } #[test] fn worker_provider_deadline_reserves_terminal_write_budget() { let started_at = Instant::now(); let (job_deadline, provider_deadline) = external_generation_worker_deadlines(started_at, Duration::from_secs(1_800)); assert_eq!( job_deadline.duration_since(started_at), Duration::from_secs(1_800) ); assert_eq!( provider_deadline.duration_since(started_at), Duration::from_secs(1_740) ); assert_eq!( job_deadline.duration_since(provider_deadline), Duration::from_secs(60) ); let (short_job_deadline, short_provider_deadline) = external_generation_worker_deadlines(started_at, Duration::from_secs(30)); assert_eq!( short_job_deadline.duration_since(started_at), Duration::from_secs(30) ); assert_eq!( short_provider_deadline.duration_since(started_at), Duration::from_secs(15) ); } #[test] fn worker_request_context_carries_provider_deadline() { let job = external_generation_job_record_fixture(Some("lease-1")); let provider_deadline = Instant::now() + Duration::from_secs(30); let request_context = worker_request_context(&job, provider_deadline); assert_eq!( request_context.external_call_deadline(), Some(provider_deadline) ); } #[test] fn worker_timeout_message_mentions_job_and_budget() { let job = external_generation_job_record_fixture(Some("lease-1")); let message = external_generation_worker_timeout_message(&job, Duration::from_secs(90)); assert!(message.contains("extgen-1")); assert!(message.contains(EDITOR_IMAGE_GENERATION_JOB_KIND)); assert!(message.contains("90 秒")); } #[test] fn terminal_worker_failure_delegates_refund_to_the_fenced_failure_transaction() { let source = include_str!("external_generation_worker.rs"); let body = source .split_once("async fn fail_job(") .and_then(|(_, tail)| { tail.split_once("async fn renew_job_lease(") .map(|(body, _)| body) }) .expect("fail_job helper"); assert!(!body.contains("settle_current_external_generation_attempt_refund(")); assert!(body.contains(".fail_external_generation_job(")); assert!(body.contains("refund_ledger_id: None")); } fn external_generation_job_record_fixture( lease_token: Option<&str>, ) -> ExternalGenerationJobRecord { ExternalGenerationJobRecord { job_id: "extgen-1".to_string(), dedupe_key: "editor:image-generation:project-1:extgen-1".to_string(), job_kind: EDITOR_IMAGE_GENERATION_JOB_KIND.to_string(), owner_user_id: "user-1".to_string(), source_module: "editor".to_string(), source_entity_id: "project-1".to_string(), request_label: "画布图片生成".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-06-03T00:00:00Z".to_string()), available_at: "2026-06-03T00:00:00Z".to_string(), result_payload_json: None, created_at: "2026-06-03T00:00:00Z".to_string(), started_at: Some("2026-06-03T00:00:00Z".to_string()), completed_at: None, updated_at: "2026-06-03T00:00:00Z".to_string(), updated_at_micros: 1_780_444_800_000_000, lease_token: lease_token.map(ToOwned::to_owned), price_mud_points: 2, refund_ledger_id: None, notification_acknowledged_at: None, notification_acknowledged_at_micros: None, phase: Some("generating".to_string()), } } }