Files
Genarrative/server-rs/crates/spacetime-module/src/external_generation.rs
T
k88936 ce464168ce docs(会员): 标注外部生成认领的跨账号队头阻塞 TODO
- claim_external_generation_jobs_tx 的 take(limit) 只检查前 limit 行,被上限拦下/lease 终结的行仍占检查名额
- 记录触发条件(负载越高 limit 越小)与未决修法(改为成功认领才计数,仍是全局 FIFO),本次不修
2026-10-04 13:29:32 +08:00

4016 lines
142 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 crate::*;
use std::cmp::Ordering;
use std::collections::HashMap;
use std::ops::RangeFrom;
const EXTERNAL_GENERATION_STATUS_PENDING: &str = "pending";
const EXTERNAL_GENERATION_STATUS_RUNNING: &str = "running";
const EXTERNAL_GENERATION_STATUS_COMPLETED: &str = "completed";
const EXTERNAL_GENERATION_STATUS_FAILED: &str = "failed";
const EXTERNAL_GENERATION_STATUS_CANCELLED: &str = "cancelled";
const EXTERNAL_GENERATION_PHASE_GENERATING: &str = "generating";
const EXTERNAL_GENERATION_PHASE_PROCESSING: &str = "processing";
const EXTERNAL_GENERATION_EVENT_ENQUEUED: &str = "enqueued";
const EXTERNAL_GENERATION_EVENT_CLAIMED: &str = "claimed";
const EXTERNAL_GENERATION_EVENT_LEASE_RENEWED: &str = "lease_renewed";
const EXTERNAL_GENERATION_EVENT_COMPLETED: &str = "completed";
const EXTERNAL_GENERATION_EVENT_FAILED: &str = "failed";
const EXTERNAL_GENERATION_EVENT_ACKNOWLEDGED: &str = "acknowledged";
const EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE: &str = "editor-canvas";
const EXTERNAL_GENERATION_FINAL_ATTEMPT_LEASE_EXPIRED_MESSAGE: &str =
"worker 最终执行次数的 lease 已过期,任务已终止";
const MAX_EXTERNAL_GENERATION_PAYLOAD_BYTES: usize = 512 * 1024;
const MAX_EXTERNAL_GENERATION_REQUEST_PROMPT_CHARS: usize = 2_048;
const MAX_EXTERNAL_GENERATION_ERROR_MESSAGE_CHARS: usize = 2_048;
const MAX_EXTERNAL_GENERATION_WARNING_MESSAGE_CHARS: usize = 2_048;
const MAX_EXTERNAL_GENERATION_MAINTENANCE_BATCH_SIZE: u32 = 25;
const EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE: usize = 256;
const INLINE_MEDIA_REMOVED_PLACEHOLDER: &str = "[inline-media-removed]";
const INLINE_MEDIA_ERROR_REDACTED_MESSAGE: &str = "外部生成失败(错误详情含内联媒体引用,已省略)";
const INLINE_MEDIA_WARNING_REDACTED_MESSAGE: &str =
"外部生成已完成(告警详情含内联媒体引用,已省略)";
#[spacetimedb::table(
accessor = external_generation_job,
index(
accessor = by_external_generation_job_status_available,
btree(columns = [status, available_at])
),
index(
accessor = by_external_generation_job_worker_id,
btree(columns = [worker_id])
),
index(
accessor = by_external_generation_job_source,
btree(columns = [source_module, source_entity_id])
),
index(
accessor = by_external_generation_job_owner_user_id,
btree(columns = [owner_user_id])
),
// 会员并发上限在认领事务里按账号现算 `status = running` 行数;不建计数表,`running` 行即真相。
index(
accessor = by_external_generation_job_owner_status,
btree(columns = [owner_user_id, status])
),
index(
accessor = by_external_generation_job_cursor,
btree(columns = [job_id, source_module])
),
index(
accessor = by_external_generation_job_source_cursor,
btree(columns = [source_module, job_id])
)
)]
#[derive(Clone)]
pub struct ExternalGenerationJob {
#[primary_key]
pub(crate) job_id: String,
#[unique]
pub(crate) dedupe_key: String,
pub(crate) job_kind: String,
pub(crate) owner_user_id: String,
pub(crate) source_module: String,
pub(crate) source_entity_id: String,
pub(crate) request_label: String,
pub(crate) request_payload_json: String,
pub(crate) status: String,
pub(crate) attempt: u32,
pub(crate) max_attempts: u32,
pub(crate) last_error_message: Option<String>,
pub(crate) worker_id: Option<String>,
pub(crate) lease_expires_at: Option<Timestamp>,
pub(crate) available_at: Timestamp,
pub(crate) result_payload_json: Option<String>,
pub(crate) created_at: Timestamp,
pub(crate) started_at: Option<Timestamp>,
pub(crate) completed_at: Option<Timestamp>,
pub(crate) updated_at: Timestamp,
#[default(None::<String>)]
pub(crate) lease_token: Option<String>,
#[default(0u64)]
pub(crate) price_mud_points: u64,
#[default(None::<String>)]
pub(crate) refund_ledger_id: Option<String>,
#[default(None::<Timestamp>)]
pub(crate) notification_acknowledged_at: Option<Timestamp>,
#[default(None::<String>)]
pub(crate) phase: Option<String>,
}
#[spacetimedb::table(
accessor = external_generation_job_event,
index(
accessor = by_external_generation_job_event_job_id,
btree(columns = [job_id, created_at])
),
index(
accessor = by_external_generation_job_event_job_id_only,
btree(columns = [job_id])
),
index(
accessor = by_external_generation_job_event_owner,
btree(columns = [owner_user_id, created_at])
)
)]
#[derive(Clone)]
pub struct ExternalGenerationJobEvent {
#[primary_key]
pub(crate) event_id: String,
pub(crate) job_id: String,
pub(crate) owner_user_id: String,
pub(crate) event_kind: String,
pub(crate) status: String,
pub(crate) message: Option<String>,
pub(crate) worker_id: Option<String>,
pub(crate) created_at: Timestamp,
}
#[spacetimedb::table(
accessor = external_generation_job_summary,
index(
accessor = by_external_generation_job_summary_owner_user_id,
btree(columns = [owner_user_id])
)
)]
#[derive(Clone)]
pub struct ExternalGenerationJobSummary {
#[primary_key]
pub(crate) job_id: String,
pub(crate) job_kind: String,
pub(crate) owner_user_id: String,
pub(crate) source_module: String,
pub(crate) source_entity_id: String,
pub(crate) request_label: String,
pub(crate) request_prompt: Option<String>,
pub(crate) status: String,
pub(crate) last_error_message: Option<String>,
pub(crate) created_at: Timestamp,
pub(crate) started_at: Option<Timestamp>,
pub(crate) completed_at: Option<Timestamp>,
pub(crate) updated_at: Timestamp,
pub(crate) price_mud_points: u64,
pub(crate) refund_ledger_id: Option<String>,
pub(crate) notification_acknowledged_at: Option<Timestamp>,
#[default(None::<String>)]
pub(crate) warning_message: Option<String>,
#[default(None::<String>)]
pub(crate) phase: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobEnqueueInput {
pub job_id: String,
pub dedupe_key: String,
pub job_kind: String,
pub owner_user_id: String,
pub source_module: String,
pub source_entity_id: String,
pub request_label: String,
pub request_payload_json: String,
pub max_attempts: u32,
pub available_at_micros: i64,
pub created_at_micros: i64,
pub price_mud_points: u64,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobClaimInput {
pub worker_id: String,
pub limit: u32,
pub lease_expires_at_micros: i64,
pub claimed_at_micros: i64,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobRenewLeaseInput {
pub job_id: String,
pub worker_id: String,
pub lease_token: String,
pub lease_expires_at_micros: i64,
pub renewed_at_micros: i64,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobPhaseUpdateInput {
pub job_id: String,
pub worker_id: String,
pub lease_token: String,
pub phase: String,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, SpacetimeType)]
pub enum ExternalGenerationJobPhaseUpdateFailureKind {
LeaseFencingRejected,
OtherRejected,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobCompleteInput {
pub job_id: String,
pub worker_id: String,
pub lease_token: String,
pub result_payload_json: Option<String>,
pub completed_at_micros: i64,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobFailInput {
pub job_id: String,
pub worker_id: String,
pub lease_token: String,
pub error_message: String,
pub retry_after_micros: i64,
pub failed_at_micros: i64,
pub refund_ledger_id: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobGetInput {
pub job_id: String,
pub owner_user_id: String,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobListInput {
pub owner_user_id: String,
pub limit: u32,
pub include_acknowledged_terminal: bool,
pub statuses: Vec<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobAcknowledgeInput {
pub owner_user_id: String,
pub job_ids: Vec<String>,
pub acknowledged_at_micros: i64,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobSummaryBackfillInput {
pub owner_user_id: Option<String>,
pub limit: u32,
pub cursor_job_id: Option<String>,
pub dry_run: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobPayloadCompactionInput {
pub dry_run: bool,
pub limit: u32,
pub cursor_job_id: Option<String>,
pub completed_before_micros: Option<i64>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobSnapshot {
pub job_id: String,
pub dedupe_key: String,
pub job_kind: String,
pub owner_user_id: String,
pub source_module: String,
pub source_entity_id: String,
pub request_label: String,
pub request_payload_json: String,
pub status: String,
pub attempt: u32,
pub max_attempts: u32,
pub last_error_message: Option<String>,
pub worker_id: Option<String>,
pub lease_expires_at_micros: Option<i64>,
pub available_at_micros: i64,
pub result_payload_json: Option<String>,
pub created_at_micros: i64,
pub started_at_micros: Option<i64>,
pub completed_at_micros: Option<i64>,
pub updated_at_micros: i64,
pub lease_token: Option<String>,
pub price_mud_points: u64,
pub refund_ledger_id: Option<String>,
pub notification_acknowledged_at_micros: Option<i64>,
pub phase: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobProcedureResult {
pub ok: bool,
pub job: Option<ExternalGenerationJobSnapshot>,
pub jobs: Vec<ExternalGenerationJobSnapshot>,
pub pending_count: u32,
pub running_count: u32,
pub unacknowledged_terminal_count: u32,
pub now_micros: i64,
pub error_message: Option<String>,
}
// Private backend read for reconciliation; it intentionally excludes job request and lease data.
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobResultSnapshot {
pub job_id: String,
pub status: String,
pub last_error_message: Option<String>,
pub result_payload_json: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobResultProcedureResult {
pub ok: bool,
pub result: Option<ExternalGenerationJobResultSnapshot>,
pub error_message: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobPhaseUpdateProcedureResult {
pub ok: bool,
pub job: Option<ExternalGenerationJobSnapshot>,
pub failure_kind: Option<ExternalGenerationJobPhaseUpdateFailureKind>,
pub error_message: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobSummarySnapshot {
pub job_id: String,
pub job_kind: String,
pub owner_user_id: String,
pub source_module: String,
pub source_entity_id: String,
pub request_label: String,
pub request_prompt: Option<String>,
pub status: String,
pub last_error_message: Option<String>,
pub created_at_micros: i64,
pub started_at_micros: Option<i64>,
pub completed_at_micros: Option<i64>,
pub updated_at_micros: i64,
pub price_mud_points: u64,
pub refund_ledger_id: Option<String>,
pub notification_acknowledged_at_micros: Option<i64>,
pub warning_message: Option<String>,
pub phase: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobSummaryProcedureResult {
pub ok: bool,
pub job: Option<ExternalGenerationJobSummarySnapshot>,
pub jobs: Vec<ExternalGenerationJobSummarySnapshot>,
pub pending_count: u32,
pub running_count: u32,
pub unacknowledged_terminal_count: u32,
pub now_micros: i64,
pub error_message: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobSummaryBackfillProcedureResult {
pub ok: bool,
pub dry_run: bool,
pub scanned_count: u64,
pub selected_count: u32,
pub upserted_count: u32,
pub next_cursor_job_id: Option<String>,
pub has_more: bool,
pub error_message: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobPayloadCompactionProcedureResult {
pub ok: bool,
pub dry_run: bool,
pub scanned_count: u64,
pub matched_count: u32,
pub updated_count: u32,
pub before_bytes: u64,
pub after_bytes: u64,
pub inline_media_count: u64,
pub invalid_json_count: u32,
pub next_cursor_job_id: Option<String>,
pub has_more: bool,
pub error_message: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobRetentionInput {
pub source_module: String,
pub limit: u32,
pub cursor_job_id: Option<String>,
pub completed_before_micros: i64,
pub dry_run: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationJobRetentionProcedureResult {
pub ok: bool,
pub dry_run: bool,
pub scanned_count: u64,
pub selected_count: u32,
pub deleted_job_count: u32,
pub deleted_summary_count: u32,
pub deleted_event_count: u32,
pub next_cursor_job_id: Option<String>,
pub has_more: bool,
pub error_message: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationQueueStatsSnapshot {
pub pending_count: u32,
pub delayed_pending_count: u32,
pub claimable_pending_count: u32,
pub running_active_count: u32,
pub expired_running_count: u32,
// 中文注释:保留字段兼容已生成 bindings;controller 只按非终态队列压力扩缩容,不每轮扫描历史终态任务。
pub terminal_count: u32,
pub claimable_count: u32,
pub oldest_claimable_age_micros: Option<i64>,
pub now_micros: i64,
}
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
pub struct ExternalGenerationQueueStatsProcedureResult {
pub ok: bool,
pub stats: Option<ExternalGenerationQueueStatsSnapshot>,
pub error_message: Option<String>,
}
#[spacetimedb::procedure]
pub fn enqueue_external_generation_job_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobEnqueueInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
enqueue_external_generation_job_tx(tx, input.clone())
}) {
Ok(job) => single_external_generation_job_result(job),
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn claim_external_generation_jobs_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobClaimInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
claim_external_generation_jobs_tx(tx, input.clone())
}) {
Ok(jobs) => ExternalGenerationJobProcedureResult {
ok: true,
job: None,
jobs,
pending_count: 0,
running_count: 0,
unacknowledged_terminal_count: 0,
now_micros: ctx.timestamp.to_micros_since_unix_epoch(),
error_message: None,
},
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn complete_external_generation_job_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobCompleteInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
complete_external_generation_job_tx(tx, input.clone())
}) {
Ok(job) => single_external_generation_job_result(job),
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn renew_external_generation_job_lease_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobRenewLeaseInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
renew_external_generation_job_lease_tx(tx, input.clone())
}) {
Ok(job) => single_external_generation_job_result(job),
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn update_external_generation_job_phase_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobPhaseUpdateInput,
) -> ExternalGenerationJobPhaseUpdateProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)
.map_err(ExternalGenerationJobPhaseUpdateError::other)?;
update_external_generation_job_phase_tx(tx, input.clone())
}) {
Ok(job) => ExternalGenerationJobPhaseUpdateProcedureResult {
ok: true,
job: Some(job),
failure_kind: None,
error_message: None,
},
Err(error) => ExternalGenerationJobPhaseUpdateProcedureResult {
ok: false,
job: None,
failure_kind: Some(error.kind),
error_message: Some(error.message),
},
}
}
#[spacetimedb::procedure]
pub fn fail_external_generation_job_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobFailInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
fail_external_generation_job_tx(tx, input.clone())
}) {
Ok(job) => single_external_generation_job_result(job),
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn get_external_generation_job_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobGetInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
get_external_generation_job_tx(tx, input.clone())
}) {
Ok(job) => single_external_generation_job_result(job),
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn get_external_generation_job_result_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobGetInput,
) -> ExternalGenerationJobResultProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
get_external_generation_job_result_tx(tx, input.clone())
}) {
Ok(result) => single_external_generation_job_result_read_result(result),
Err(message) => failed_external_generation_job_result_read_result(message),
}
}
#[spacetimedb::procedure]
pub fn list_external_generation_jobs_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobListInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
list_external_generation_jobs_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => failed_external_generation_job_result(message),
}
}
#[spacetimedb::procedure]
pub fn acknowledge_external_generation_jobs_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobAcknowledgeInput,
) -> ExternalGenerationJobProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
acknowledge_external_generation_jobs_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => failed_external_generation_job_result(message),
}
}
// 正式任务列表、详情与通知确认只返回轻量投影,禁止把持久任务 payload 带入 UI 读取链路。
#[spacetimedb::procedure]
pub fn get_external_generation_job_summary_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobGetInput,
) -> ExternalGenerationJobSummaryProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
get_external_generation_job_summary_tx(tx, input.clone())
}) {
Ok(job) => single_external_generation_job_summary_result(job),
Err(message) => failed_external_generation_job_summary_result(message),
}
}
#[spacetimedb::procedure]
pub fn list_external_generation_job_summaries_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobListInput,
) -> ExternalGenerationJobSummaryProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
list_external_generation_job_summaries_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => failed_external_generation_job_summary_result(message),
}
}
#[spacetimedb::procedure]
pub fn acknowledge_external_generation_job_summaries_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobAcknowledgeInput,
) -> ExternalGenerationJobSummaryProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
acknowledge_external_generation_job_summaries_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => failed_external_generation_job_summary_result(message),
}
}
// 历史投影回填是显式维护动作;正式 list 不回扫大 payload 表。
#[spacetimedb::procedure]
pub fn backfill_external_generation_job_summaries_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobSummaryBackfillInput,
) -> ExternalGenerationJobSummaryBackfillProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::migration::require_migration_operator(tx, caller)?;
backfill_external_generation_job_summaries_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => {
failed_external_generation_job_summary_backfill_result(input.dry_run, message)
}
}
}
#[spacetimedb::procedure]
pub fn compact_external_generation_job_payloads_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobPayloadCompactionInput,
) -> ExternalGenerationJobPayloadCompactionProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::migration::require_migration_operator(tx, caller)?;
compact_external_generation_job_payloads_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => {
failed_external_generation_job_payload_compaction_result(input.dry_run, message)
}
}
}
#[spacetimedb::procedure]
pub fn prune_external_generation_job_history_and_return(
ctx: &mut ProcedureContext,
input: ExternalGenerationJobRetentionInput,
) -> ExternalGenerationJobRetentionProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::migration::require_migration_operator(tx, caller)?;
prune_external_generation_job_history_tx(tx, input.clone())
}) {
Ok(result) => result,
Err(message) => failed_external_generation_job_retention_result(input.dry_run, message),
}
}
#[spacetimedb::procedure]
pub fn get_external_generation_queue_stats_and_return(
ctx: &mut ProcedureContext,
) -> ExternalGenerationQueueStatsProcedureResult {
let caller = ctx.sender();
match ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
get_external_generation_queue_stats_tx(tx)
}) {
Ok(stats) => ExternalGenerationQueueStatsProcedureResult {
ok: true,
stats: Some(stats),
error_message: None,
},
Err(message) => ExternalGenerationQueueStatsProcedureResult {
ok: false,
stats: None,
error_message: Some(message),
},
}
}
fn enqueue_external_generation_job_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobEnqueueInput,
) -> Result<ExternalGenerationJobSnapshot, String> {
validate_required("external_generation_job.job_id", &input.job_id)?;
validate_required("external_generation_job.dedupe_key", &input.dedupe_key)?;
validate_required("external_generation_job.job_kind", &input.job_kind)?;
validate_required(
"external_generation_job.owner_user_id",
&input.owner_user_id,
)?;
validate_required(
"external_generation_job.source_module",
&input.source_module,
)?;
validate_required(
"external_generation_job.source_entity_id",
&input.source_entity_id,
)?;
validate_required(
"external_generation_job.request_label",
&input.request_label,
)?;
let request_payload_json = validate_external_generation_persisted_payload_for_source(
&input.source_module,
"external_generation_job.request_payload_json",
&input.request_payload_json,
)?;
if let Some(row) = ctx
.db
.external_generation_job()
.dedupe_key()
.find(&input.dedupe_key)
{
persist_external_generation_job_summary(ctx, &row);
return Ok(map_external_generation_job_row(row));
}
if ctx
.db
.external_generation_job()
.job_id()
.find(&input.job_id)
.is_some()
{
return Err("external_generation_job.job_id 已存在".to_string());
}
let now = Timestamp::from_micros_since_unix_epoch(input.created_at_micros);
let available_at = Timestamp::from_micros_since_unix_epoch(input.available_at_micros);
let row = ExternalGenerationJob {
job_id: input.job_id.trim().to_string(),
dedupe_key: input.dedupe_key.trim().to_string(),
job_kind: input.job_kind.trim().to_string(),
owner_user_id: input.owner_user_id.trim().to_string(),
source_module: input.source_module.trim().to_string(),
source_entity_id: input.source_entity_id.trim().to_string(),
request_label: input.request_label.trim().to_string(),
request_payload_json,
status: EXTERNAL_GENERATION_STATUS_PENDING.to_string(),
attempt: 0,
max_attempts: input.max_attempts.max(1),
last_error_message: None,
worker_id: None,
lease_expires_at: None,
available_at,
result_payload_json: None,
created_at: now,
started_at: None,
completed_at: None,
updated_at: now,
lease_token: None,
price_mud_points: input.price_mud_points,
refund_ledger_id: None,
notification_acknowledged_at: None,
phase: None,
};
persist_external_generation_job_row(ctx, row.clone());
insert_external_generation_job_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_ENQUEUED,
Some("任务已入队".to_string()),
None,
now,
);
Ok(map_external_generation_job_row(row))
}
fn claim_external_generation_jobs_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobClaimInput,
) -> Result<Vec<ExternalGenerationJobSnapshot>, String> {
validate_required("external_generation_job.worker_id", &input.worker_id)?;
if input.limit == 0 {
return Ok(Vec::new());
}
let claim_time = ctx.timestamp;
let lease_duration_micros = duration_between_micros(
input.lease_expires_at_micros,
input.claimed_at_micros,
"external_generation_job.lease_duration",
)?;
let lease_expires_at = timestamp_after_micros(claim_time, lease_duration_micros);
let worker_id = input.worker_id.trim().to_string();
let limit = input.limit.min(64) as usize;
let mut candidates = Vec::new();
candidates.extend(
ctx.db
.external_generation_job()
.by_external_generation_job_status_available()
.filter(&EXTERNAL_GENERATION_STATUS_PENDING.to_string())
.filter(|row| is_external_generation_job_worker_claimable(row, claim_time)),
);
candidates.extend(
ctx.db
.external_generation_job()
.by_external_generation_job_status_available()
.filter(&EXTERNAL_GENERATION_STATUS_RUNNING.to_string())
.filter(|row| is_external_generation_job_worker_claimable(row, claim_time)),
);
candidates.sort_by(|left, right| {
left.available_at
.to_micros_since_unix_epoch()
.cmp(&right.available_at.to_micros_since_unix_epoch())
.then_with(|| {
left.created_at
.to_micros_since_unix_epoch()
.cmp(&right.created_at.to_micros_since_unix_epoch())
})
.then_with(|| left.job_id.cmp(&right.job_id))
});
// TODO(并发认领公平性): `take(limit)` 只检查按 `available_at / created_at` 排序的前 `limit` 行,
// 而 `limit` 就是 worker 空闲槽数(api-server 传 `available`,池子默认 2 / release 8),被并发上限
// 拦下或被 lease 耗尽终结的行只 `continue`,却已吃掉一个检查名额。若队头若干行都属已饱和账号,
// 本轮 `claimed` 可能为空,后面其它账号的可认领任务看不到 —— 跨账号队头阻塞;负载越高 `limit`
// 越小(7/8 worker 占用时 `limit = 1`)越容易触发。
// 修法(未决,改前需确认单事务写集上界):遍历全部候选,成功认领一行才计数,`claimed.len() >= limit`
// 才 break,被跳过 / 终结的行不占名额,仍是全局 FIFO;若要「不同账号平均穿插」则需另做按账号轮转。
let mut claimed = Vec::new();
let mut concurrency = ExternalGenerationOwnerConcurrencyCache::default();
for mut row in candidates.into_iter().take(limit) {
if external_generation_job_has_exhausted_attempts(&row) {
let owner_user_id = row.owner_user_id.clone();
finalize_external_generation_job_after_lease_exhaustion(ctx, row, claim_time)?;
// 终结可能把 running 行改成 failed,账号计数缓存随之失效。
concurrency.invalidate_running(&owner_user_id);
continue;
}
// 会员并发上限:只约束 `pending → running` 的认领。达上限的任务留在 `pending` 天然排队;
// 账号自己那条已过期的 `running` 回收不受限,否则超限账号的卡死任务永远无法回收。
if !external_generation_owner_has_concurrency_capacity(
ctx,
&mut concurrency,
&row,
claim_time,
) {
continue;
}
let next_attempt = row.attempt.saturating_add(1);
let lease_token = build_external_generation_lease_token(
&row.job_id,
&worker_id,
next_attempt,
claim_time,
);
row.status = EXTERNAL_GENERATION_STATUS_RUNNING.to_string();
row.phase = Some(EXTERNAL_GENERATION_PHASE_GENERATING.to_string());
row.worker_id = Some(worker_id.clone());
row.lease_expires_at = Some(lease_expires_at);
row.lease_token = Some(lease_token);
row.attempt = next_attempt;
if row.started_at.is_none() {
row.started_at = Some(claim_time);
}
row.updated_at = claim_time;
persist_external_generation_job_row(ctx, row.clone());
insert_external_generation_job_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_CLAIMED,
Some("worker 已领取任务".to_string()),
Some(worker_id.clone()),
claim_time,
);
// 本次认领把一行推进到 running,账号计数缓存失效,后续候选行重新现查。
concurrency.invalidate_running(&row.owner_user_id);
claimed.push(map_external_generation_job_row(row));
}
Ok(claimed)
}
/// 单次 claim 事务内按账号缓存有效并发上限与 running 行数。
///
/// 上限在一次事务内是常量;running 计数在账号的任一行被认领或终结后立即失效、下次重查,
/// 因此判定结果与逐候选行现查一致,只是把同一账号的重复查询折叠为一次。
#[derive(Default)]
struct ExternalGenerationOwnerConcurrencyCache {
limits: HashMap<String, u32>,
running: HashMap<String, u32>,
}
impl ExternalGenerationOwnerConcurrencyCache {
fn limit(&mut self, ctx: &ReducerContext, owner_user_id: &str) -> u32 {
if let Some(limit) = self.limits.get(owner_user_id) {
return *limit;
}
let limit = crate::effective_profile_concurrent_job_limit(ctx, owner_user_id);
self.limits.insert(owner_user_id.to_string(), limit);
limit
}
fn running_count(&mut self, ctx: &ReducerContext, owner_user_id: &str) -> u32 {
if let Some(count) = self.running.get(owner_user_id) {
return *count;
}
let count = count_running_external_generation_jobs_for_owner(ctx, owner_user_id);
self.running.insert(owner_user_id.to_string(), count);
count
}
/// 账号任一 `running` 行发生状态变化(认领 / 终结)后调用,保证下次读取重新现查。
fn invalidate_running(&mut self, owner_user_id: &str) {
self.running.remove(owner_user_id);
}
fn allows(&mut self, ctx: &ReducerContext, owner_user_id: &str, is_recycling: bool) -> bool {
let limit = self.limit(ctx, owner_user_id);
if is_unlimited_concurrency(limit) {
return true;
}
// 回收豁免不需要计数:自身过期 running 已计入在飞数,放行不会推高它。
let running = if is_recycling {
0
} else {
self.running_count(ctx, owner_user_id)
};
external_generation_claim_within_concurrency_limit(is_recycling, running, limit)
}
}
/// 认领前的并发上限判定。
///
/// 计数只算 `status = running`(含 lease 过期待回收的 `expired_running`):`pending`(含延时重试)
/// 不占名额,失败退回 `pending` 会自然释放。`128` 哨兵由 [`is_unlimited_concurrency`] 解释。
fn external_generation_owner_has_concurrency_capacity(
ctx: &ReducerContext,
cache: &mut ExternalGenerationOwnerConcurrencyCache,
row: &ExternalGenerationJob,
now: Timestamp,
) -> bool {
// 回收豁免:候选本身已是该账号过期的 `running`,它已计入在飞数;放行回收不会推高在飞任务数。
let is_recycling = row.status == EXTERNAL_GENERATION_STATUS_RUNNING
&& is_external_generation_job_claimable(row, now);
cache.allows(ctx, &row.owner_user_id, is_recycling)
}
/// 并发上限判定(纯函数):`running_count >= limit` 时只放行对账号自身过期 `running` 的回收。
fn external_generation_claim_within_concurrency_limit(
is_recycling: bool,
running_count: u32,
limit: u32,
) -> bool {
is_unlimited_concurrency(limit) || is_recycling || running_count < limit
}
/// 账号当前在飞的 `running` 行数;事务内的认领写入会立即反映到下一次计数。
fn count_running_external_generation_jobs_for_owner(
ctx: &ReducerContext,
owner_user_id: &str,
) -> u32 {
ctx.db
.external_generation_job()
.by_external_generation_job_owner_status()
.filter(&owner_user_id.to_string())
.filter(|row| row.status == EXTERNAL_GENERATION_STATUS_RUNNING)
.count()
.try_into()
.unwrap_or(u32::MAX)
}
fn finalize_external_generation_job_after_lease_exhaustion(
ctx: &ReducerContext,
row: ExternalGenerationJob,
failed_at: Timestamp,
) -> Result<(), String> {
let expired_worker_id = row.worker_id.clone();
let refund_ledger_id = crate::settle_external_generation_attempt_refund(
ctx,
&row.job_id,
row.attempt,
&row.owner_user_id,
row.price_mud_points,
failed_at,
"final_attempt_lease_expired",
&row.job_kind,
&row.source_entity_id,
)?;
let row = mark_external_generation_job_lease_exhausted(row, failed_at, refund_ledger_id);
persist_external_generation_job_row(ctx, row.clone());
insert_external_generation_job_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_FAILED,
Some(EXTERNAL_GENERATION_FINAL_ATTEMPT_LEASE_EXPIRED_MESSAGE.to_string()),
expired_worker_id,
failed_at,
);
Ok(())
}
fn mark_external_generation_job_lease_exhausted(
mut row: ExternalGenerationJob,
failed_at: Timestamp,
refund_ledger_id: Option<String>,
) -> ExternalGenerationJob {
row.refund_ledger_id = refund_ledger_id;
row.status = EXTERNAL_GENERATION_STATUS_FAILED.to_string();
row.last_error_message =
Some(EXTERNAL_GENERATION_FINAL_ATTEMPT_LEASE_EXPIRED_MESSAGE.to_string());
row.worker_id = None;
row.lease_expires_at = None;
row.lease_token = None;
row.completed_at = Some(failed_at);
row.updated_at = failed_at;
row
}
pub(crate) fn complete_external_generation_job_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobCompleteInput,
) -> Result<ExternalGenerationJobSnapshot, String> {
let mut row = get_worker_owned_external_generation_job(
ctx,
&input.job_id,
&input.worker_id,
&input.lease_token,
)?;
let result_payload_json = if is_external_generation_editor_source(&row.source_module) {
validate_optional_external_generation_payload_json(
"external_generation_job.result_payload_json",
input.result_payload_json.as_deref(),
)?
} else {
input
.result_payload_json
.as_deref()
.and_then(normalize_optional_text)
};
let completed_at = ctx.timestamp;
row.status = EXTERNAL_GENERATION_STATUS_COMPLETED.to_string();
row.last_error_message = None;
row.result_payload_json = result_payload_json;
row.lease_expires_at = None;
row.completed_at = Some(completed_at);
row.updated_at = completed_at;
persist_external_generation_job_row(ctx, row.clone());
insert_external_generation_job_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_COMPLETED,
Some("任务已完成".to_string()),
Some(input.worker_id),
completed_at,
);
Ok(map_external_generation_job_row(row))
}
fn get_external_generation_job_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobGetInput,
) -> Result<ExternalGenerationJobSnapshot, String> {
get_external_generation_job_summary_tx(ctx, input)
.map(map_external_generation_job_summary_to_compat_snapshot)
}
fn get_external_generation_job_result_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobGetInput,
) -> Result<ExternalGenerationJobResultSnapshot, String> {
validate_required("external_generation_job.job_id", &input.job_id)?;
validate_required(
"external_generation_job.owner_user_id",
&input.owner_user_id,
)?;
let job_id = input.job_id.trim().to_string();
let owner_user_id = input.owner_user_id.trim();
let row = ctx
.db
.external_generation_job()
.job_id()
.find(&job_id)
.ok_or_else(|| "external_generation_job 不存在".to_string())?;
if row.owner_user_id.trim() != owner_user_id {
return Err("external_generation_job 不存在".to_string());
}
Ok(ExternalGenerationJobResultSnapshot {
job_id: row.job_id,
status: row.status,
last_error_message: row.last_error_message,
result_payload_json: row.result_payload_json,
})
}
fn list_external_generation_jobs_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobListInput,
) -> Result<ExternalGenerationJobProcedureResult, String> {
let result = list_external_generation_job_summaries_tx(ctx, input)?;
Ok(ExternalGenerationJobProcedureResult {
ok: result.ok,
job: None,
jobs: result
.jobs
.into_iter()
.map(map_external_generation_job_summary_to_compat_snapshot)
.collect(),
pending_count: result.pending_count,
running_count: result.running_count,
unacknowledged_terminal_count: result.unacknowledged_terminal_count,
now_micros: result.now_micros,
error_message: result.error_message,
})
}
fn acknowledge_external_generation_jobs_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobAcknowledgeInput,
) -> Result<ExternalGenerationJobProcedureResult, String> {
let result = acknowledge_external_generation_job_summaries_tx(ctx, input)?;
Ok(ExternalGenerationJobProcedureResult {
ok: result.ok,
job: None,
jobs: result
.jobs
.into_iter()
.map(map_external_generation_job_summary_to_compat_snapshot)
.collect(),
pending_count: result.pending_count,
running_count: result.running_count,
unacknowledged_terminal_count: result.unacknowledged_terminal_count,
now_micros: result.now_micros,
error_message: result.error_message,
})
}
fn get_external_generation_job_summary_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobGetInput,
) -> Result<ExternalGenerationJobSummarySnapshot, String> {
validate_required("external_generation_job.job_id", &input.job_id)?;
validate_required(
"external_generation_job.owner_user_id",
&input.owner_user_id,
)?;
let job_id = input.job_id.trim().to_string();
let owner_user_id = input.owner_user_id.trim();
if let Some(summary) = ctx
.db
.external_generation_job_summary()
.job_id()
.find(&job_id)
{
if summary.owner_user_id.trim() != owner_user_id {
return Err("external_generation_job 不存在".to_string());
}
return Ok(map_external_generation_job_summary_row(summary));
}
// 详情兼容旧任务时只允许按主键回填一行,绝不按 owner 扫描完整 payload 表。
let row = ctx
.db
.external_generation_job()
.job_id()
.find(&job_id)
.ok_or_else(|| "external_generation_job 不存在".to_string())?;
if row.owner_user_id.trim() != owner_user_id {
return Err("external_generation_job 不存在".to_string());
}
let summary = persist_external_generation_job_summary(ctx, &row);
Ok(map_external_generation_job_summary_row(summary))
}
fn list_external_generation_job_summaries_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobListInput,
) -> Result<ExternalGenerationJobSummaryProcedureResult, String> {
validate_required(
"external_generation_job.owner_user_id",
&input.owner_user_id,
)?;
let owner_user_id = input.owner_user_id.trim().to_string();
let now_micros = ctx.timestamp.to_micros_since_unix_epoch();
let status_filter = normalize_external_generation_job_status_filter(&input.statuses);
let limit = input.limit.clamp(1, 100) as usize;
let mut rows = Vec::with_capacity(limit);
let mut pending_count = 0u32;
let mut running_count = 0u32;
let mut unacknowledged_terminal_count = 0u32;
// 这里故意只读轻量投影,并在单次 owner 扫描中同时计数和维护固定大小 top-N;
// 历史任务由 operator maintenance procedure 显式分批回填。
for row in ctx
.db
.external_generation_job_summary()
.by_external_generation_job_summary_owner_user_id()
.filter(&owner_user_id)
{
match row.status.as_str() {
EXTERNAL_GENERATION_STATUS_PENDING => pending_count = pending_count.saturating_add(1),
EXTERNAL_GENERATION_STATUS_RUNNING => running_count = running_count.saturating_add(1),
EXTERNAL_GENERATION_STATUS_COMPLETED | EXTERNAL_GENERATION_STATUS_FAILED => {
if row.notification_acknowledged_at.is_none() {
unacknowledged_terminal_count = unacknowledged_terminal_count.saturating_add(1);
}
}
_ => {}
}
let should_include = input.include_acknowledged_terminal
|| !is_external_generation_job_summary_terminal(&row)
|| row.notification_acknowledged_at.is_none();
let should_include_status =
status_filter.is_empty() || status_filter.iter().any(|status| row.status == *status);
if should_include && should_include_status {
retain_external_generation_job_summary_top_n(&mut rows, row, limit);
}
}
Ok(ExternalGenerationJobSummaryProcedureResult {
ok: true,
job: None,
jobs: rows
.into_iter()
.map(map_external_generation_job_summary_row)
.collect(),
pending_count,
running_count,
unacknowledged_terminal_count,
now_micros,
error_message: None,
})
}
fn acknowledge_external_generation_job_summaries_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobAcknowledgeInput,
) -> Result<ExternalGenerationJobSummaryProcedureResult, String> {
validate_required(
"external_generation_job.owner_user_id",
&input.owner_user_id,
)?;
let owner_user_id = input.owner_user_id.trim().to_string();
let acknowledged_at = Timestamp::from_micros_since_unix_epoch(input.acknowledged_at_micros);
let mut acknowledged = Vec::new();
for job_id in input.job_ids.iter().take(100) {
let normalized_job_id = job_id.trim().to_string();
let Some(mut row) = ctx
.db
.external_generation_job_summary()
.job_id()
.find(&normalized_job_id)
else {
continue;
};
if row.owner_user_id.trim() != owner_user_id
|| !is_external_generation_job_summary_terminal(&row)
{
continue;
}
if row.notification_acknowledged_at.is_some() {
continue;
}
row.notification_acknowledged_at = Some(acknowledged_at);
row.updated_at = acknowledged_at;
ctx.db
.external_generation_job_summary()
.job_id()
.update(row.clone());
insert_external_generation_job_summary_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_ACKNOWLEDGED,
Some("用户已确认任务通知".to_string()),
acknowledged_at,
);
acknowledged.push(map_external_generation_job_summary_row(row));
}
let (pending_count, running_count, unacknowledged_terminal_count) =
count_external_generation_job_summaries_for_owner(ctx, &owner_user_id);
Ok(ExternalGenerationJobSummaryProcedureResult {
ok: true,
job: None,
jobs: acknowledged,
pending_count,
running_count,
unacknowledged_terminal_count,
now_micros: ctx.timestamp.to_micros_since_unix_epoch(),
error_message: None,
})
}
fn backfill_external_generation_job_summaries_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobSummaryBackfillInput,
) -> Result<ExternalGenerationJobSummaryBackfillProcedureResult, String> {
let owner_user_id = input
.owner_user_id
.as_deref()
.and_then(normalize_optional_text);
let cursor_job_id = input
.cursor_job_id
.as_deref()
.and_then(normalize_optional_text);
let limit = input
.limit
.clamp(1, MAX_EXTERNAL_GENERATION_MAINTENANCE_BATCH_SIZE) as usize;
let cursor_range = external_generation_job_maintenance_cursor_range(cursor_job_id.as_deref());
let cursor_to_skip = cursor_job_id.clone();
let rows = ctx
.db
.external_generation_job()
.by_external_generation_job_cursor()
.filter(cursor_range)
.filter(move |row| {
cursor_to_skip
.as_deref()
.is_none_or(|cursor| row.job_id != cursor)
});
let (job_ids, next_cursor_job_id, has_more, scanned_count) =
select_external_generation_job_ids_for_maintenance(rows, limit, |row| {
owner_user_id
.as_deref()
.is_none_or(|owner| row.owner_user_id.trim() == owner)
&& ctx
.db
.external_generation_job_summary()
.job_id()
.find(&row.job_id)
.is_none()
});
let mut upserted_count = 0u32;
if !input.dry_run {
for job_id in &job_ids {
if let Some(row) = ctx.db.external_generation_job().job_id().find(job_id) {
persist_external_generation_job_summary(ctx, &row);
upserted_count = upserted_count.saturating_add(1);
}
}
}
Ok(ExternalGenerationJobSummaryBackfillProcedureResult {
ok: true,
dry_run: input.dry_run,
scanned_count,
selected_count: job_ids.len() as u32,
upserted_count,
next_cursor_job_id,
has_more,
error_message: None,
})
}
fn compact_external_generation_job_payloads_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobPayloadCompactionInput,
) -> Result<ExternalGenerationJobPayloadCompactionProcedureResult, String> {
let cursor_job_id = input
.cursor_job_id
.as_deref()
.and_then(normalize_optional_text);
let limit = input
.limit
.clamp(1, MAX_EXTERNAL_GENERATION_MAINTENANCE_BATCH_SIZE) as usize;
let cursor_range = external_generation_job_maintenance_cursor_range(cursor_job_id.as_deref());
let cursor_to_skip = cursor_job_id.clone();
let rows = ctx
.db
.external_generation_job()
.by_external_generation_job_source_cursor()
.filter((EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE, cursor_range))
.filter(move |row| {
cursor_to_skip
.as_deref()
.is_none_or(|cursor| row.job_id != cursor)
});
let (job_ids, next_cursor_job_id, has_more, scanned_count) =
select_external_generation_job_ids_for_maintenance(rows, limit, |row| {
should_compact_external_generation_job_payloads(row, input.completed_before_micros)
});
let mut matched_count = 0u32;
let mut updated_count = 0u32;
let mut before_bytes = 0u64;
let mut after_bytes = 0u64;
let mut inline_media_count = 0u64;
let mut invalid_json_count = 0u32;
for job_id in &job_ids {
let Some(mut row) = ctx.db.external_generation_job().job_id().find(job_id) else {
continue;
};
if !should_compact_external_generation_job_payloads(&row, input.completed_before_micros) {
continue;
}
let request_outcome = compact_external_generation_payload_json(&row.request_payload_json);
let result_outcome = row
.result_payload_json
.as_deref()
.map(compact_external_generation_payload_json);
invalid_json_count = invalid_json_count
.saturating_add(u32::from(request_outcome.invalid_json))
.saturating_add(u32::from(
result_outcome
.as_ref()
.is_some_and(|outcome| outcome.invalid_json),
));
let job_inline_media_count = request_outcome.inline_media_count.saturating_add(
result_outcome
.as_ref()
.map(|outcome| outcome.inline_media_count)
.unwrap_or(0),
);
if job_inline_media_count > 0 {
matched_count = matched_count.saturating_add(1);
inline_media_count = inline_media_count.saturating_add(job_inline_media_count);
for outcome in std::iter::once(&request_outcome).chain(result_outcome.iter()) {
if outcome.inline_media_count > 0 {
before_bytes = before_bytes.saturating_add(outcome.before_bytes);
after_bytes = after_bytes.saturating_add(outcome.after_bytes);
}
}
}
if input.dry_run {
continue;
}
if let Some(compacted_json) = request_outcome.compacted_json {
row.request_payload_json = compacted_json;
}
if let Some(outcome) = result_outcome {
if let Some(compacted_json) = outcome.compacted_json {
row.result_payload_json = Some(compacted_json);
}
}
if job_inline_media_count > 0 {
persist_external_generation_job_row(ctx, row);
updated_count = updated_count.saturating_add(1);
} else {
// 即使无需压缩,正式执行也顺带补齐该终态任务的轻量投影。
persist_external_generation_job_summary(ctx, &row);
}
}
Ok(ExternalGenerationJobPayloadCompactionProcedureResult {
ok: true,
dry_run: input.dry_run,
scanned_count,
matched_count,
updated_count,
before_bytes,
after_bytes,
inline_media_count,
invalid_json_count,
next_cursor_job_id,
has_more,
error_message: None,
})
}
fn prune_external_generation_job_history_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobRetentionInput,
) -> Result<ExternalGenerationJobRetentionProcedureResult, String> {
let source_module = input.source_module.trim().to_string();
validate_required("external_generation_job.source_module", &source_module)?;
let now_micros = ctx.timestamp.to_micros_since_unix_epoch();
if input.completed_before_micros > now_micros {
return Err(
"external_generation_job.completed_before_micros 不能晚于数据库当前时间".to_string(),
);
}
let cursor_job_id = input
.cursor_job_id
.as_deref()
.and_then(normalize_optional_text);
let limit = input
.limit
.clamp(1, MAX_EXTERNAL_GENERATION_MAINTENANCE_BATCH_SIZE) as usize;
let cursor_range = external_generation_job_maintenance_cursor_range(cursor_job_id.as_deref());
let cursor_to_skip = cursor_job_id.clone();
// 若 cursor 对应的任务仍存在,说明上一事务只删完了事件的一部分,或 dry-run
// 尚未执行 apply;下一次必须包含该任务继续清理,成功删除后它会自然消失。
let include_existing_cursor = cursor_job_id.as_deref().is_some_and(|cursor| {
ctx.db
.external_generation_job()
.job_id()
.find(&cursor.to_string())
.is_some()
});
let rows = ctx
.db
.external_generation_job()
.by_external_generation_job_source_cursor()
.filter((source_module.as_str(), cursor_range))
.filter(move |row| {
cursor_to_skip
.as_deref()
.is_none_or(|cursor| row.job_id != cursor || include_existing_cursor)
});
let (job_ids, next_cursor_job_id, has_more, scanned_count) =
select_external_generation_job_ids_for_maintenance(rows, limit, |row| {
ctx.db
.external_generation_job_summary()
.job_id()
.find(&row.job_id)
.is_some_and(|summary| {
is_external_generation_job_retention_candidate(
row,
&summary,
&source_module,
input.completed_before_micros,
)
})
});
let mut deleted_job_count = 0u32;
let mut deleted_summary_count = 0u32;
let mut deleted_event_count = 0u32;
let mut pending_event_cursor_job_id = None;
if !input.dry_run {
for job_id in &job_ids {
let Some(row) = ctx.db.external_generation_job().job_id().find(job_id) else {
continue;
};
let Some(summary) = ctx
.db
.external_generation_job_summary()
.job_id()
.find(job_id)
else {
continue;
};
if !is_external_generation_job_retention_candidate(
&row,
&summary,
&source_module,
input.completed_before_micros,
) {
continue;
}
let (deleted_for_job, has_more_events) =
delete_external_generation_job_events_for_job(ctx, job_id);
deleted_event_count = deleted_event_count.saturating_add(deleted_for_job);
if has_more_events {
pending_event_cursor_job_id = Some(job_id.clone());
break;
}
ctx.db
.external_generation_job_summary()
.job_id()
.delete(job_id);
deleted_summary_count = deleted_summary_count.saturating_add(1);
ctx.db.external_generation_job().job_id().delete(job_id);
deleted_job_count = deleted_job_count.saturating_add(1);
}
}
let (next_cursor_job_id, has_more) = pending_event_cursor_job_id
.map(|job_id| (Some(job_id), true))
.unwrap_or((next_cursor_job_id, has_more));
Ok(ExternalGenerationJobRetentionProcedureResult {
ok: true,
dry_run: input.dry_run,
scanned_count,
selected_count: job_ids.len() as u32,
deleted_job_count,
deleted_summary_count,
deleted_event_count,
next_cursor_job_id,
has_more,
error_message: None,
})
}
fn renew_external_generation_job_lease_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobRenewLeaseInput,
) -> Result<ExternalGenerationJobSnapshot, String> {
let mut row = get_worker_owned_external_generation_job(
ctx,
&input.job_id,
&input.worker_id,
&input.lease_token,
)?;
let renewed_at = ctx.timestamp;
let lease_duration_micros = duration_between_micros(
input.lease_expires_at_micros,
input.renewed_at_micros,
"external_generation_job.lease_duration",
)?;
row.lease_expires_at = Some(timestamp_after_micros(renewed_at, lease_duration_micros));
row.updated_at = renewed_at;
persist_external_generation_job_row(ctx, row.clone());
insert_external_generation_job_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_LEASE_RENEWED,
Some("worker 已续租任务".to_string()),
Some(input.worker_id),
renewed_at,
);
Ok(map_external_generation_job_row(row))
}
fn update_external_generation_job_phase_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobPhaseUpdateInput,
) -> Result<ExternalGenerationJobSnapshot, ExternalGenerationJobPhaseUpdateError> {
let phase = normalize_external_generation_job_phase(&input.phase)
.map_err(ExternalGenerationJobPhaseUpdateError::other)?;
let mut row = get_worker_owned_external_generation_job_for_phase_update(
ctx,
&input.job_id,
&input.worker_id,
&input.lease_token,
)?;
row.phase = Some(phase);
row.updated_at = ctx.timestamp;
persist_external_generation_job_row(ctx, row.clone());
Ok(map_external_generation_job_row(row))
}
fn fail_external_generation_job_tx(
ctx: &ReducerContext,
input: ExternalGenerationJobFailInput,
) -> Result<ExternalGenerationJobSnapshot, String> {
let error_message = normalize_external_generation_error_message(&input.error_message)
.ok_or_else(|| "external_generation_job.error_message 不能为空".to_string())?;
let mut row = get_worker_owned_external_generation_job(
ctx,
&input.job_id,
&input.worker_id,
&input.lease_token,
)?;
let failed_at = ctx.timestamp;
let refund_ledger_id = crate::settle_external_generation_attempt_refund(
ctx,
&row.job_id,
row.attempt,
&row.owner_user_id,
row.price_mud_points,
failed_at,
"worker_attempt_failed",
&row.job_kind,
&row.source_entity_id,
)?;
let requested_refund_ledger_id = input
.refund_ledger_id
.and_then(|value| normalize_optional_text(value.as_str()));
if requested_refund_ledger_id.is_some() && requested_refund_ledger_id != refund_ledger_id {
return Err("external_generation_job.refund_ledger_id 与当前 attempt 不匹配".to_string());
}
let retry_delay_micros = duration_between_micros(
input.retry_after_micros,
input.failed_at_micros,
"external_generation_job.retry_delay",
)?;
row.last_error_message = Some(error_message.clone());
row.refund_ledger_id = refund_ledger_id;
row.lease_expires_at = None;
row.worker_id = None;
row.lease_token = None;
row.updated_at = failed_at;
if row.attempt < row.max_attempts {
row.status = EXTERNAL_GENERATION_STATUS_PENDING.to_string();
row.available_at = timestamp_after_micros(failed_at, retry_delay_micros);
} else {
row.status = EXTERNAL_GENERATION_STATUS_FAILED.to_string();
row.completed_at = Some(failed_at);
}
persist_external_generation_job_row(ctx, row.clone());
insert_external_generation_job_event(
ctx,
&row,
EXTERNAL_GENERATION_EVENT_FAILED,
Some(error_message),
Some(input.worker_id),
failed_at,
);
Ok(map_external_generation_job_row(row))
}
fn get_external_generation_queue_stats_tx(
ctx: &ReducerContext,
) -> Result<ExternalGenerationQueueStatsSnapshot, String> {
let now = ctx.timestamp;
let now_micros = now.to_micros_since_unix_epoch();
let mut stats = ExternalGenerationQueueStatsSnapshot {
pending_count: 0,
delayed_pending_count: 0,
claimable_pending_count: 0,
running_active_count: 0,
expired_running_count: 0,
terminal_count: 0,
claimable_count: 0,
oldest_claimable_age_micros: None,
now_micros,
};
for row in ctx
.db
.external_generation_job()
.by_external_generation_job_status_available()
.filter(&EXTERNAL_GENERATION_STATUS_PENDING.to_string())
{
stats.pending_count = stats.pending_count.saturating_add(1);
if is_external_generation_job_claimable(&row, now) {
stats.claimable_pending_count = stats.claimable_pending_count.saturating_add(1);
record_external_generation_claimable_age(&mut stats, &row, now_micros);
} else {
stats.delayed_pending_count = stats.delayed_pending_count.saturating_add(1);
}
}
for row in ctx
.db
.external_generation_job()
.by_external_generation_job_status_available()
.filter(&EXTERNAL_GENERATION_STATUS_RUNNING.to_string())
{
if is_external_generation_job_claimable(&row, now) {
stats.expired_running_count = stats.expired_running_count.saturating_add(1);
record_external_generation_claimable_age(&mut stats, &row, now_micros);
} else {
stats.running_active_count = stats.running_active_count.saturating_add(1);
}
}
stats.claimable_count = stats
.claimable_pending_count
.saturating_add(stats.expired_running_count);
Ok(stats)
}
pub(crate) fn validate_editor_generation_job_lease_for_result_tx(
ctx: &ReducerContext,
input: &ExternalGenerationJobCompleteInput,
expected_owner_user_id: &str,
expected_operation_kind: &str,
expected_operation_fingerprint: &str,
expected_project_id: Option<&str>,
) -> Result<ExternalGenerationJobSnapshot, String> {
let row = get_worker_owned_external_generation_job(
ctx,
&input.job_id,
&input.worker_id,
&input.lease_token,
)?;
validate_editor_generation_job_business_identity(
&row,
expected_owner_user_id,
expected_operation_kind,
expected_operation_fingerprint,
expected_project_id,
)?;
Ok(map_external_generation_job_row(row))
}
pub(crate) fn validate_completed_editor_generation_job_replay_tx(
ctx: &ReducerContext,
input: &ExternalGenerationJobCompleteInput,
expected_owner_user_id: &str,
expected_operation_kind: &str,
expected_operation_fingerprint: &str,
expected_project_id: Option<&str>,
) -> Result<ExternalGenerationJobSnapshot, String> {
validate_required("external_generation_job.job_id", &input.job_id)?;
let job_id = input.job_id.trim().to_string();
let row = ctx
.db
.external_generation_job()
.job_id()
.find(&job_id)
.ok_or_else(|| "external_generation_job 不存在".to_string())?;
validate_editor_generation_job_business_identity(
&row,
expected_owner_user_id,
expected_operation_kind,
expected_operation_fingerprint,
expected_project_id,
)?;
if row.status != EXTERNAL_GENERATION_STATUS_COMPLETED {
return Err("external_generation_job 尚未完成,不能作为生成结果重放".to_string());
}
let expected_result_payload_json = if is_external_generation_editor_source(&row.source_module) {
validate_optional_external_generation_payload_json(
"external_generation_job.result_payload_json",
input.result_payload_json.as_deref(),
)?
} else {
input
.result_payload_json
.as_deref()
.and_then(normalize_optional_text)
};
if row.result_payload_json != expected_result_payload_json {
return Err("external_generation_job 已完成结果与重放内容不一致".to_string());
}
Ok(map_external_generation_job_row(row))
}
fn validate_editor_generation_job_business_identity(
row: &ExternalGenerationJob,
expected_owner_user_id: &str,
expected_operation_kind: &str,
expected_operation_fingerprint: &str,
expected_project_id: Option<&str>,
) -> Result<(), String> {
if row.owner_user_id.trim() != expected_owner_user_id.trim() {
return Err("external_generation_job owner_user_id 与业务写回不匹配".to_string());
}
if row.job_kind.trim() != expected_operation_kind.trim() {
return Err("external_generation_job job_kind 与业务写回不匹配".to_string());
}
if !is_external_generation_editor_source(&row.source_module) {
return Err("external_generation_job source_module 不是编辑器生成链路".to_string());
}
if let Some(expected_project_id) = expected_project_id
&& row.source_entity_id.trim() != expected_project_id.trim()
{
return Err("external_generation_job source_entity_id 与生成结果项目不匹配".to_string());
}
let actual_fingerprint = shared_contracts::editor_generation_request_fingerprint(
row.job_kind.as_str(),
row.request_payload_json.as_str(),
);
if actual_fingerprint != expected_operation_fingerprint.trim() {
return Err("external_generation_job 请求 fingerprint 与业务写回不匹配".to_string());
}
Ok(())
}
fn get_worker_owned_external_generation_job(
ctx: &ReducerContext,
job_id: &str,
worker_id: &str,
lease_token: &str,
) -> Result<ExternalGenerationJob, String> {
get_worker_owned_external_generation_job_for_phase_update(ctx, job_id, worker_id, lease_token)
.map_err(|error| error.message)
}
#[derive(Debug)]
struct ExternalGenerationJobPhaseUpdateError {
kind: ExternalGenerationJobPhaseUpdateFailureKind,
message: String,
}
impl ExternalGenerationJobPhaseUpdateError {
fn lease_fencing(message: impl Into<String>) -> Self {
Self {
kind: ExternalGenerationJobPhaseUpdateFailureKind::LeaseFencingRejected,
message: message.into(),
}
}
fn other(message: impl Into<String>) -> Self {
Self {
kind: ExternalGenerationJobPhaseUpdateFailureKind::OtherRejected,
message: message.into(),
}
}
}
impl std::fmt::Display for ExternalGenerationJobPhaseUpdateError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(self.message.as_str())
}
}
impl std::error::Error for ExternalGenerationJobPhaseUpdateError {}
fn get_worker_owned_external_generation_job_for_phase_update(
ctx: &ReducerContext,
job_id: &str,
worker_id: &str,
lease_token: &str,
) -> Result<ExternalGenerationJob, ExternalGenerationJobPhaseUpdateError> {
validate_required("external_generation_job.job_id", job_id)
.map_err(ExternalGenerationJobPhaseUpdateError::other)?;
validate_required("external_generation_job.worker_id", worker_id)
.map_err(ExternalGenerationJobPhaseUpdateError::other)?;
validate_required("external_generation_job.lease_token", lease_token)
.map_err(ExternalGenerationJobPhaseUpdateError::other)?;
let row = ctx
.db
.external_generation_job()
.job_id()
.find(&job_id.trim().to_string())
.ok_or_else(|| {
ExternalGenerationJobPhaseUpdateError::lease_fencing("external_generation_job 不存在")
})?;
validate_external_generation_job_phase_update_lease(
&row,
worker_id,
lease_token,
ctx.timestamp,
)?;
Ok(row)
}
fn validate_external_generation_job_phase_update_lease(
row: &ExternalGenerationJob,
worker_id: &str,
lease_token: &str,
now: Timestamp,
) -> Result<(), ExternalGenerationJobPhaseUpdateError> {
if row.status != EXTERNAL_GENERATION_STATUS_RUNNING {
return Err(ExternalGenerationJobPhaseUpdateError::lease_fencing(
"external_generation_job 当前不是 running 状态",
));
}
if !is_external_generation_job_owned_by_worker(&row, worker_id) {
return Err(ExternalGenerationJobPhaseUpdateError::lease_fencing(
"external_generation_job worker lease 不匹配",
));
}
if !is_external_generation_job_owned_by_lease_token(&row, lease_token) {
return Err(ExternalGenerationJobPhaseUpdateError::lease_fencing(
"external_generation_job lease token 不匹配",
));
}
if !is_external_generation_job_lease_active(row, now) {
return Err(ExternalGenerationJobPhaseUpdateError::lease_fencing(
"external_generation_job lease 已过期",
));
}
Ok(())
}
fn is_external_generation_job_owned_by_worker(
row: &ExternalGenerationJob,
worker_id: &str,
) -> bool {
row.worker_id.as_deref() == Some(worker_id.trim())
}
fn is_external_generation_job_owned_by_lease_token(
row: &ExternalGenerationJob,
lease_token: &str,
) -> bool {
row.lease_token.as_deref() == Some(lease_token.trim())
}
fn is_external_generation_job_lease_active(row: &ExternalGenerationJob, now: Timestamp) -> bool {
row.lease_expires_at
.map(|lease_expires_at| lease_expires_at > now)
.unwrap_or(false)
}
fn is_external_generation_job_claimable(row: &ExternalGenerationJob, now: Timestamp) -> bool {
match row.status.as_str() {
EXTERNAL_GENERATION_STATUS_PENDING => row.available_at <= now,
EXTERNAL_GENERATION_STATUS_RUNNING => row
.lease_expires_at
.map(|lease_expires_at| lease_expires_at <= now)
.unwrap_or(true),
EXTERNAL_GENERATION_STATUS_COMPLETED
| EXTERNAL_GENERATION_STATUS_FAILED
| EXTERNAL_GENERATION_STATUS_CANCELLED => false,
_ => false,
}
}
fn external_generation_job_has_exhausted_attempts(row: &ExternalGenerationJob) -> bool {
row.attempt >= row.max_attempts
}
fn is_external_generation_job_terminal(row: &ExternalGenerationJob) -> bool {
matches!(
row.status.as_str(),
EXTERNAL_GENERATION_STATUS_COMPLETED
| EXTERNAL_GENERATION_STATUS_FAILED
| EXTERNAL_GENERATION_STATUS_CANCELLED
)
}
fn is_external_generation_job_summary_terminal(row: &ExternalGenerationJobSummary) -> bool {
matches!(
row.status.as_str(),
EXTERNAL_GENERATION_STATUS_COMPLETED
| EXTERNAL_GENERATION_STATUS_FAILED
| EXTERNAL_GENERATION_STATUS_CANCELLED
)
}
fn should_compact_external_generation_job_payloads(
row: &ExternalGenerationJob,
completed_before_micros: Option<i64>,
) -> bool {
if !is_external_generation_editor_source(&row.source_module)
|| !is_external_generation_job_terminal(row)
{
return false;
}
completed_before_micros.is_none_or(|cutoff| {
row.completed_at
.unwrap_or(row.updated_at)
.to_micros_since_unix_epoch()
<= cutoff
})
}
fn is_external_generation_job_retention_candidate(
row: &ExternalGenerationJob,
summary: &ExternalGenerationJobSummary,
source_module: &str,
completed_before_micros: i64,
) -> bool {
row.source_module.trim() == source_module.trim()
&& summary.job_id == row.job_id
&& summary.status == row.status
&& is_external_generation_job_terminal(row)
&& is_external_generation_job_summary_terminal(summary)
&& summary.notification_acknowledged_at.is_some()
&& row
.completed_at
.unwrap_or(row.updated_at)
.to_micros_since_unix_epoch()
<= completed_before_micros
}
fn external_generation_job_maintenance_cursor_range(
cursor_job_id: Option<&str>,
) -> RangeFrom<&str> {
cursor_job_id.unwrap_or_default()..
}
fn select_external_generation_job_ids_for_maintenance(
rows: impl Iterator<Item = ExternalGenerationJob>,
limit: usize,
mut should_select: impl FnMut(&ExternalGenerationJob) -> bool,
) -> (Vec<String>, Option<String>, bool, u64) {
let mut rows = rows.peekable();
let mut selected_job_ids = Vec::with_capacity(limit);
let mut next_cursor_job_id = None;
let mut scanned_count = 0u64;
// 游标选择阶段最多反序列化 limit + 1 条大 payload row;apply 随后按主键逐条
// 重新读取选中行,避免把整批大 payload 同时保留在事务内存中。
for row in rows.by_ref().take(limit) {
scanned_count = scanned_count.saturating_add(1);
next_cursor_job_id = Some(row.job_id.clone());
if should_select(&row) {
selected_job_ids.push(row.job_id);
}
}
let has_more = rows.peek().is_some();
(
selected_job_ids,
next_cursor_job_id,
has_more,
scanned_count,
)
}
fn delete_external_generation_job_events_for_job(
ctx: &ReducerContext,
job_id: &str,
) -> (u32, bool) {
// 每次 procedure 最多删除一个固定批次;若仍有事件,保留 job/summary,调用方
// 通过同一个 job cursor 重试,避免单个任务把整段审计历史塞进一个事务写集。
let event_ids = ctx
.db
.external_generation_job_event()
.by_external_generation_job_event_job_id_only()
.filter(job_id)
.take(EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE + 1)
.map(|event| event.event_id.clone())
.collect::<Vec<_>>();
let has_more = event_ids.len() > EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE;
let deleted_count = event_ids
.len()
.min(EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE) as u32;
for event_id in event_ids
.into_iter()
.take(EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE)
{
ctx.db
.external_generation_job_event()
.event_id()
.delete(&event_id);
}
(deleted_count, has_more)
}
fn count_external_generation_job_summaries_for_owner(
ctx: &ReducerContext,
owner_user_id: &str,
) -> (u32, u32, u32) {
let normalized_owner_user_id = owner_user_id.trim().to_string();
let mut pending_count = 0u32;
let mut running_count = 0u32;
let mut unacknowledged_terminal_count = 0u32;
for row in ctx
.db
.external_generation_job_summary()
.by_external_generation_job_summary_owner_user_id()
.filter(&normalized_owner_user_id)
{
match row.status.as_str() {
EXTERNAL_GENERATION_STATUS_PENDING => pending_count = pending_count.saturating_add(1),
EXTERNAL_GENERATION_STATUS_RUNNING => running_count = running_count.saturating_add(1),
EXTERNAL_GENERATION_STATUS_COMPLETED | EXTERNAL_GENERATION_STATUS_FAILED => {
if row.notification_acknowledged_at.is_none() {
unacknowledged_terminal_count = unacknowledged_terminal_count.saturating_add(1);
}
}
_ => {}
}
}
(pending_count, running_count, unacknowledged_terminal_count)
}
#[cfg(test)]
fn external_generation_job_sort_bucket(row: &ExternalGenerationJob) -> u8 {
match row.status.as_str() {
EXTERNAL_GENERATION_STATUS_RUNNING => 0,
EXTERNAL_GENERATION_STATUS_PENDING => 1,
EXTERNAL_GENERATION_STATUS_COMPLETED | EXTERNAL_GENERATION_STATUS_FAILED
if row.notification_acknowledged_at.is_none() =>
{
2
}
EXTERNAL_GENERATION_STATUS_COMPLETED | EXTERNAL_GENERATION_STATUS_FAILED => 3,
EXTERNAL_GENERATION_STATUS_CANCELLED if row.notification_acknowledged_at.is_none() => 4,
EXTERNAL_GENERATION_STATUS_CANCELLED => 5,
_ => 6,
}
}
#[cfg(test)]
fn external_generation_job_sort_time_micros(row: &ExternalGenerationJob) -> i64 {
if is_external_generation_job_terminal(row) {
return row
.completed_at
.unwrap_or(row.updated_at)
.to_micros_since_unix_epoch();
}
row.updated_at.to_micros_since_unix_epoch()
}
fn external_generation_job_summary_sort_bucket(row: &ExternalGenerationJobSummary) -> u8 {
match row.status.as_str() {
EXTERNAL_GENERATION_STATUS_RUNNING => 0,
EXTERNAL_GENERATION_STATUS_PENDING => 1,
EXTERNAL_GENERATION_STATUS_COMPLETED | EXTERNAL_GENERATION_STATUS_FAILED
if row.notification_acknowledged_at.is_none() =>
{
2
}
EXTERNAL_GENERATION_STATUS_COMPLETED | EXTERNAL_GENERATION_STATUS_FAILED => 3,
EXTERNAL_GENERATION_STATUS_CANCELLED if row.notification_acknowledged_at.is_none() => 4,
EXTERNAL_GENERATION_STATUS_CANCELLED => 5,
_ => 6,
}
}
fn external_generation_job_summary_sort_time_micros(row: &ExternalGenerationJobSummary) -> i64 {
if is_external_generation_job_summary_terminal(row) {
return row
.completed_at
.unwrap_or(row.updated_at)
.to_micros_since_unix_epoch();
}
row.updated_at.to_micros_since_unix_epoch()
}
fn compare_external_generation_job_summaries(
left: &ExternalGenerationJobSummary,
right: &ExternalGenerationJobSummary,
) -> Ordering {
external_generation_job_summary_sort_bucket(left)
.cmp(&external_generation_job_summary_sort_bucket(right))
.then_with(|| {
external_generation_job_summary_sort_time_micros(right)
.cmp(&external_generation_job_summary_sort_time_micros(left))
})
.then_with(|| left.job_id.cmp(&right.job_id))
}
fn retain_external_generation_job_summary_top_n(
rows: &mut Vec<ExternalGenerationJobSummary>,
row: ExternalGenerationJobSummary,
limit: usize,
) {
if limit == 0 {
return;
}
let insert_at = rows
.binary_search_by(|existing| compare_external_generation_job_summaries(existing, &row))
.unwrap_or_else(|index| index);
if insert_at >= limit {
return;
}
rows.insert(insert_at, row);
if rows.len() > limit {
rows.pop();
}
}
fn normalize_external_generation_job_status_filter(statuses: &[String]) -> Vec<&'static str> {
statuses
.iter()
.filter_map(|status| match status.trim() {
"queued" | EXTERNAL_GENERATION_STATUS_PENDING => {
Some(EXTERNAL_GENERATION_STATUS_PENDING)
}
EXTERNAL_GENERATION_STATUS_RUNNING => Some(EXTERNAL_GENERATION_STATUS_RUNNING),
EXTERNAL_GENERATION_STATUS_COMPLETED => Some(EXTERNAL_GENERATION_STATUS_COMPLETED),
EXTERNAL_GENERATION_STATUS_FAILED => Some(EXTERNAL_GENERATION_STATUS_FAILED),
EXTERNAL_GENERATION_STATUS_CANCELLED => Some(EXTERNAL_GENERATION_STATUS_CANCELLED),
_ => None,
})
.collect()
}
fn normalize_external_generation_job_phase(phase: &str) -> Result<String, String> {
match phase.trim() {
EXTERNAL_GENERATION_PHASE_GENERATING => {
Ok(EXTERNAL_GENERATION_PHASE_GENERATING.to_string())
}
EXTERNAL_GENERATION_PHASE_PROCESSING => {
Ok(EXTERNAL_GENERATION_PHASE_PROCESSING.to_string())
}
_ => Err("external_generation_job.phase 只支持 generating 或 processing".to_string()),
}
}
fn record_external_generation_claimable_age(
stats: &mut ExternalGenerationQueueStatsSnapshot,
row: &ExternalGenerationJob,
now_micros: i64,
) {
let age = now_micros
.saturating_sub(row.available_at.to_micros_since_unix_epoch())
.max(0);
stats.oldest_claimable_age_micros = Some(
stats
.oldest_claimable_age_micros
.map(|current| current.max(age))
.unwrap_or(age),
);
}
fn persist_external_generation_job_row(ctx: &ReducerContext, row: ExternalGenerationJob) {
ctx.db
.external_generation_job()
.job_id()
.delete(&row.job_id);
ctx.db.external_generation_job().insert(row.clone());
persist_external_generation_job_summary(ctx, &row);
}
fn persist_external_generation_job_summary(
ctx: &ReducerContext,
row: &ExternalGenerationJob,
) -> ExternalGenerationJobSummary {
let existing = ctx
.db
.external_generation_job_summary()
.job_id()
.find(&row.job_id);
let cached_request_prompt = existing
.as_ref()
.map(|summary| summary.request_prompt.clone());
let cached_notification_acknowledged_at = existing
.as_ref()
.and_then(|summary| summary.notification_acknowledged_at);
if existing.is_some() {
ctx.db
.external_generation_job_summary()
.job_id()
.delete(&row.job_id);
}
let mut summary = build_external_generation_job_summary_row(row, cached_request_prompt);
if summary.notification_acknowledged_at.is_none() {
summary.notification_acknowledged_at = cached_notification_acknowledged_at;
if let Some(acknowledged_at) = cached_notification_acknowledged_at {
summary.updated_at = summary.updated_at.max(acknowledged_at);
}
}
ctx.db
.external_generation_job_summary()
.insert(summary.clone());
summary
}
fn build_external_generation_job_summary_row(
row: &ExternalGenerationJob,
cached_request_prompt: Option<Option<String>>,
) -> ExternalGenerationJobSummary {
let request_payload = serde_json::from_str::<serde_json::Value>(&row.request_payload_json).ok();
let is_scene_generation = row.source_module.trim() == EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE
&& row.job_kind.trim() == "editor_image_generation"
&& request_payload
.as_ref()
.is_some_and(external_generation_payload_is_scene);
ExternalGenerationJobSummary {
job_id: row.job_id.clone(),
job_kind: row.job_kind.clone(),
owner_user_id: row.owner_user_id.clone(),
source_module: row.source_module.clone(),
source_entity_id: row.source_entity_id.clone(),
request_label: if is_scene_generation {
"图片画布生成游戏场景".to_string()
} else {
row.request_label.clone()
},
request_prompt: match cached_request_prompt.filter(|_| !is_scene_generation) {
Some(prompt) => prompt
.as_deref()
.and_then(normalize_external_generation_request_prompt_text),
None => request_payload.as_ref().and_then(|payload| {
extract_external_generation_request_prompt(payload, is_scene_generation)
}),
},
status: row.status.clone(),
last_error_message: row
.last_error_message
.as_deref()
.and_then(normalize_external_generation_error_message),
created_at: row.created_at,
started_at: row.started_at,
completed_at: row.completed_at,
updated_at: row.updated_at,
price_mud_points: row.price_mud_points,
refund_ledger_id: row.refund_ledger_id.clone(),
notification_acknowledged_at: row.notification_acknowledged_at,
warning_message: extract_external_generation_warning_message(
row.result_payload_json.as_deref(),
),
phase: row.phase.clone(),
}
}
fn external_generation_payload_is_scene(payload: &serde_json::Value) -> bool {
payload
.get("kind")
.and_then(serde_json::Value::as_str)
.is_some_and(|kind| kind.trim() == "scene")
}
fn extract_external_generation_warning_message(
result_payload_json: Option<&str>,
) -> Option<String> {
let payload: serde_json::Value = serde_json::from_str(result_payload_json?.trim()).ok()?;
let warning = payload
.get("warning")
.or_else(|| payload.get("sliceWarning"))?;
let reason = warning.get("reason").and_then(serde_json::Value::as_str)?;
normalize_external_generation_warning_message(reason)
}
fn extract_external_generation_request_prompt(
payload: &serde_json::Value,
is_scene_generation: bool,
) -> Option<String> {
if is_scene_generation {
return payload
.get("generationInputs")
.and_then(|value| value.get("fields"))
.and_then(serde_json::Value::as_array)
.and_then(|fields| {
fields.iter().find_map(|field| {
if field
.get("title")
.and_then(serde_json::Value::as_str)?
.trim()
!= "画面内容"
{
return None;
}
field
.get("value")
.and_then(serde_json::Value::as_str)
.and_then(normalize_external_generation_request_prompt_text)
})
});
}
for key in ["prompt", "promptText", "spritesheetLabel"] {
if let Some(prompt) = payload
.get(key)
.and_then(serde_json::Value::as_str)
.and_then(normalize_external_generation_request_prompt_text)
{
return Some(prompt);
}
}
if let Some(prompt) = payload
.get("iconDescriptions")
.and_then(serde_json::Value::as_array)
.map(|items| {
items
.iter()
.filter_map(serde_json::Value::as_str)
.filter_map(normalize_external_generation_request_prompt_text)
.collect::<Vec<_>>()
.join("、")
})
.and_then(|value| normalize_external_generation_request_prompt_text(&value))
{
return Some(prompt);
}
payload
.get("generationInputs")
.and_then(|value| value.get("fields"))
.and_then(serde_json::Value::as_array)
.and_then(|fields| {
fields.iter().find_map(|field| {
let title = field
.get("title")
.and_then(serde_json::Value::as_str)?
.trim();
if !matches!(title, "prompt" | "gpt_description_prompt") {
return None;
}
field
.get("value")
.and_then(serde_json::Value::as_str)
.and_then(normalize_external_generation_request_prompt_text)
})
})
}
fn normalize_external_generation_request_prompt_text(value: &str) -> Option<String> {
let trimmed = value.trim();
if trimmed.is_empty() || is_external_generation_inline_media_reference(trimmed) {
return None;
}
let mut chars = trimmed.chars();
let mut normalized = chars
.by_ref()
.take(MAX_EXTERNAL_GENERATION_REQUEST_PROMPT_CHARS)
.collect::<String>();
if chars.next().is_some() {
normalized.push('…');
}
Some(normalized)
}
fn normalize_external_generation_error_message(value: &str) -> Option<String> {
let trimmed = value.trim();
if trimmed.is_empty() {
return None;
}
if contains_external_generation_inline_media_reference_text(trimmed) {
return Some(INLINE_MEDIA_ERROR_REDACTED_MESSAGE.to_string());
}
let mut chars = trimmed.chars();
let mut normalized = chars
.by_ref()
.take(MAX_EXTERNAL_GENERATION_ERROR_MESSAGE_CHARS)
.collect::<String>();
if chars.next().is_some() {
normalized.push('…');
}
Some(normalized)
}
fn normalize_external_generation_warning_message(value: &str) -> Option<String> {
let trimmed = value.trim();
if trimmed.is_empty() {
return None;
}
if contains_external_generation_inline_media_reference_text(trimmed) {
return Some(INLINE_MEDIA_WARNING_REDACTED_MESSAGE.to_string());
}
let mut chars = trimmed.chars();
let mut normalized = chars
.by_ref()
.take(MAX_EXTERNAL_GENERATION_WARNING_MESSAGE_CHARS)
.collect::<String>();
if chars.next().is_some() {
normalized.push('…');
}
Some(normalized)
}
fn contains_external_generation_inline_media_reference_text(value: &str) -> bool {
let normalized = value.to_ascii_lowercase();
["data:", "blob:"].iter().any(|scheme| {
normalized.match_indices(scheme).any(|(index, _)| {
index == 0
|| normalized
.as_bytes()
.get(index - 1)
.is_some_and(|previous| {
matches!(
previous,
b' ' | b'\t'
| b'\r'
| b'\n'
| b'"'
| b'\''
| b'='
| b'('
| b'['
| b'{'
| b','
)
})
})
})
}
fn insert_external_generation_job_event(
ctx: &ReducerContext,
row: &ExternalGenerationJob,
event_kind: &str,
message: Option<String>,
worker_id: Option<String>,
created_at: Timestamp,
) {
let event_id = build_external_generation_event_id(row, event_kind, created_at);
if ctx
.db
.external_generation_job_event()
.event_id()
.find(&event_id)
.is_some()
{
return;
}
ctx.db
.external_generation_job_event()
.insert(ExternalGenerationJobEvent {
event_id,
job_id: row.job_id.clone(),
owner_user_id: row.owner_user_id.clone(),
event_kind: event_kind.to_string(),
status: row.status.clone(),
message,
worker_id,
created_at,
});
}
fn insert_external_generation_job_summary_event(
ctx: &ReducerContext,
row: &ExternalGenerationJobSummary,
event_kind: &str,
message: Option<String>,
created_at: Timestamp,
) {
let event_id = format!(
"{}:{}:{}:summary:{}",
row.job_id.trim(),
event_kind.trim(),
row.status.trim(),
created_at.to_micros_since_unix_epoch()
);
if ctx
.db
.external_generation_job_event()
.event_id()
.find(&event_id)
.is_some()
{
return;
}
ctx.db
.external_generation_job_event()
.insert(ExternalGenerationJobEvent {
event_id,
job_id: row.job_id.clone(),
owner_user_id: row.owner_user_id.clone(),
event_kind: event_kind.to_string(),
status: row.status.clone(),
message,
worker_id: None,
created_at,
});
}
fn map_external_generation_job_row(row: ExternalGenerationJob) -> ExternalGenerationJobSnapshot {
let notification_acknowledged_at_micros = row
.notification_acknowledged_at
.map(|value| value.to_micros_since_unix_epoch());
ExternalGenerationJobSnapshot {
job_id: row.job_id,
dedupe_key: row.dedupe_key,
job_kind: row.job_kind,
owner_user_id: row.owner_user_id,
source_module: row.source_module,
source_entity_id: row.source_entity_id,
request_label: row.request_label,
request_payload_json: row.request_payload_json,
status: row.status,
attempt: row.attempt,
max_attempts: row.max_attempts,
last_error_message: row.last_error_message,
worker_id: row.worker_id,
lease_expires_at_micros: row
.lease_expires_at
.map(|value| value.to_micros_since_unix_epoch()),
available_at_micros: row.available_at.to_micros_since_unix_epoch(),
result_payload_json: row.result_payload_json,
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
started_at_micros: row
.started_at
.map(|value| value.to_micros_since_unix_epoch()),
completed_at_micros: row
.completed_at
.map(|value| value.to_micros_since_unix_epoch()),
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
lease_token: row.lease_token,
price_mud_points: row.price_mud_points,
refund_ledger_id: row.refund_ledger_id,
notification_acknowledged_at_micros,
phase: row.phase,
}
}
fn map_external_generation_job_summary_row(
row: ExternalGenerationJobSummary,
) -> ExternalGenerationJobSummarySnapshot {
ExternalGenerationJobSummarySnapshot {
job_id: row.job_id,
job_kind: row.job_kind,
owner_user_id: row.owner_user_id,
source_module: row.source_module,
source_entity_id: row.source_entity_id,
request_label: row.request_label,
request_prompt: row.request_prompt,
status: row.status,
last_error_message: row.last_error_message,
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
started_at_micros: row
.started_at
.map(|value| value.to_micros_since_unix_epoch()),
completed_at_micros: row
.completed_at
.map(|value| value.to_micros_since_unix_epoch()),
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
price_mud_points: row.price_mud_points,
refund_ledger_id: row.refund_ledger_id,
notification_acknowledged_at_micros: row
.notification_acknowledged_at
.map(|value| value.to_micros_since_unix_epoch()),
warning_message: row.warning_message,
phase: row.phase,
}
}
fn map_external_generation_job_summary_to_compat_snapshot(
summary: ExternalGenerationJobSummarySnapshot,
) -> ExternalGenerationJobSnapshot {
let request_payload_json = summary
.request_prompt
.as_ref()
.map(|prompt| serde_json::json!({ "prompt": prompt }).to_string())
.unwrap_or_else(|| "{}".to_string());
ExternalGenerationJobSnapshot {
job_id: summary.job_id,
dedupe_key: String::new(),
job_kind: summary.job_kind,
owner_user_id: summary.owner_user_id,
source_module: summary.source_module,
source_entity_id: summary.source_entity_id,
request_label: summary.request_label,
request_payload_json,
status: summary.status,
attempt: 0,
max_attempts: 0,
last_error_message: summary.last_error_message,
worker_id: None,
lease_expires_at_micros: None,
available_at_micros: summary.updated_at_micros,
result_payload_json: None,
created_at_micros: summary.created_at_micros,
started_at_micros: summary.started_at_micros,
completed_at_micros: summary.completed_at_micros,
updated_at_micros: summary.updated_at_micros,
lease_token: None,
price_mud_points: summary.price_mud_points,
refund_ledger_id: summary.refund_ledger_id,
notification_acknowledged_at_micros: summary.notification_acknowledged_at_micros,
phase: summary.phase,
}
}
fn single_external_generation_job_result(
job: ExternalGenerationJobSnapshot,
) -> ExternalGenerationJobProcedureResult {
ExternalGenerationJobProcedureResult {
ok: true,
job: Some(job),
jobs: Vec::new(),
pending_count: 0,
running_count: 0,
unacknowledged_terminal_count: 0,
now_micros: 0,
error_message: None,
}
}
fn failed_external_generation_job_result(message: String) -> ExternalGenerationJobProcedureResult {
ExternalGenerationJobProcedureResult {
ok: false,
job: None,
jobs: Vec::new(),
pending_count: 0,
running_count: 0,
unacknowledged_terminal_count: 0,
now_micros: 0,
error_message: Some(message),
}
}
fn single_external_generation_job_result_read_result(
result: ExternalGenerationJobResultSnapshot,
) -> ExternalGenerationJobResultProcedureResult {
ExternalGenerationJobResultProcedureResult {
ok: true,
result: Some(result),
error_message: None,
}
}
fn failed_external_generation_job_result_read_result(
message: String,
) -> ExternalGenerationJobResultProcedureResult {
ExternalGenerationJobResultProcedureResult {
ok: false,
result: None,
error_message: Some(message),
}
}
fn single_external_generation_job_summary_result(
job: ExternalGenerationJobSummarySnapshot,
) -> ExternalGenerationJobSummaryProcedureResult {
ExternalGenerationJobSummaryProcedureResult {
ok: true,
job: Some(job),
jobs: Vec::new(),
pending_count: 0,
running_count: 0,
unacknowledged_terminal_count: 0,
now_micros: 0,
error_message: None,
}
}
fn failed_external_generation_job_summary_result(
message: String,
) -> ExternalGenerationJobSummaryProcedureResult {
ExternalGenerationJobSummaryProcedureResult {
ok: false,
job: None,
jobs: Vec::new(),
pending_count: 0,
running_count: 0,
unacknowledged_terminal_count: 0,
now_micros: 0,
error_message: Some(message),
}
}
fn failed_external_generation_job_summary_backfill_result(
dry_run: bool,
message: String,
) -> ExternalGenerationJobSummaryBackfillProcedureResult {
ExternalGenerationJobSummaryBackfillProcedureResult {
ok: false,
dry_run,
scanned_count: 0,
selected_count: 0,
upserted_count: 0,
next_cursor_job_id: None,
has_more: false,
error_message: Some(message),
}
}
fn failed_external_generation_job_payload_compaction_result(
dry_run: bool,
message: String,
) -> ExternalGenerationJobPayloadCompactionProcedureResult {
ExternalGenerationJobPayloadCompactionProcedureResult {
ok: false,
dry_run,
scanned_count: 0,
matched_count: 0,
updated_count: 0,
before_bytes: 0,
after_bytes: 0,
inline_media_count: 0,
invalid_json_count: 0,
next_cursor_job_id: None,
has_more: false,
error_message: Some(message),
}
}
fn failed_external_generation_job_retention_result(
dry_run: bool,
message: String,
) -> ExternalGenerationJobRetentionProcedureResult {
ExternalGenerationJobRetentionProcedureResult {
ok: false,
dry_run,
scanned_count: 0,
selected_count: 0,
deleted_job_count: 0,
deleted_summary_count: 0,
deleted_event_count: 0,
next_cursor_job_id: None,
has_more: false,
error_message: Some(message),
}
}
fn validate_required(field: &str, value: &str) -> Result<(), String> {
if value.trim().is_empty() {
return Err(format!("{field} 不能为空"));
}
Ok(())
}
fn validate_external_generation_persisted_payload_for_source(
source_module: &str,
field: &str,
value: &str,
) -> Result<String, String> {
if is_external_generation_editor_source(source_module) {
validate_external_generation_payload_json(field, value)
} else {
validate_required(field, value)?;
Ok(value.trim().to_string())
}
}
fn is_external_generation_editor_source(source_module: &str) -> bool {
source_module.trim() == EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE
}
fn is_external_generation_job_worker_claimable(
row: &ExternalGenerationJob,
now: Timestamp,
) -> bool {
is_external_generation_editor_source(&row.source_module)
&& is_external_generation_job_claimable(row, now)
}
fn validate_external_generation_payload_json(field: &str, value: &str) -> Result<String, String> {
let normalized = value.trim();
if normalized.is_empty() {
return Err(format!("{field} 不能为空"));
}
let payload_bytes = normalized.len();
if payload_bytes > MAX_EXTERNAL_GENERATION_PAYLOAD_BYTES {
return Err(format!(
"{field} JSON 大小为 {payload_bytes} 字节,超过持久化上限 {MAX_EXTERNAL_GENERATION_PAYLOAD_BYTES} 字节;请移除冗余数据并改传 objectKey 或 resourceId"
));
}
let payload = serde_json::from_str::<serde_json::Value>(normalized)
.map_err(|error| format!("{field} 不是合法 JSON: {error}"))?;
if contains_external_generation_inline_media_reference(&payload) {
return Err(format!(
"{field} 禁止包含 data: 或 blob: 内联媒体引用,请先上传对象存储并改传 objectKey 或 resourceId"
));
}
Ok(normalized.to_string())
}
fn validate_optional_external_generation_payload_json(
field: &str,
value: Option<&str>,
) -> Result<Option<String>, String> {
let Some(value) = value else {
return Ok(None);
};
if value.trim().is_empty() {
return Ok(None);
}
validate_external_generation_payload_json(field, value).map(Some)
}
fn contains_external_generation_inline_media_reference(value: &serde_json::Value) -> bool {
match value {
serde_json::Value::String(value) => is_external_generation_inline_media_reference(value),
serde_json::Value::Array(values) => values
.iter()
.any(contains_external_generation_inline_media_reference),
serde_json::Value::Object(values) => values.iter().any(|(key, value)| {
is_external_generation_inline_media_reference(key)
|| contains_external_generation_inline_media_reference(value)
}),
serde_json::Value::Null | serde_json::Value::Bool(_) | serde_json::Value::Number(_) => {
false
}
}
}
fn is_external_generation_inline_media_reference(value: &str) -> bool {
value
.trim_start()
.as_bytes()
.get(..5)
.is_some_and(|prefix| {
prefix.eq_ignore_ascii_case(b"data:") || prefix.eq_ignore_ascii_case(b"blob:")
})
}
struct ExternalGenerationPayloadCompactionOutcome {
compacted_json: Option<String>,
before_bytes: u64,
after_bytes: u64,
inline_media_count: u64,
invalid_json: bool,
}
fn compact_external_generation_payload_json(
payload_json: &str,
) -> ExternalGenerationPayloadCompactionOutcome {
let before_bytes = payload_json.len() as u64;
let Ok(mut payload) = serde_json::from_str::<serde_json::Value>(payload_json) else {
return ExternalGenerationPayloadCompactionOutcome {
compacted_json: None,
before_bytes,
after_bytes: before_bytes,
inline_media_count: 0,
invalid_json: true,
};
};
let inline_media_count = compact_external_generation_inline_media_references(&mut payload);
if inline_media_count == 0 {
return ExternalGenerationPayloadCompactionOutcome {
compacted_json: None,
before_bytes,
after_bytes: before_bytes,
inline_media_count: 0,
invalid_json: false,
};
}
let Ok(compacted_json) = serde_json::to_string(&payload) else {
return ExternalGenerationPayloadCompactionOutcome {
compacted_json: None,
before_bytes,
after_bytes: before_bytes,
inline_media_count: 0,
invalid_json: true,
};
};
ExternalGenerationPayloadCompactionOutcome {
after_bytes: compacted_json.len() as u64,
compacted_json: Some(compacted_json),
before_bytes,
inline_media_count,
invalid_json: false,
}
}
fn compact_external_generation_inline_media_references(value: &mut serde_json::Value) -> u64 {
match value {
serde_json::Value::String(text) => {
if is_external_generation_inline_media_reference(text) {
*text = INLINE_MEDIA_REMOVED_PLACEHOLDER.to_string();
1
} else {
0
}
}
serde_json::Value::Array(values) => values.iter_mut().fold(0u64, |count, value| {
count.saturating_add(compact_external_generation_inline_media_references(value))
}),
serde_json::Value::Object(values) => {
let source = std::mem::take(values);
let mut compacted = serde_json::Map::new();
let mut count = 0u64;
for (index, (key, mut nested)) in source.into_iter().enumerate() {
let mut next_key = key;
if is_external_generation_inline_media_reference(&next_key) {
count = count.saturating_add(1);
next_key = format!("__inline_media_key_removed_{index}");
while compacted.contains_key(&next_key) {
next_key.push('_');
}
}
count = count.saturating_add(compact_external_generation_inline_media_references(
&mut nested,
));
compacted.insert(next_key, nested);
}
*values = compacted;
count
}
serde_json::Value::Null | serde_json::Value::Bool(_) | serde_json::Value::Number(_) => 0,
}
}
fn duration_between_micros(later: i64, earlier: i64, field: &str) -> Result<i64, String> {
let duration = later.saturating_sub(earlier);
if duration <= 0 {
return Err(format!("{field} 必须大于 0"));
}
Ok(duration)
}
fn timestamp_after_micros(timestamp: Timestamp, duration_micros: i64) -> Timestamp {
Timestamp::from_micros_since_unix_epoch(
timestamp
.to_micros_since_unix_epoch()
.saturating_add(duration_micros.max(0)),
)
}
fn build_external_generation_lease_token(
job_id: &str,
worker_id: &str,
attempt: u32,
claimed_at: Timestamp,
) -> String {
format!(
"{}:{}:{}:{}",
job_id.trim(),
worker_id.trim(),
attempt,
claimed_at.to_micros_since_unix_epoch()
)
}
fn build_external_generation_event_id(
row: &ExternalGenerationJob,
event_kind: &str,
created_at: Timestamp,
) -> String {
format!(
"{}:{}:{}:{}:{}",
row.job_id.trim(),
event_kind.trim(),
row.status.trim(),
row.attempt,
created_at.to_micros_since_unix_epoch()
)
}
fn normalize_optional_text(value: &str) -> Option<String> {
let normalized = value.trim();
if normalized.is_empty() {
None
} else {
Some(normalized.to_string())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn external_generation_phase_only_accepts_known_execution_phases() {
assert_eq!(
normalize_external_generation_job_phase(" processing ").as_deref(),
Ok(EXTERNAL_GENERATION_PHASE_PROCESSING)
);
assert!(normalize_external_generation_job_phase("uploading").is_err());
}
#[test]
fn external_generation_phase_rejection_kind_is_machine_readable() {
let lease = ExternalGenerationJobPhaseUpdateError::lease_fencing("lease 已过期");
assert_eq!(
lease.kind,
ExternalGenerationJobPhaseUpdateFailureKind::LeaseFencingRejected
);
assert_eq!(lease.message, "lease 已过期");
let other = ExternalGenerationJobPhaseUpdateError::other("phase 非法");
assert_eq!(
other.kind,
ExternalGenerationJobPhaseUpdateFailureKind::OtherRejected
);
assert_eq!(other.message, "phase 非法");
}
#[test]
fn external_generation_phase_lease_guard_classifies_every_fencing_rejection() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.worker_id = Some("worker-a".to_string());
row.lease_token = Some("lease-1".to_string());
row.lease_expires_at = Some(micros(2_000));
assert!(
validate_external_generation_job_phase_update_lease(
&row,
"worker-a",
"lease-1",
micros(1_999),
)
.is_ok()
);
let mut terminal = row.clone();
terminal.status = EXTERNAL_GENERATION_STATUS_COMPLETED.to_string();
let cases = [
validate_external_generation_job_phase_update_lease(
&terminal,
"worker-a",
"lease-1",
micros(1_999),
),
validate_external_generation_job_phase_update_lease(
&row,
"worker-b",
"lease-1",
micros(1_999),
),
validate_external_generation_job_phase_update_lease(
&row,
"worker-a",
"lease-2",
micros(1_999),
),
validate_external_generation_job_phase_update_lease(
&row,
"worker-a",
"lease-1",
micros(2_000),
),
];
for result in cases {
let error = result.expect_err("stale worker 必须被 fencing 拒绝");
assert_eq!(
error.kind,
ExternalGenerationJobPhaseUpdateFailureKind::LeaseFencingRejected
);
}
}
#[test]
fn external_generation_job_result_failure_is_structured() {
let result = failed_external_generation_job_result("失败".to_string());
assert!(!result.ok);
assert_eq!(result.error_message.as_deref(), Some("失败"));
assert!(result.jobs.is_empty());
}
#[test]
fn worker_failure_refund_is_fenced_and_committed_with_the_job_transition() {
let source = include_str!("external_generation.rs");
let body = source
.split_once("fn fail_external_generation_job_tx(")
.and_then(|(_, tail)| {
tail.split_once("fn get_external_generation_queue_stats_tx(")
.map(|(body, _)| body)
})
.expect("worker failure transaction");
let lease_guard = body
.find("get_worker_owned_external_generation_job(")
.expect("lease guard");
let refund = body
.find("settle_external_generation_attempt_refund(")
.expect("wallet settlement");
let persist = body
.find("persist_external_generation_job_row(")
.expect("job transition");
assert!(lease_guard < refund);
assert!(refund < persist);
assert!(body.contains("row.refund_ledger_id = refund_ledger_id"));
}
#[test]
fn pending_job_is_claimable_only_after_available_time() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_PENDING);
row.request_label = "图片画布生成图片".to_string();
row.available_at = micros(1_000);
assert!(!is_external_generation_job_claimable(&row, micros(999)));
assert!(is_external_generation_job_claimable(&row, micros(1_000)));
}
#[test]
fn worker_claims_only_editor_jobs_and_leaves_legacy_rows_untouched() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_PENDING);
row.available_at = micros(1_000);
assert!(!is_external_generation_job_worker_claimable(
&row,
micros(1_000)
));
row.source_module = EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE.to_string();
assert!(is_external_generation_job_worker_claimable(
&row,
micros(1_000)
));
}
#[test]
fn running_job_is_claimable_only_after_lease_expires() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.lease_expires_at = Some(micros(2_000));
assert!(!is_external_generation_job_claimable(&row, micros(1_999)));
assert!(is_external_generation_job_claimable(&row, micros(2_000)));
}
#[test]
fn expired_final_attempt_is_finalized_instead_of_reclaimed() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.attempt = 1;
row.max_attempts = 1;
row.worker_id = Some("worker-a".to_string());
row.lease_expires_at = Some(micros(2_000));
row.lease_token = Some("lease-a".to_string());
assert!(is_external_generation_job_claimable(&row, micros(2_000)));
assert!(external_generation_job_has_exhausted_attempts(&row));
let finalized = mark_external_generation_job_lease_exhausted(
row.clone(),
micros(2_000),
Some(
"asset_operation_refund:external_generation_job:extgen-test:attempt:1".to_string(),
),
);
assert_eq!(finalized.status, EXTERNAL_GENERATION_STATUS_FAILED);
assert_eq!(finalized.attempt, 1);
assert_eq!(finalized.completed_at, Some(micros(2_000)));
assert!(finalized.worker_id.is_none());
assert!(finalized.lease_expires_at.is_none());
assert!(finalized.lease_token.is_none());
assert!(finalized.refund_ledger_id.is_some());
assert_eq!(
finalized.last_error_message.as_deref(),
Some(EXTERNAL_GENERATION_FINAL_ATTEMPT_LEASE_EXPIRED_MESSAGE)
);
row.max_attempts = 2;
assert!(!external_generation_job_has_exhausted_attempts(&row));
}
#[test]
fn terminal_job_is_not_claimable() {
for status in [
EXTERNAL_GENERATION_STATUS_COMPLETED,
EXTERNAL_GENERATION_STATUS_FAILED,
EXTERNAL_GENERATION_STATUS_CANCELLED,
] {
let row = external_generation_job_fixture(status);
assert!(!is_external_generation_job_claimable(&row, micros(10_000)));
}
}
#[test]
fn terminal_helper_includes_completed_failed_and_cancelled() {
for status in [
EXTERNAL_GENERATION_STATUS_COMPLETED,
EXTERNAL_GENERATION_STATUS_FAILED,
EXTERNAL_GENERATION_STATUS_CANCELLED,
] {
let row = external_generation_job_fixture(status);
assert!(is_external_generation_job_terminal(&row));
}
let pending = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_PENDING);
assert!(!is_external_generation_job_terminal(&pending));
}
#[test]
fn task_list_sort_bucket_prioritizes_active_then_unacknowledged_terminal() {
let running = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
let pending = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_PENDING);
let completed_unack = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
let mut completed_ack =
external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
completed_ack.notification_acknowledged_at = Some(micros(2_000));
assert!(
external_generation_job_sort_bucket(&running)
< external_generation_job_sort_bucket(&pending)
);
assert!(
external_generation_job_sort_bucket(&pending)
< external_generation_job_sort_bucket(&completed_unack)
);
assert!(
external_generation_job_sort_bucket(&completed_unack)
< external_generation_job_sort_bucket(&completed_ack)
);
}
#[test]
fn task_list_sort_time_uses_completion_time_for_terminal_jobs() {
let mut completed = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
completed.completed_at = Some(micros(2_000));
completed.updated_at = micros(9_000);
let mut older_completed =
external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
older_completed.completed_at = Some(micros(1_000));
older_completed.updated_at = micros(10_000);
assert!(external_generation_job_sort_time_micros(&completed) > 1_000);
assert!(
external_generation_job_sort_time_micros(&completed)
> external_generation_job_sort_time_micros(&older_completed)
);
}
#[test]
fn task_list_status_filter_accepts_public_queued_alias() {
assert_eq!(
normalize_external_generation_job_status_filter(&[
"queued".to_string(),
"running".to_string(),
"unknown".to_string(),
]),
vec![
EXTERNAL_GENERATION_STATUS_PENDING,
EXTERNAL_GENERATION_STATUS_RUNNING,
],
);
}
#[test]
fn worker_ownership_requires_matching_trimmed_worker_id() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.worker_id = Some("worker-a".to_string());
assert!(is_external_generation_job_owned_by_worker(
&row,
" worker-a "
));
assert!(!is_external_generation_job_owned_by_worker(
&row, "worker-b"
));
}
#[test]
fn worker_ownership_requires_matching_trimmed_lease_token() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.lease_token = Some("job-1:worker-a:1:1000".to_string());
assert!(is_external_generation_job_owned_by_lease_token(
&row,
" job-1:worker-a:1:1000 "
));
assert!(!is_external_generation_job_owned_by_lease_token(
&row,
"job-1:worker-a:2:2000"
));
}
#[test]
fn worker_lease_is_active_only_before_expiry() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.lease_expires_at = Some(micros(2_000));
assert!(is_external_generation_job_lease_active(&row, micros(1_999)));
assert!(!is_external_generation_job_lease_active(
&row,
micros(2_000)
));
}
#[test]
fn lease_token_changes_with_claim_attempt() {
let first =
build_external_generation_lease_token("extgen-test", "worker-a", 1, micros(1_000));
let second =
build_external_generation_lease_token("extgen-test", "worker-a", 2, micros(2_000));
assert_ne!(first, second);
}
#[test]
fn claimable_age_keeps_oldest_available_job() {
let mut stats = ExternalGenerationQueueStatsSnapshot {
pending_count: 0,
delayed_pending_count: 0,
claimable_pending_count: 0,
running_active_count: 0,
expired_running_count: 0,
terminal_count: 0,
claimable_count: 0,
oldest_claimable_age_micros: None,
now_micros: 10_000,
};
let mut old_job = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_PENDING);
old_job.available_at = micros(1_000);
let mut newer_job = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
newer_job.available_at = micros(8_000);
record_external_generation_claimable_age(&mut stats, &newer_job, 10_000);
record_external_generation_claimable_age(&mut stats, &old_job, 10_000);
assert_eq!(stats.oldest_claimable_age_micros, Some(9_000));
}
#[test]
fn positive_duration_between_client_times_is_preserved() {
assert_eq!(
duration_between_micros(3_500, 1_000, "external_generation_job.lease_duration"),
Ok(2_500),
);
assert!(duration_between_micros(1_000, 1_000, "duration").is_err());
}
#[test]
fn persisted_payload_validation_rejects_inline_media_recursively() {
let error = validate_external_generation_payload_json(
"external_generation_job.request_payload_json",
r#"{"prompt":"保留 data 和 blob 普通文本","nested":[{"url":" \nDaTa:image/png;base64,AAAA"}]}"#,
)
.expect_err("嵌套 Data URL 不得进入持久任务");
assert!(error.contains("禁止包含 data: 或 blob:"));
assert!(
validate_external_generation_payload_json(
"external_generation_job.request_payload_json",
r#"{"sourceImageObjectKey":"users/user-1/source.png","resourceId":"resource-1"}"#,
)
.is_ok()
);
}
#[test]
fn persisted_payload_validation_rejects_invalid_or_oversized_json() {
assert!(
validate_external_generation_payload_json(
"external_generation_job.request_payload_json",
"not-json",
)
.expect_err("持久任务参数必须是合法 JSON")
.contains("不是合法 JSON")
);
let oversized = format!(
r#"{{"prompt":"{}"}}"#,
"x".repeat(MAX_EXTERNAL_GENERATION_PAYLOAD_BYTES)
);
assert!(
validate_external_generation_payload_json(
"external_generation_job.result_payload_json",
&oversized,
)
.expect_err("request 和 result 使用同一持久化上限")
.contains("超过持久化上限")
);
}
#[test]
fn persisted_payload_guard_does_not_break_non_editor_transient_references() {
let puzzle_payload = r#"{"reference_image_src":"data:image/png;base64,AAAA"}"#;
assert_eq!(
validate_external_generation_persisted_payload_for_source(
"puzzle",
"external_generation_job.request_payload_json",
puzzle_payload,
),
Ok(puzzle_payload.to_string())
);
assert!(
validate_external_generation_persisted_payload_for_source(
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
"external_generation_job.request_payload_json",
puzzle_payload,
)
.is_err()
);
}
#[test]
fn job_summary_contract_never_serializes_request_or_result_payload() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.request_payload_json =
r#"{"prompt":"一只橙色陶罐猫","source":"data:image/png;base64,AAAA"}"#.to_string();
row.result_payload_json = Some(r#"{"image":"data:image/png;base64,BBBB"}"#.to_string());
let summary_row = build_external_generation_job_summary_row(&row, None);
let summary = map_external_generation_job_summary_row(summary_row);
let serialized =
serde_json::to_value(spacetimedb::sats::ser::serde::SerializeWrapper(&summary))
.expect("summary 应可序列化");
assert_eq!(summary.request_prompt.as_deref(), Some("一只橙色陶罐猫"));
assert!(serialized.get("request_payload_json").is_none());
assert!(serialized.get("result_payload_json").is_none());
for internal_field in [
"dedupe_key",
"attempt",
"max_attempts",
"worker_id",
"lease_expires_at_micros",
"available_at_micros",
] {
assert!(
serialized.get(internal_field).is_none(),
"summary 不应复制 worker 内部字段 {internal_field}"
);
}
}
#[test]
fn scene_job_summary_uses_user_scene_content_instead_of_internal_prompt() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_PENDING);
row.source_module = EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE.to_string();
row.job_kind = "editor_image_generation".to_string();
row.request_label = "图片画布生成图片".to_string();
row.request_payload_json = serde_json::json!({
"kind": "scene",
"prompt": "【内部场景规则】仅生成环境背景,禁止 UI 和文字",
"generationInputs": {
"fields": [
{ "title": "画面内容", "value": "雨夜中的海边车站" },
{ "title": "视觉风格", "value": "日系动画" }
],
"references": []
}
})
.to_string();
let summary = build_external_generation_job_summary_row(&row, None);
let refreshed = build_external_generation_job_summary_row(
&row,
Some(Some("【内部场景规则】旧摘要缓存".to_string())),
);
assert_eq!(summary.request_prompt.as_deref(), Some("雨夜中的海边车站"));
assert_eq!(summary.request_label, "图片画布生成游戏场景");
assert_eq!(
refreshed.request_prompt.as_deref(),
Some("雨夜中的海边车站")
);
assert_eq!(refreshed.request_label, "图片画布生成游戏场景");
row.request_payload_json = serde_json::json!({
"kind": "scene",
"prompt": "【内部场景规则】缺少用户画面内容时也不得展示",
"generationInputs": { "fields": [], "references": [] }
})
.to_string();
let missing_scene_content = build_external_generation_job_summary_row(&row, None);
assert!(missing_scene_content.request_prompt.is_none());
assert_eq!(missing_scene_content.request_label, "图片画布生成游戏场景");
row.source_module = "puzzle".to_string();
row.job_kind = "puzzle_compile_draft".to_string();
row.request_label = "拼图场景生成".to_string();
let non_editor = build_external_generation_job_summary_row(&row, None);
assert_eq!(non_editor.request_label, "拼图场景生成");
assert_eq!(
non_editor.request_prompt.as_deref(),
Some("【内部场景规则】缺少用户画面内容时也不得展示")
);
}
#[test]
fn job_summary_never_copies_inline_media_as_request_prompt() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.request_payload_json = r#"{"prompt":"data:image/png;base64,AAAA"}"#.to_string();
let extracted = build_external_generation_job_summary_row(&row, None);
let cached = build_external_generation_job_summary_row(
&row,
Some(Some(" blob:https://example.test/id".to_string())),
);
assert!(extracted.request_prompt.is_none());
assert!(cached.request_prompt.is_none());
}
#[test]
fn job_summary_bounds_and_redacts_error_messages() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_FAILED);
row.last_error_message =
Some(r#"provider response: {"image":"data:image/png;base64,AAAA"}"#.to_string());
let redacted = build_external_generation_job_summary_row(&row, None);
assert_eq!(
redacted.last_error_message.as_deref(),
Some(INLINE_MEDIA_ERROR_REDACTED_MESSAGE)
);
row.last_error_message = Some("x".repeat(MAX_EXTERNAL_GENERATION_ERROR_MESSAGE_CHARS + 10));
let bounded = build_external_generation_job_summary_row(&row, None)
.last_error_message
.expect("非空错误应保留有界摘要");
assert_eq!(
bounded.chars().count(),
MAX_EXTERNAL_GENERATION_ERROR_MESSAGE_CHARS + 1
);
assert!(bounded.ends_with('…'));
assert_eq!(
normalize_external_generation_error_message("metadata: unavailable").as_deref(),
Some("metadata: unavailable")
);
}
#[test]
fn job_summary_extracts_lightweight_completion_warning() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.result_payload_json = Some(
r#"{"sourceModule":"editor-canvas","warning":{"code":"insufficient-connected-components","reason":"连通域数量不足"}}"#
.to_string(),
);
let summary = build_external_generation_job_summary_row(&row, None);
assert_eq!(summary.warning_message.as_deref(), Some("连通域数量不足"));
assert!(summary.last_error_message.is_none());
}
#[test]
fn job_summary_bounds_and_redacts_completion_warnings() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.result_payload_json = Some(
serde_json::json!({
"warning": {
"code": "slice-persistence-failed",
"reason": "provider returned data:image/png;base64,AAAA"
}
})
.to_string(),
);
let redacted = build_external_generation_job_summary_row(&row, None);
assert_eq!(
redacted.warning_message.as_deref(),
Some(INLINE_MEDIA_WARNING_REDACTED_MESSAGE)
);
row.result_payload_json = Some(
serde_json::json!({
"warning": {
"code": "slice-persistence-failed",
"reason": "x".repeat(MAX_EXTERNAL_GENERATION_WARNING_MESSAGE_CHARS + 10)
}
})
.to_string(),
);
let bounded = build_external_generation_job_summary_row(&row, None)
.warning_message
.expect("非空告警应保留有界摘要");
assert_eq!(
bounded.chars().count(),
MAX_EXTERNAL_GENERATION_WARNING_MESSAGE_CHARS + 1
);
assert!(bounded.ends_with('…'));
row.result_payload_json = Some(r#"{"warning":{"code":"missing-reason"}}"#.to_string());
assert!(
build_external_generation_job_summary_row(&row, None)
.warning_message
.is_none()
);
}
#[test]
fn job_summary_top_n_keeps_memory_bounded_and_ordered() {
let mut rows = Vec::new();
for index in 1..=5 {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.job_id = format!("extgen-{index}");
row.completed_at = Some(micros(index));
row.updated_at = micros(index);
retain_external_generation_job_summary_top_n(
&mut rows,
build_external_generation_job_summary_row(&row, None),
2,
);
assert!(rows.len() <= 2);
}
assert_eq!(
rows.into_iter().map(|row| row.job_id).collect::<Vec<_>>(),
vec!["extgen-5".to_string(), "extgen-4".to_string()]
);
}
#[test]
fn active_jobs_are_never_payload_compaction_candidates() {
for status in [
EXTERNAL_GENERATION_STATUS_PENDING,
EXTERNAL_GENERATION_STATUS_RUNNING,
] {
let mut row = external_generation_job_fixture(status);
row.source_module = EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE.to_string();
assert!(!should_compact_external_generation_job_payloads(&row, None));
}
let mut completed = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
completed.source_module = EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE.to_string();
completed.completed_at = Some(micros(2_000));
assert!(should_compact_external_generation_job_payloads(
&completed,
Some(2_000)
));
assert!(!should_compact_external_generation_job_payloads(
&completed,
Some(1_999)
));
let non_editor = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
assert!(!should_compact_external_generation_job_payloads(
&non_editor,
None
));
}
#[test]
fn retention_only_selects_acknowledged_terminal_rows_before_cutoff() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.source_module = EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE.to_string();
row.completed_at = Some(micros(1_000));
row.updated_at = micros(1_000);
let mut summary = build_external_generation_job_summary_row(&row, None);
summary.notification_acknowledged_at = Some(micros(2_000));
assert!(is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
summary.notification_acknowledged_at = None;
assert!(!is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
summary.notification_acknowledged_at = Some(micros(2_000));
row.status = EXTERNAL_GENERATION_STATUS_RUNNING.to_string();
summary.status = EXTERNAL_GENERATION_STATUS_RUNNING.to_string();
assert!(!is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
row.status = EXTERNAL_GENERATION_STATUS_COMPLETED.to_string();
summary.status = EXTERNAL_GENERATION_STATUS_COMPLETED.to_string();
row.completed_at = Some(micros(1_001));
assert!(!is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
row.completed_at = Some(micros(1_000));
row.source_module = "puzzle".to_string();
assert!(!is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
}
#[test]
fn retention_rejects_mismatched_summary_identity_or_status() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_FAILED);
row.source_module = EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE.to_string();
row.completed_at = Some(micros(1_000));
let mut summary = build_external_generation_job_summary_row(&row, None);
summary.notification_acknowledged_at = Some(micros(2_000));
summary.job_id = "different-job".to_string();
assert!(!is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
summary.job_id = row.job_id.clone();
summary.status = EXTERNAL_GENERATION_STATUS_CANCELLED.to_string();
assert!(!is_external_generation_job_retention_candidate(
&row,
&summary,
EXTERNAL_GENERATION_EDITOR_SOURCE_MODULE,
1_000,
));
}
#[test]
fn maintenance_selector_bounds_scanned_rows_and_advances_by_last_scanned_job() {
let rows = (1..=4).map(|index| {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_COMPLETED);
row.job_id = format!("extgen-{index}");
row
});
let (selected, next_cursor, has_more, scanned_count) =
select_external_generation_job_ids_for_maintenance(rows, 2, |row| {
row.job_id != "extgen-1"
});
assert_eq!(selected, vec!["extgen-2".to_string()]);
assert_eq!(next_cursor.as_deref(), Some("extgen-2"));
assert!(has_more);
assert_eq!(scanned_count, 2);
}
#[test]
fn editor_generation_job_identity_binds_result_project() {
let mut row = external_generation_job_fixture(EXTERNAL_GENERATION_STATUS_RUNNING);
row.job_kind = "editor_image_generation".to_string();
row.source_module = "editor-canvas".to_string();
row.source_entity_id = "project-1".to_string();
row.request_payload_json = r#"{"prompt":"test"}"#.to_string();
let fingerprint = shared_contracts::editor_generation_request_fingerprint(
row.job_kind.as_str(),
row.request_payload_json.as_str(),
);
assert!(
validate_editor_generation_job_business_identity(
&row,
"user-1",
"editor_image_generation",
fingerprint.as_str(),
Some("project-1"),
)
.is_ok()
);
assert!(
validate_editor_generation_job_business_identity(
&row,
"user-1",
"editor_image_generation",
fingerprint.as_str(),
Some("project-2"),
)
.expect_err("job source entity must match the written project")
.contains("source_entity_id")
);
}
#[test]
fn payload_compaction_replaces_nested_inline_media_and_preserves_prompt() {
let payload = r#"{"prompt":"保留这段展示提示","nested":[{"url":"data:image/png;base64,AAAA"},{"deep":{"source":" BLOB:https://example.test/id"}}]}"#;
let outcome = compact_external_generation_payload_json(payload);
let compacted = outcome
.compacted_json
.as_deref()
.expect("含内联媒体的 JSON 应生成压缩结果");
let value: serde_json::Value = serde_json::from_str(compacted).expect("压缩结果仍是 JSON");
assert_eq!(outcome.inline_media_count, 2);
assert!(outcome.after_bytes < outcome.before_bytes);
assert_eq!(value["prompt"], "保留这段展示提示");
assert_eq!(value["nested"][0]["url"], INLINE_MEDIA_REMOVED_PLACEHOLDER);
assert_eq!(
value["nested"][1]["deep"]["source"],
INLINE_MEDIA_REMOVED_PLACEHOLDER
);
}
fn external_generation_job_fixture(status: &str) -> ExternalGenerationJob {
ExternalGenerationJob {
job_id: "extgen-test".to_string(),
dedupe_key: "puzzle:compile:test".to_string(),
job_kind: "puzzle_compile_draft".to_string(),
owner_user_id: "user-1".to_string(),
source_module: "puzzle".to_string(),
source_entity_id: "session-1".to_string(),
request_label: "拼图首关草稿生成".to_string(),
request_payload_json: r#"{"sessionId":"session-1"}"#.to_string(),
status: status.to_string(),
attempt: 0,
max_attempts: 1,
last_error_message: None,
worker_id: None,
lease_expires_at: None,
available_at: micros(0),
result_payload_json: None,
created_at: micros(0),
started_at: None,
completed_at: None,
updated_at: micros(0),
lease_token: None,
price_mud_points: 10,
refund_ledger_id: None,
notification_acknowledged_at: None,
phase: None,
}
}
#[test]
fn concurrency_limit_blocks_pending_at_limit_but_recycles_expired_running() {
// pending 在飞数达到上限后不再认领,任务保持 pending 排队。
assert!(external_generation_claim_within_concurrency_limit(
false, 0, 1
));
assert!(!external_generation_claim_within_concurrency_limit(
false, 1, 1
));
assert!(external_generation_claim_within_concurrency_limit(
false, 2, 3
));
assert!(!external_generation_claim_within_concurrency_limit(
false, 3, 3
));
// 回收账号自身过期的 running 不受上限限制,否则超限账号的卡死任务永远无法回收。
assert!(external_generation_claim_within_concurrency_limit(
true, 3, 1
));
// 128 哨兵表示不设上限。
assert!(external_generation_claim_within_concurrency_limit(
false,
10_000,
MEMBERSHIP_UNLIMITED_CONCURRENCY,
));
}
#[test]
fn owner_concurrency_cache_only_invalidates_running_count() {
let mut cache = ExternalGenerationOwnerConcurrencyCache::default();
cache.running.insert("owner-a".to_string(), 2);
cache.limits.insert("owner-a".to_string(), 3);
// 认领 / 终结后 running 计数必须重查,否则会把上一行的变更带到下一候选行。
cache.invalidate_running("owner-a");
assert!(cache.running.get("owner-a").is_none());
// 上限在一次事务内是常量,不受 running 计数失效影响。
assert_eq!(cache.limits.get("owner-a"), Some(&3));
}
fn micros(value: i64) -> Timestamp {
Timestamp::from_micros_since_unix_epoch(value)
}
}