Files
Genarrative/server-rs/crates/api-server/src/external_generation_worker.rs
T
kdletters 4e7bcb8f24 补齐生成中间产物持久化与拆图提示
图标和UI图集拆分失败时保留主图集并在内联、队列及刷新路径提示
角色、图片修改、图标UI和角色动作链路把可恢复中间产物写入项目与素材库
外部生成摘要新增独立告警字段并清理成功任务的历史错误
同步共享契约、SpacetimeDB迁移与绑定、OpenAPI、设计文档和回归测试
2026-07-13 21:41:46 +08:00

1454 lines
53 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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<ExternalGenerationQueueWakeSubscription>,
) {
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<Box<dyn Future<Output = ()> + 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<ExternalGenerationQueueWakeSubscription>,
sleeper: impl Future<Output = ()>,
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<ExternalGenerationQueueWakeSubscription>,
) {
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::<PuzzleCompileDraftWorkerPayload>(
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::<PuzzleGenerateImagesWorkerPayload>(
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::<PuzzleGenerateUiBackgroundWorkerPayload>(
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::<JumpHopCompileDraftWorkerPayload>(
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::<PuzzleClearCompileDraftWorkerPayload>(
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::<WoodenFishGenerateImageAssetsWorkerPayload>(
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::<EditorImageGenerationRequest>(
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::<EditorImageEditRequest>(
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::<EditorBackgroundRemovalRequest>(
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::<EditorIconSpritesheetGenerationRequest>(
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::<EditorUiDesignAssetExtractionRequest>(
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::<serde_json::Value>(&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<String> {
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<String> {
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<Value> {
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::<String>();
if chars.next().is_some() {
bounded.push('…');
}
bounded
}
async fn complete_job(
state: &AppState,
worker_id: &str,
job: &ExternalGenerationJobRecord,
result_payload_json: Option<String>,
) -> 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<String, String> {
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<ExternalGenerationWriteLeaseGuard, String> {
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, String> {
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<u32, String> {
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,
}
}
}