Files
Genarrative/server-rs/crates/api-server/src/external_generation_worker.rs
T
k88936 e35429bb88 撤销画布 Agent 图片编辑门禁改动
撤销基于调用来源全面绕过图片编辑素材白名单的实现。

为改用普通图片空素材类型及有限历史兼容恢复基线。
2026-08-06 15:53:45 +08:00

2533 lines
92 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, Instant},
};
use axum::Json;
use serde_json::{Value, json};
use shared_kernel::offset_datetime_to_unix_micros;
use spacetime_client::{
ExternalGenerationJobClaimRecordInput, ExternalGenerationJobCompleteRecordInput,
ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord,
ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription,
};
use tokio::{
task::{JoinHandle, JoinSet},
time::sleep,
};
use tracing::{error, info, warn};
const MAX_EDITOR_GENERATION_WARNING_CHARS: usize = 2_048;
// provider 必须先结束,给失败审计、计费结算和队列终态写回保留有效 lease 内的收尾窗口。
const EXTERNAL_GENERATION_WORKER_TERMINAL_WRITE_RESERVE: Duration = Duration::from_secs(60);
const EDITOR_GENERATION_SLICE_WARNING_PREFIX: &str = "图集已生成,但自动拆分未完成:";
const EDITOR_GENERATION_WARNING_REDACTED_MESSAGE: &str =
"自动拆分未完成(告警详情含内联媒体引用,已省略)";
use crate::{
asset_billing::with_external_generation_billing_attempt_context,
character_animation_assets::{
generate_editor_character_animation_for_owner, generate_editor_video_for_owner,
},
config::AppConfig,
editor_generation_queue::{
EDITOR_BACKGROUND_MUSIC_GENERATION_JOB_KIND, EDITOR_BACKGROUND_REMOVAL_JOB_KIND,
EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND,
EDITOR_ICON_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_GENERATION_MULTIPLE_WARNINGS_CODE, EditorBackgroundRemovalRequest,
EditorGenerationCaller, EditorGenerationPhaseReporter,
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,
},
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;
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 (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,
"任务超过执行预算",
);
Err(message)
}
ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error) => {
detach_external_generation_work_until_lease_expiry(
work_handle,
&job,
lease,
"任务租约续期失败",
);
Err(error)
}
}
}
async fn join_external_generation_work(
work_handle: &mut JoinHandle<Result<(), String>>,
) -> 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<Result<(), String>>,
job: &ExternalGenerationJobRecord,
lease: Duration,
reason: &'static str,
) {
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();
warn!(
job_id = %job_id,
job_kind = %job_kind,
reason,
grace_ms = grace.as_millis() as u64,
"external generation worker 脱管任务超过租约仲裁窗口仍未结束,已取消;其后续写回将被 lease fencing 拒绝"
);
}
}
});
}
#[derive(Debug, PartialEq, Eq)]
enum ExternalGenerationJobExecutionOutcome {
Finished(Result<(), String>),
TimedOut,
LeaseRenewalFailed(String),
}
async fn await_external_generation_job_execution<W, H>(
work: W,
heartbeat: H,
job_deadline: Instant,
) -> ExternalGenerationJobExecutionOutcome
where
W: Future<Output = Result<(), String>>,
H: Future<Output = Result<(), String>>,
{
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::<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)
}
}
}
#[cfg(any())]
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)
}
}
}
#[cfg(any())]
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)
}
}
}
#[cfg(any())]
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)
}
}
}
#[cfg(any())]
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, provider_deadline);
match generate_editor_image_for_owner(
&state,
&request_context,
editor_generation_worker_caller(&worker_id, &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_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, provider_deadline);
match edit_editor_image_for_owner(
&state,
&request_context,
editor_generation_worker_caller(&worker_id, &job)?,
payload,
)
.await
{
Ok(result) => {
complete_editor_generation_job(&state, &worker_id, &job, result.0).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, provider_deadline);
match remove_editor_image_background_for_owner(
&state,
&request_context,
editor_generation_worker_caller(&worker_id, &job)?,
payload,
)
.await
{
Ok(result) => {
complete_editor_generation_job(&state, &worker_id, &job, result.0).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, provider_deadline);
match generate_editor_icon_spritesheet_for_owner(
&state,
&request_context,
editor_generation_worker_caller(&worker_id, &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, provider_deadline);
match extract_editor_ui_design_assets_for_owner(
&state,
&request_context,
editor_generation_worker_caller(&worker_id, &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, provider_deadline);
match generate_editor_character_animation_for_owner(
state.clone(),
request_context,
job.owner_user_id.clone(),
Ok(Json(payload)),
Some(editor_generation_phase_reporter(&worker_id, &job)?),
)
.await
{
Ok(result) => {
complete_editor_generation_job(&state, &worker_id, &job, result.0).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, provider_deadline);
match generate_editor_video_for_owner(
state.clone(),
request_context,
job.owner_user_id.clone(),
Ok(Json(payload)),
)
.await
{
Ok(result) => {
complete_editor_generation_job(&state, &worker_id, &job, result.0).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, provider_deadline);
match generate_editor_sound_effect_for_owner(
state.clone(),
request_context,
job.owner_user_id.clone(),
Ok(Json(payload)),
)
.await
{
Ok(result) => {
complete_editor_generation_job(&state, &worker_id, &job, result.0).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, provider_deadline);
match generate_editor_background_music_for_owner(
state.clone(),
request_context,
job.owner_user_id.clone(),
Ok(Json(payload)),
)
.await
{
Ok(result) => {
complete_editor_generation_job(&state, &worker_id, &job, result.0).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)
}
#[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)
}
fn editor_generation_worker_caller(
worker_id: &str,
job: &ExternalGenerationJobRecord,
) -> Result<EditorGenerationCaller, String> {
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)?),
})
}
fn editor_generation_phase_reporter(
worker_id: &str,
job: &ExternalGenerationJobRecord,
) -> Result<EditorGenerationPhaseReporter, String> {
Ok(EditorGenerationPhaseReporter::new(
job.job_id.clone(),
worker_id.to_string(),
require_job_lease_token(job)?,
))
}
async fn complete_editor_generation_job(
state: &AppState,
worker_id: &str,
job: &ExternalGenerationJobRecord,
response: Value,
) -> Result<(), String> {
complete_job(
state,
worker_id,
job,
Some(editor_generation_result_payload_json(job, &response)),
)
.await
}
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 is_editor_agent_generation_job(job)
&& let Some(object) = payload.as_object_mut()
{
// The Agent needs this compact result to restore its tool-call card. Other jobs keep
// master's metadata-only completion payload to avoid turning the queue into an asset API.
object.insert(
// TODO extract const
"editor-agent-tool-call-result".to_string(),
compact_editor_generation_result(response.clone()),
);
}
if is_external_api_generation_job(job)
&& let Some(object) = payload.as_object_mut()
{
object.insert(
"result".to_string(),
compact_external_api_generation_result(response.clone()),
);
}
if let Some(warning) = extract_editor_generation_warning(response)
&& let Some(object) = payload.as_object_mut()
{
object.insert("warning".to_string(), warning);
}
payload.to_string()
}
fn is_external_api_generation_job(job: &ExternalGenerationJobRecord) -> bool {
job.dedupe_key
.trim()
.starts_with("external-api-generation:")
}
fn is_editor_agent_generation_job(job: &ExternalGenerationJobRecord) -> bool {
serde_json::from_str::<Value>(job.request_payload_json.as_str())
.ok()
.is_some_and(|payload| {
payload
.pointer("/generationInputs/source")
.and_then(Value::as_str)
// TODO extract const
.is_some_and(|source| source.trim() == "editor-agent")
})
}
fn compact_editor_generation_result(mut result: Value) -> Value {
let Some(object) = result.as_object_mut() else {
return result;
};
// 紧凑结果会回传给普通用户的 Agent 工具调用卡片,不能把供应商或内部后处理实现
// 当作可见生成信息下发。正常的用户可见模型仍然保留,以便卡片恢复原有展示。
object.remove("provider");
if object
.get("model")
.and_then(Value::as_str)
.is_some_and(is_editor_internal_processing_model)
{
object.remove("model");
}
if let Some(generation_inputs) = object
.get_mut("generationInputs")
.and_then(Value::as_object_mut)
{
for field in ["screenColorHex", "mattingProvider", "mattingModel"] {
generation_inputs.remove(field);
}
}
object.remove("project");
object.remove("asset");
object.remove("spritesheetAsset");
for field in ["resource", "spritesheetResource"] {
let Some(resource) = object.get_mut(field).and_then(Value::as_object_mut) else {
continue;
};
resource.retain(|key, _| {
matches!(
key.as_str(),
"resourceId" | "objectKey" | "assetObjectId" | "sourceResourceId"
)
});
}
if let Some(icon_image_srcs) = object
.get_mut("iconImageSrcs")
.and_then(Value::as_array_mut)
{
for icon in icon_image_srcs {
let Some(icon) = icon.as_object_mut() else {
continue;
};
if let Some(resource) = icon.get_mut("resource").and_then(Value::as_object_mut) {
resource.retain(|key, _| {
matches!(
key.as_str(),
"resourceId" | "objectKey" | "assetObjectId" | "sourceResourceId"
)
});
}
icon.retain(|key, _| {
matches!(
key.as_str(),
"name" | "imageSrc" | "width" | "height" | "resource"
)
});
}
}
result
}
fn compact_external_api_generation_result(result: Value) -> Value {
let mut result = result.get("data").cloned().unwrap_or(result);
let Some(object) = result.as_object_mut() else {
return Value::Null;
};
object.retain(|key, _| {
matches!(
key.as_str(),
"ok" | "imageSrc"
| "videoSrc"
| "audioSrc"
| "previewVideoPath"
| "thumbnailSrc"
| "objectKey"
| "assetObjectId"
| "width"
| "height"
| "sourceType"
| "model"
| "taskId"
| "durationSeconds"
| "resolution"
| "priceMudPoints"
| "audioKind"
| "spritesheetImageSrc"
| "spritesheetWidth"
| "spritesheetHeight"
| "iconImageSrcs"
| "frames"
| "frameCount"
| "frameWidth"
| "frameHeight"
| "fps"
| "resource"
| "asset"
| "spritesheetResource"
| "spritesheetAsset"
| "warning"
| "sliceWarning"
)
});
for field in ["resource", "spritesheetResource"] {
if let Some(resource) = object.get_mut(field).and_then(Value::as_object_mut) {
compact_external_generation_resource(resource);
}
}
for field in ["asset", "spritesheetAsset"] {
if let Some(asset) = object.get_mut(field).and_then(Value::as_object_mut) {
compact_external_generation_asset(asset);
}
}
if let Some(icons) = object
.get_mut("iconImageSrcs")
.and_then(Value::as_array_mut)
{
for icon in icons {
let Some(icon) = icon.as_object_mut() else {
continue;
};
icon.retain(|key, _| {
matches!(
key.as_str(),
"name" | "imageSrc" | "objectKey" | "width" | "height" | "resource" | "asset"
)
});
if let Some(resource) = icon.get_mut("resource").and_then(Value::as_object_mut) {
compact_external_generation_resource(resource);
}
if let Some(asset) = icon.get_mut("asset").and_then(Value::as_object_mut) {
compact_external_generation_asset(asset);
}
remove_unstable_external_generation_media_fields(icon);
}
}
if let Some(frames) = object.get_mut("frames").and_then(Value::as_array_mut) {
for frame in frames {
let Some(frame) = frame.as_object_mut() else {
continue;
};
frame.retain(|key, _| {
matches!(
key.as_str(),
"frameIndex" | "imageSrc" | "objectKey" | "width" | "height"
)
});
remove_unstable_external_generation_media_fields(frame);
}
}
for field in ["warning", "sliceWarning"] {
if let Some(warning) = object.get_mut(field).and_then(Value::as_object_mut)
&& let Some(reason) = warning.get_mut("reason")
&& let Some(value) = reason.as_str()
{
*reason = Value::String(normalize_editor_generation_warning_reason(value));
}
}
remove_unstable_external_generation_media_fields(object);
result
}
fn compact_external_generation_resource(resource: &mut serde_json::Map<String, Value>) {
resource.retain(|key, _| {
matches!(
key.as_str(),
"resourceId"
| "projectId"
| "objectKey"
| "assetObjectId"
| "imageSrc"
| "width"
| "height"
| "sourceType"
| "assetKind"
| "taskId"
| "sourceResourceId"
)
});
remove_unstable_external_generation_media_fields(resource);
}
fn compact_external_generation_asset(asset: &mut serde_json::Map<String, Value>) {
asset.retain(|key, _| {
matches!(
key.as_str(),
"assetId"
| "folderId"
| "objectKey"
| "assetObjectId"
| "imageSrc"
| "thumbnailSrc"
| "width"
| "height"
| "sourceType"
| "assetKind"
| "taskId"
)
});
remove_unstable_external_generation_media_fields(asset);
}
fn remove_unstable_external_generation_media_fields(object: &mut serde_json::Map<String, Value>) {
object.retain(|key, value| {
if !matches!(
key.as_str(),
"imageSrc"
| "videoSrc"
| "audioSrc"
| "previewVideoPath"
| "thumbnailSrc"
| "spritesheetImageSrc"
) {
return true;
}
value
.as_str()
.is_some_and(is_stable_external_generation_media_reference)
});
}
fn is_stable_external_generation_media_reference(value: &str) -> bool {
let value = value.trim();
!value.is_empty()
&& value.starts_with('/')
&& !value.starts_with("//")
&& !value.contains('?')
&& !value.contains('#')
&& !value.to_ascii_lowercase().starts_with("data:")
&& !value.to_ascii_lowercase().starts_with("blob:")
&& !value.to_ascii_lowercase().starts_with("http://")
&& !value.to_ascii_lowercase().starts_with("https://")
}
fn is_editor_internal_processing_model(model: &str) -> bool {
matches!(
model.trim().to_ascii_lowercase().as_str(),
"anime-seg"
| "bgfilter complex"
| "birefnet"
| "connected-components"
| "screen-color-keying"
| "segment-common-image"
)
}
fn extract_editor_generation_warning_fields(
warning: Option<&Value>,
is_slice_warning: bool,
) -> Option<(String, String)> {
let warning = warning?;
let code = warning.get("code")?.as_str()?.trim();
let reason = warning.get("reason")?.as_str()?.trim();
if code.is_empty() || reason.is_empty() {
return None;
}
let reason = if is_slice_warning {
format!("{EDITOR_GENERATION_SLICE_WARNING_PREFIX}{reason}")
} else {
reason.to_string()
};
Some((code.to_string(), reason))
}
fn extract_editor_generation_warning(response: &Value) -> Option<Value> {
let data = response.get("data").unwrap_or(response);
// 中文注释:风格归一化和像素规整产生的通用 warning 可以与 sliceWarning 并存。
// 队列结果只有一个有界 warning 字段,因此按与 inline 响应相同的策略归一:
// code 不同时收敛为 multiple-generation-warningsreason 按“通用在前、拆分在后”
// 顺序拼接,再交给既有上界收敛,不允许其中任何一条被静默丢弃。
let common = extract_editor_generation_warning_fields(data.get("warning"), false);
let slice = extract_editor_generation_warning_fields(data.get("sliceWarning"), true);
let (code, reason) = match (common, slice) {
(None, None) => return None,
(Some(warning), None) | (None, Some(warning)) => warning,
(Some((common_code, common_reason)), Some((slice_code, slice_reason))) => {
let code = if common_code == slice_code {
common_code
} else {
EDITOR_GENERATION_MULTIPLE_WARNINGS_CODE.to_string()
};
(code, format!("{common_reason} {slice_reason}"))
}
};
let reason = normalize_editor_generation_warning_reason(reason.as_str());
Some(json!({
"code": code,
"reason": reason,
}))
}
fn normalize_editor_generation_warning_reason(reason: &str) -> String {
let normalized = reason.to_ascii_lowercase();
if normalized.contains("data:")
|| normalized.contains("blob:")
|| normalized.contains("http://")
|| normalized.contains("https://")
|| normalized.contains("x-amz-")
|| normalized.contains("signature=")
{
return EDITOR_GENERATION_WARNING_REDACTED_MESSAGE.to_string();
}
let mut chars = reason.chars();
let mut bounded = chars
.by_ref()
.take(MAX_EDITOR_GENERATION_WARNING_CHARS)
.collect::<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))
}
#[cfg(any())]
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_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_SPRITESHEET_GENERATION_JOB_KIND
| EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND
| EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND
| EDITOR_VIDEO_GENERATION_JOB_KIND => config.external_generation_worker_long_job_timeout,
_ => config.external_generation_worker_job_timeout,
}
}
fn external_generation_worker_timeout_message(
job: &ExternalGenerationJobRecord,
timeout: Duration,
) -> String {
format!(
"外部生成任务 {}{})超过 worker 执行预算 {} 秒,已停止当前尝试",
job.job_id,
job.job_kind,
timeout.as_secs()
)
}
fn current_utc_micros() -> i64 {
offset_datetime_to_unix_micros(time::OffsetDateTime::now_utc())
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(any())]
#[test]
fn worker_write_guard_uses_claimed_job_lease_token() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let guard = build_external_generation_write_lease_guard("worker-a", &job)
.expect("guard should build");
assert_eq!(guard.job_id.as_deref(), Some("extgen-1"));
assert_eq!(guard.worker_id.as_deref(), Some("worker-a"));
assert_eq!(guard.lease_token.as_deref(), Some("lease-1"));
}
#[test]
fn editor_worker_caller_carries_phase_reporter() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let caller = editor_generation_worker_caller("worker-a", &job)
.expect("worker caller should include claimed job phase reporter");
assert!(caller.phase_reporter.is_some());
}
#[cfg(any())]
#[test]
fn worker_write_guard_requires_claimed_job_lease_token() {
let job = external_generation_job_record_fixture(None);
let error = build_external_generation_write_lease_guard("worker-a", &job)
.expect_err("missing token should fail");
assert!(error.contains("缺少 lease token"));
}
#[test]
fn worker_billing_price_rejects_values_above_u32() {
let mut job = external_generation_job_record_fixture(Some("lease-1"));
job.price_mud_points = u64::from(u32::MAX);
assert_eq!(
external_generation_billing_price_mud_points(&job),
Ok(u32::MAX)
);
job.price_mud_points += 1;
let error = external_generation_billing_price_mud_points(&job)
.expect_err("out-of-range queued price should fail");
assert!(error.contains("extgen-1"));
assert!(error.contains("price_mud_points=4294967296"));
assert!(error.contains("4294967295"));
}
#[test]
fn worker_billing_requires_claim_attempt() {
let mut job = external_generation_job_record_fixture(Some("lease-1"));
assert_eq!(external_generation_billing_claim_attempt(&job), Ok(1));
job.attempt = 0;
let error = external_generation_billing_claim_attempt(&job)
.expect_err("unclaimed job should not enter billing context");
assert!(error.contains("extgen-1"));
assert!(error.contains("claim attempt"));
}
#[tokio::test]
async fn worker_response_error_message_prefers_api_error_details() {
let response = axum::http::Response::builder()
.status(axum::http::StatusCode::BAD_GATEWAY)
.body(axum::body::Body::from(
json!({
"ok": false,
"data": null,
"error": {
"code": "UPSTREAM_ERROR",
"message": "上游服务请求失败",
"details": {
"provider": "vector-engine",
"reason": "VECTOR_ENGINE_API_KEY 未配置",
"message": "提交编辑器音效任务失败:missing field sound"
}
},
"meta": {}
})
.to_string(),
))
.expect("response should build");
let message = response_error_message(response).await;
assert_eq!(message, "VECTOR_ENGINE_API_KEY 未配置");
}
#[tokio::test]
async fn worker_heartbeat_wait_does_not_stop_work_progress() {
let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
let work_started = std::sync::Arc::new(tokio::sync::Notify::new());
let work_connection = connection.clone();
let work_started_signal = work_started.clone();
let work = async move {
let permit = work_connection
.acquire_owned()
.await
.expect("work should acquire the only connection");
work_started_signal.notify_one();
tokio::time::sleep(Duration::from_millis(25)).await;
drop(permit);
Ok(())
};
let heartbeat_connection = connection.clone();
let heartbeat = async move {
work_started.notified().await;
let _permit = heartbeat_connection
.acquire_owned()
.await
.expect("heartbeat should acquire after work releases the connection");
std::future::pending::<Result<(), String>>().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::<Result<(), String>>().await
};
let heartbeat = std::future::pending::<Result<(), String>>();
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),
"任务超过执行预算",
);
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
finished.load(std::sync::atomic::Ordering::SeqCst),
"超时后在途 work 应继续执行完成,而不是被取消"
);
}
#[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::<Result<(), String>>().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),
"任务超过执行预算",
);
let reacquired =
tokio::time::timeout(Duration::from_secs(2), connection.acquire_owned()).await;
assert!(
reacquired.is_ok(),
"超过租约仲裁窗口后应取消仍未结束的 work 并释放其持有的资源"
);
}
#[test]
fn editor_generation_result_payload_keeps_only_lightweight_slice_warning() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let response = json!({
"spritesheetImageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST",
"iconImageSrcs": [{"imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST"}],
"sliceWarning": {
"code": "insufficient-connected-components",
"reason": "连通域数量不足"
}
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("worker 结果应是合法 JSON");
assert_eq!(payload["sourceModule"], json!("editor"));
assert_eq!(payload["sourceEntityId"], json!("project-1"));
assert_eq!(
payload["warning"],
json!({
"code": "insufficient-connected-components",
"reason": "图集已生成,但自动拆分未完成:连通域数量不足"
})
);
assert!(payload.get("spritesheetImageSrc").is_none());
assert!(payload.get("iconImageSrcs").is_none());
assert!(payload.get("editor-agent-tool-call-result").is_none());
}
#[test]
fn editor_agent_result_payload_keeps_compact_response() {
let mut job = external_generation_job_record_fixture(Some("lease-1"));
job.request_payload_json = json!({
"generationInputs": { "source": "editor-agent" },
})
.to_string();
let response = json!({
"imageSrc": "/api/assets/object/generated.png",
"objectKey": "users/user-1/generated.png",
"assetObjectId": "asset-object-1",
"width": 1024,
"height": 1024,
"sourceType": "generated",
"prompt": "castle",
"actualPrompt": null,
"model": "gpt-image-2",
"provider": "VectorEngine",
"generationInputs": {
"fields": { "prompt": "castle" },
"characterAnimation": { "frameCount": 8 },
"screenColorHex": "#CFEFFF",
"mattingProvider": "BgFilter",
"mattingModel": "birefnet",
},
"taskId": "provider-task-1",
"resource": {
"resourceId": "resource-1",
"objectKey": "users/user-1/generated.png",
"assetObjectId": "asset-object-1",
"imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST",
},
"asset": { "assetId": "asset-1" },
"project": { "projectId": "project-1" },
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("worker result should be valid JSON");
assert_eq!(
payload["editor-agent-tool-call-result"]["imageSrc"],
json!("/api/assets/object/generated.png")
);
assert_eq!(
payload["editor-agent-tool-call-result"]["resource"],
json!({
"resourceId": "resource-1",
"objectKey": "users/user-1/generated.png",
"assetObjectId": "asset-object-1",
})
);
assert!(
payload["editor-agent-tool-call-result"]
.get("asset")
.is_none()
);
assert!(
payload["editor-agent-tool-call-result"]
.get("project")
.is_none()
);
assert_eq!(
payload["editor-agent-tool-call-result"]["model"],
json!("gpt-image-2")
);
assert!(
payload["editor-agent-tool-call-result"]
.get("provider")
.is_none()
);
assert_eq!(
payload["editor-agent-tool-call-result"]["generationInputs"],
json!({
"fields": { "prompt": "castle" },
"characterAnimation": { "frameCount": 8 },
})
);
assert!(!payload.to_string().contains("data:image"));
}
#[test]
fn compact_editor_generation_result_removes_internal_models_case_insensitively() {
for model in [
"anime-seg",
" BgFilter Complex ",
"BIREFNET",
"connected-components",
"screen-color-keying",
"segment-common-image",
] {
let compact = compact_editor_generation_result(json!({
"model": model,
"provider": "internal-provider",
}));
assert!(compact.get("model").is_none(), "model={model}");
assert!(compact.get("provider").is_none(), "model={model}");
}
}
#[test]
fn editor_agent_spritesheet_result_keeps_all_persisted_slices() {
let mut job = external_generation_job_record_fixture(Some("lease-1"));
job.request_payload_json = json!({
"generationInputs": { "source": "editor-agent" },
})
.to_string();
let response = json!({
"spritesheetImageSrc": "/api/assets/object/sheet.png",
"spritesheetWidth": 512,
"spritesheetHeight": 512,
"taskId": "provider-task-1",
"spritesheetResource": {
"resourceId": "sheet-resource",
"objectKey": "users/user-1/sheet.png",
"assetObjectId": "sheet-object",
},
"iconImageSrcs": [
{
"name": "backpack",
"imageSrc": "/api/assets/object/backpack.png",
"width": 64,
"height": 64,
"resource": {
"resourceId": "icon-resource-1",
"objectKey": "users/user-1/backpack.png",
"assetObjectId": "icon-object-1",
"sourceResourceId": "sheet-resource",
"imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST",
},
"asset": { "assetId": "icon-asset-1" },
},
{
"name": "map",
"imageSrc": "/api/assets/object/map.png",
"width": 64,
"height": 64,
"resource": {
"resourceId": "icon-resource-2",
"objectKey": "users/user-1/map.png",
"assetObjectId": "icon-object-2",
"sourceResourceId": "sheet-resource",
},
},
],
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("worker result should be valid JSON");
assert_eq!(
payload["editor-agent-tool-call-result"]["iconImageSrcs"]
.as_array()
.map(Vec::len),
Some(2)
);
assert_eq!(
payload["editor-agent-tool-call-result"]["iconImageSrcs"][0]["resource"],
json!({
"resourceId": "icon-resource-1",
"objectKey": "users/user-1/backpack.png",
"assetObjectId": "icon-object-1",
"sourceResourceId": "sheet-resource",
})
);
assert!(
payload["editor-agent-tool-call-result"]["iconImageSrcs"][0]
.get("asset")
.is_none()
);
assert!(!payload.to_string().contains("data:image"));
}
#[test]
fn editor_generation_result_payload_merges_common_and_slice_warnings() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let response = json!({
"data": {
"imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST",
"warning": {
"code": "unsupported-image-style",
"reason": "不支持的图片风格,已按无风格继续生成。"
},
"sliceWarning": {
"code": "insufficient-connected-components",
"reason": "有效连通域不足"
}
}
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("worker 结果应是合法 JSON");
// 中文注释:风格归一化告警与拆分告警可以并存,队列只有一个 warning 字段,
// 必须拼接后收敛 code,不能让其中任何一条消失。
assert_eq!(
payload["warning"],
json!({
"code": "multiple-generation-warnings",
"reason": "不支持的图片风格,已按无风格继续生成。 图集已生成,但自动拆分未完成:有效连通域不足"
})
);
assert!(payload.get("imageSrc").is_none());
}
#[test]
fn editor_generation_result_payload_keeps_single_warning_untouched() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let response = json!({
"data": {
"warning": {
"code": "postprocess-failed-source-preserved",
"reason": "生成任务成功,后处理失败。"
}
}
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("worker 结果应是合法 JSON");
// 中文注释:透明背景最终失败不会进入拆分,此时仍是单条告警,原样保留。
assert_eq!(
payload["warning"],
json!({
"code": "postprocess-failed-source-preserved",
"reason": "生成任务成功,后处理失败。"
})
);
}
#[test]
fn editor_image_job_completion_keeps_inline_response_warning() {
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];
for snippet in [
"Ok(response)",
"complete_editor_generation_job_with_response",
"&response.0",
] {
assert!(
branch.contains(snippet),
"editor image completion should preserve response warning via {snippet}"
);
}
assert!(
!branch.contains("Ok(_) => complete_editor_generation_job"),
"editor image completion must not discard the inline fallback warning"
);
}
#[test]
fn editor_generation_result_payload_accepts_envelope_and_redacts_inline_media() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let response = json!({
"data": {
"sliceWarning": {
"code": "slice-persistence-failed",
"reason": "provider returned data:image/png;base64,AAAA"
}
}
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("worker 结果应是合法 JSON");
assert_eq!(
payload["warning"]["reason"],
json!(EDITOR_GENERATION_WARNING_REDACTED_MESSAGE)
);
assert!(!payload.to_string().contains("data:image"));
}
#[test]
fn external_api_result_keeps_stable_artifacts_and_removes_unstable_media() {
let mut job = external_generation_job_record_fixture(Some("lease-1"));
job.dedupe_key = "external-api-generation:editor_image_generation:fingerprint".to_string();
let response = json!({
"data": {
"imageSrc": "data:image/png;base64,SHOULD_NOT_PERSIST",
"videoSrc": "blob:https://example.test/video",
"audioSrc": "https://cdn.example.test/audio.mp3?X-Amz-Signature=secret",
"previewVideoPath": "https://cdn.example.test/stable-looking-but-external.mp4",
"thumbnailSrc": "/api/assets/object/thumbnail.png?expires=1&signature=secret",
"objectKey": "users/user-1/generated/main.png",
"assetObjectId": "asset-object-main",
"width": 1024,
"height": 1024,
"provider": "internal-provider-must-not-persist",
"resource": {
"resourceId": "resource-main",
"projectId": "project-1",
"objectKey": "users/user-1/generated/main.png",
"assetObjectId": "asset-object-main",
"sourceResourceId": "source-resource-main",
"imageSrc": "https://cdn.example.test/main.png?signature=secret",
"width": 1024,
"height": 1024,
"prompt": "不应复制完整资源元数据"
},
"asset": {
"assetId": "asset-main",
"folderId": "folder-1",
"objectKey": "users/user-1/generated/main.png",
"assetObjectId": "asset-object-main",
"imageSrc": "/api/assets/object/main.png",
"thumbnailSrc": "https://cdn.example.test/thumb.png?signature=secret",
"width": 1024,
"height": 1024,
"generationInputs": {"private": true}
},
"project": {
"projectId": "project-1",
"canvas": {"layers": ["large-layout-must-not-persist"]}
},
"warning": {
"code": "dimension-restore-fallback",
"reason": "已保留 provider 实际输出尺寸。"
}
},
"meta": {
"requestId": "worker-envelope-must-not-persist"
}
});
let payload: Value =
serde_json::from_str(&editor_generation_result_payload_json(&job, &response))
.expect("外部生成结果应是合法 JSON");
let result = &payload["result"];
assert!(result.get("project").is_none());
assert!(result.get("provider").is_none());
for unstable_field in [
"imageSrc",
"videoSrc",
"audioSrc",
"previewVideoPath",
"thumbnailSrc",
] {
assert!(
result.get(unstable_field).is_none(),
"不稳定媒体字段 {unstable_field} 不得持久化"
);
}
assert_eq!(
result["objectKey"],
json!("users/user-1/generated/main.png")
);
assert_eq!(result["assetObjectId"], json!("asset-object-main"));
assert_eq!(result["resource"]["resourceId"], json!("resource-main"));
assert_eq!(
result["resource"]["sourceResourceId"],
json!("source-resource-main")
);
assert_eq!(
result["resource"]["objectKey"],
json!("users/user-1/generated/main.png")
);
assert!(result["resource"].get("imageSrc").is_none());
assert!(result["resource"].get("prompt").is_none());
assert_eq!(result["asset"]["assetId"], json!("asset-main"));
assert_eq!(
result["asset"]["imageSrc"],
json!("/api/assets/object/main.png")
);
assert!(result["asset"].get("thumbnailSrc").is_none());
assert!(result["asset"].get("generationInputs").is_none());
assert_eq!(
result["warning"],
json!({
"code": "dimension-restore-fallback",
"reason": "已保留 provider 实际输出尺寸。"
})
);
assert_eq!(payload["warning"], result["warning"]);
assert!(result.get("prompt").is_none());
assert!(result.get("actualPrompt").is_none());
let serialized = payload.to_string().to_ascii_lowercase();
for forbidden in [
"data:",
"blob:",
"x-amz-signature",
"?signature=",
"large-layout",
] {
assert!(
!serialized.contains(forbidden),
"compact result 不应包含 {forbidden}"
);
}
}
#[test]
fn non_external_job_does_not_publish_query_result() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let payload: Value = serde_json::from_str(&editor_generation_result_payload_json(
&job,
&json!({
"objectKey": "users/user-1/generated/main.png",
"resource": {"resourceId": "resource-main"}
}),
))
.expect("普通编辑器任务结果应为合法 JSON");
assert!(payload.get("result").is_none());
}
#[test]
fn worker_job_timeout_uses_long_budget_for_image_and_video_jobs() {
let config = AppConfig {
external_generation_worker_job_timeout: Duration::from_secs(60),
external_generation_worker_long_job_timeout: Duration::from_secs(600),
..AppConfig::default()
};
assert_eq!(
external_generation_worker_job_timeout(&config, EDITOR_IMAGE_GENERATION_JOB_KIND),
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(&config, EDITOR_IMAGE_EDIT_JOB_KIND),
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(
&config,
EDITOR_ICON_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 秒"));
}
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()),
}
}
}