use std::{future::Future, io, pin::Pin, time::Duration}; use axum::{Json, extract::FromRef}; use serde_json::{Value, json}; use shared_kernel::offset_datetime_to_unix_micros; use spacetime_client::{ ExternalGenerationJobClaimRecordInput, ExternalGenerationJobCompleteRecordInput, ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord, ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription, }; use tokio::{ task::JoinSet, time::{Instant, sleep}, }; use tracing::{error, info, warn}; const MAX_EDITOR_GENERATION_WARNING_CHARS: usize = 2_048; 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_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::{ EditorBackgroundRemovalRequest, EditorGenerationCaller, EditorIconSpritesheetGenerationRequest, EditorImageEditRequest, EditorImageGenerationRequest, EditorUiDesignAssetExtractionRequest, edit_editor_image_for_owner, extract_editor_ui_design_assets_for_owner, generate_editor_icon_spritesheet_for_owner, generate_editor_image_for_owner, remove_editor_image_background_for_owner, }, jump_hop::{ JUMP_HOP_COMPILE_DRAFT_JOB_KIND, JumpHopCompileDraftWorkerPayload, execute_jump_hop_compile_draft_worker_job, }, 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, }, request_context::RequestContext, state::{AppState, PuzzleApiState}, vector_engine_audio_generation::{ generate_editor_background_music_for_owner, generate_editor_sound_effect_for_owner, }, wooden_fish::{ WOODEN_FISH_GENERATE_IMAGE_ASSETS_JOB_KIND, WoodenFishGenerateImageAssetsWorkerPayload, execute_wooden_fish_generate_image_assets_worker_job, }, }; pub(crate) const PUZZLE_COMPILE_DRAFT_JOB_KIND: &str = "puzzle_compile_draft"; pub(crate) const PUZZLE_GENERATE_IMAGES_JOB_KIND: &str = "puzzle_generate_images"; 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; 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 { ensure_external_generation_queue_wake_subscription(&state, &mut queue_wake).await; while tasks.len() >= concurrency { if await_worker_task_or_shutdown(&mut tasks, &mut shutdown).await { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); } } let available = concurrency.saturating_sub(tasks.len()).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(); tasks.spawn(async move { if let Err(error) = process_external_generation_job(state, worker_id, lease, job).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"); } } async fn await_worker_task_or_shutdown( tasks: &mut JoinSet<()>, shutdown: &mut ExternalGenerationShutdownSignal, ) -> bool { tokio::select! { _ = shutdown.as_mut() => true, _ = await_worker_task(tasks) => false, } } 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, ) -> 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 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); } }; let work = 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()), ); tokio::pin!(work); let heartbeat = sleep(heartbeat_interval); tokio::pin!(heartbeat); let job_deadline = sleep(job_timeout); tokio::pin!(job_deadline); loop { tokio::select! { biased; result = &mut work => return result, _ = &mut job_deadline => { 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 槽位" ); fail_job(&state, &worker_id, &job, message.clone()).await?; return Err(message); } _ = &mut heartbeat => { renew_job_lease(&state, &worker_id, &job, lease).await?; heartbeat.as_mut().reset(Instant::now() + heartbeat_interval); } } } } async fn process_external_generation_job_once( state: AppState, worker_id: String, job: ExternalGenerationJobRecord, ) -> Result<(), String> { match job.job_kind.as_str() { 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) } } } 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) } } } 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) } } } JUMP_HOP_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_jump_hop_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) } } } 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) } } } 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); match generate_editor_image_for_owner( &state, &request_context, editor_generation_worker_caller(&job), payload, ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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 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); match edit_editor_image_for_owner( &state, &request_context, editor_generation_worker_caller(&job), payload, ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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); match remove_editor_image_background_for_owner( &state, &request_context, editor_generation_worker_caller(&job), payload, ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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); match generate_editor_icon_spritesheet_for_owner( &state, &request_context, editor_generation_worker_caller(&job), payload, ) .await { Ok(response) => { complete_editor_generation_job_with_response( &state, &worker_id, &job, &response.0, ) .await } 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); match extract_editor_ui_design_assets_for_owner( &state, &request_context, editor_generation_worker_caller(&job), payload, ) .await { Ok(response) => { complete_editor_generation_job_with_response( &state, &worker_id, &job, &response.0, ) .await } 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); match generate_editor_character_animation_for_owner( state.clone(), request_context, job.owner_user_id.clone(), Ok(Json(payload)), ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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); match generate_editor_video_for_owner( state.clone(), request_context, job.owner_user_id.clone(), Ok(Json(payload)), ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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); match generate_editor_sound_effect_for_owner( state.clone(), request_context, job.owner_user_id.clone(), Ok(Json(payload)), ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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); match generate_editor_background_music_for_owner( state.clone(), request_context, job.owner_user_id.clone(), Ok(Json(payload)), ) .await { Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await, 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) } 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) -> RequestContext { RequestContext::new( format!("external-generation-worker-{}", job.job_id), format!("external-generation-worker {}", job.job_kind), std::time::Duration::ZERO, false, ) } fn editor_generation_worker_caller(job: &ExternalGenerationJobRecord) -> EditorGenerationCaller { 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()), } } async fn complete_editor_generation_job( state: &AppState, worker_id: &str, job: &ExternalGenerationJobRecord, ) -> Result<(), String> { complete_job( state, worker_id, job, Some( json!({ "sourceModule": job.source_module.clone(), "sourceEntityId": job.source_entity_id.clone(), }) .to_string(), ), ) .await } async fn complete_editor_generation_job_with_response( 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 let Some(warning) = extract_editor_generation_slice_warning(response) && let Some(object) = payload.as_object_mut() { object.insert("warning".to_string(), warning); } payload.to_string() } fn extract_editor_generation_slice_warning(response: &Value) -> Option { let data = response.get("data").unwrap_or(response); let warning = data.get("sliceWarning")?; 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 = normalize_editor_generation_warning_reason(reason); 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:") { 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, 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)) } 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_job_timeout(config: &AppConfig, job_kind: &str) -> Duration { match 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::*; #[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 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 未配置"); } #[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!("puzzle")); assert_eq!(payload["sourceEntityId"], json!("session-1:puzzle-level-1")); assert_eq!( payload["warning"], json!({ "code": "insufficient-connected-components", "reason": "连通域数量不足" }) ); assert!(payload.get("spritesheetImageSrc").is_none()); assert!(payload.get("iconImageSrcs").is_none()); } #[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 worker_job_timeout_uses_long_budget_for_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(60) ); 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) ); } #[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(PUZZLE_GENERATE_IMAGES_JOB_KIND)); assert!(message.contains("90 秒")); } fn external_generation_job_record_fixture( lease_token: Option<&str>, ) -> ExternalGenerationJobRecord { ExternalGenerationJobRecord { job_id: "extgen-1".to_string(), dedupe_key: "puzzle:generate_puzzle_images:session-1:extgen-1".to_string(), job_kind: PUZZLE_GENERATE_IMAGES_JOB_KIND.to_string(), owner_user_id: "user-1".to_string(), source_module: "puzzle".to_string(), source_entity_id: "session-1:puzzle-level-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, } } }