扩展外部生成任务结果与稳定入队
支持调用方提供任务 ID 和去重键,复用既有外部任务队列。 保存并投影编辑器生成结果,供任务消费者恢复展示。
This commit is contained in:
@@ -23,6 +23,7 @@ export interface ExternalGenerationJobStatusRecord {
|
||||
progress: number;
|
||||
error?: string | null;
|
||||
updatedAtMicros: number;
|
||||
result?: unknown;
|
||||
}
|
||||
|
||||
export interface ExternalGenerationJobStatusResponse {
|
||||
|
||||
@@ -46,6 +46,35 @@ where
|
||||
T: Serialize,
|
||||
{
|
||||
let job_id = build_prefixed_uuid_id("task-");
|
||||
enqueue_editor_generation_job_with_identity(
|
||||
state,
|
||||
owner_user_id,
|
||||
job_kind,
|
||||
source_entity_id,
|
||||
request_label,
|
||||
price_mud_points,
|
||||
payload,
|
||||
job_id.clone(),
|
||||
format!("editor-canvas:{job_kind}:{job_id}"),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(crate) async fn enqueue_editor_generation_job_with_identity<T>(
|
||||
state: &AppState,
|
||||
owner_user_id: &str,
|
||||
job_kind: &str,
|
||||
source_entity_id: impl Into<String>,
|
||||
request_label: impl Into<String>,
|
||||
price_mud_points: u64,
|
||||
payload: &T,
|
||||
job_id: String,
|
||||
dedupe_key: String,
|
||||
) -> Result<ExternalGenerationJobRecord, AppError>
|
||||
where
|
||||
T: Serialize,
|
||||
{
|
||||
let request_payload_json = serde_json::to_string(payload).map_err(|error| {
|
||||
AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_details(json!({
|
||||
"provider": EDITOR_GENERATION_QUEUE_PROVIDER,
|
||||
@@ -56,7 +85,7 @@ where
|
||||
state
|
||||
.spacetime_client()
|
||||
.enqueue_external_generation_job(ExternalGenerationJobEnqueueRecordInput {
|
||||
dedupe_key: format!("editor-canvas:{job_kind}:{job_id}"),
|
||||
dedupe_key,
|
||||
job_id,
|
||||
job_kind: job_kind.to_string(),
|
||||
owner_user_id: owner_user_id.to_string(),
|
||||
@@ -89,6 +118,7 @@ pub(crate) fn editor_generation_queue_state(
|
||||
progress: 8,
|
||||
error: job.last_error_message,
|
||||
updated_at_micros: job.updated_at_micros,
|
||||
result: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -178,6 +178,11 @@ fn map_external_generation_job_status(
|
||||
progress,
|
||||
error: job.last_error_message.clone(),
|
||||
updated_at_micros: job.updated_at_micros,
|
||||
result: job
|
||||
.result_payload_json
|
||||
.as_deref()
|
||||
.and_then(|payload| serde_json::from_str::<Value>(payload).ok())
|
||||
.and_then(|payload| payload.get("response").cloned()),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -671,7 +671,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(error) => {
|
||||
let message = error.body_text();
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -699,7 +701,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(error) => {
|
||||
let message = error.body_text();
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -728,7 +732,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(error) => {
|
||||
let message = error.body_text();
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -756,7 +762,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(error) => {
|
||||
let message = error.body_text();
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -784,7 +792,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(error) => {
|
||||
let message = error.body_text();
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -813,7 +823,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(response) => {
|
||||
let message = response_error_message(response).await;
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -842,7 +854,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(response) => {
|
||||
let message = response_error_message(response).await;
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -871,7 +885,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(response) => {
|
||||
let message = response_error_message(response).await;
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -900,7 +916,9 @@ async fn process_external_generation_job_once(
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => complete_editor_generation_job(&state, &worker_id, &job).await,
|
||||
Ok(result) => {
|
||||
complete_editor_generation_job(&state, &worker_id, &job, result.0).await
|
||||
}
|
||||
Err(response) => {
|
||||
let message = response_error_message(response).await;
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
@@ -1005,6 +1023,7 @@ async fn complete_editor_generation_job(
|
||||
state: &AppState,
|
||||
worker_id: &str,
|
||||
job: &ExternalGenerationJobRecord,
|
||||
result: serde_json::Value,
|
||||
) -> Result<(), String> {
|
||||
complete_job(
|
||||
state,
|
||||
@@ -1014,6 +1033,7 @@ async fn complete_editor_generation_job(
|
||||
json!({
|
||||
"sourceModule": job.source_module.clone(),
|
||||
"sourceEntityId": job.source_entity_id.clone(),
|
||||
"response": result,
|
||||
})
|
||||
.to_string(),
|
||||
),
|
||||
|
||||
@@ -354,6 +354,7 @@ fn map_jump_hop_queue_job_status(
|
||||
progress: 8,
|
||||
error: job.last_error_message,
|
||||
updated_at_micros: job.updated_at_micros,
|
||||
result: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -262,6 +262,7 @@ fn map_puzzle_queue_job_status(
|
||||
.unwrap_or(100),
|
||||
error: job.last_error_message.clone(),
|
||||
updated_at_micros: job.updated_at_micros,
|
||||
result: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -331,6 +331,7 @@ fn map_puzzle_clear_queue_job_status(
|
||||
progress: 8,
|
||||
error: job.last_error_message,
|
||||
updated_at_micros: job.updated_at_micros,
|
||||
result: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -340,6 +340,7 @@ fn map_wooden_fish_queue_job_status(
|
||||
progress: 8,
|
||||
error: job.last_error_message,
|
||||
updated_at_micros: job.updated_at_micros,
|
||||
result: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -34,6 +34,8 @@ pub struct ExternalGenerationJobStatusRecord {
|
||||
pub progress: u8,
|
||||
pub error: Option<String>,
|
||||
pub updated_at_micros: i64,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub result: Option<serde_json::Value>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
|
||||
Reference in New Issue
Block a user