use std::{ future::Future, io, pin::Pin, time::{Duration, Instant}, }; use axum::Json; use serde_json::Value; use shared_kernel::offset_datetime_to_unix_micros; use spacetime_client::{ ExternalGenerationJobClaimRecordInput, ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord, ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription, }; use tokio::{ sync::{OwnedSemaphorePermit, Semaphore}, task::{JoinHandle, JoinSet}, time::sleep, }; use tracing::{error, info, warn}; // provider 必须先结束,给失败审计、计费结算和队列终态写回保留有效 lease 内的收尾窗口。 const EXTERNAL_GENERATION_WORKER_TERMINAL_WRITE_RESERVE: Duration = Duration::from_secs(60); 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, }, editor_project::{ EDITOR_IMAGE_EDIT_QUEUE_PAYLOAD_VERSION, EditorBackgroundRemovalRequest, EditorGenerationCaller, EditorGenerationOperationContext, EditorGenerationPhaseReporter, EditorGenerationQueueResultContext, EditorImageEditQueuePayload, EditorImageEditRequest, EditorImageEditResolvedSource, EditorImageGenerationRequest, EditorUiDesignAssetExtractionRequest, 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 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::*; use serde_json::json; #[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": "tiantoken", "reason": "TIANTOKEN_API_KEY 未配置", "message": "提交编辑器音效任务失败:missing field sound" } }, "meta": {} }) .to_string(), )) .expect("response should build"); let message = response_error_message(response).await; assert_eq!(message, "TIANTOKEN_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_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 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()), } } }