From ce2806dea8b461fd8b6cf96d4ae3680016c499aa Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Thu, 16 Jul 2026 17:41:06 +0800 Subject: [PATCH] add retry in result reconciliation logic --- .../shared-memory/decision-log.md | 2 +- ...】server-rs与SpacetimeDB数据契约-2026-05-15.md | 1 + .../【编辑器】画布Agent对话面板-2026-07-03.md | 2 +- .../api-server/src/editor_agent/reconcile.rs | 183 +++++++++++++----- 4 files changed, 142 insertions(+), 46 deletions(-) diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 6b2f8ad6b..952c00fe4 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -3165,7 +3165,7 @@ ## 2026-07-13 画布 Agent 工具执行状态复用外部生成任务 - 决策:画布 Agent 的 OSS 工具消息使用 `status=not_completed|completed|failed|cancelled` 和可选 `externalJobId`;不使用 `cancelledAt`,不新增关联表。`external_generation_job` 是排队、执行、lease 与计费结算的唯一真相;OSS status 只表达该消息回填结果,不复制 queued / running。确认接口按 `conversationId + messageId + toolName` 稳定去重并复用既有编辑器 worker job kind。 -- 懒回填:`GET /conversation` 会在持有 conversation lock 后扫描 `status=not_completed` 且有 `externalJobId` 的工具消息,按 job id 定向读取主任务;完成时复用原工具 formatter 更新 system text、写入轻量媒体引用并标记 `completed`,失败时写入 `error` 并标记 `failed`。排队 / 执行保持 `not_completed`。前端轮询发现终态后调用该 GET,取得正常工具卡结果并刷新画布 / 任务列表。worker 的 `result_payload_json` 只保留 formatter 与媒体引用所需的轻量生成回包,不向通用 summary 状态接口投影。 +- 懒回填:`GET /conversation` 会在持有 conversation lock 后扫描 `status=not_completed` 且有 `externalJobId` 的工具消息,按 job id 定向读取主任务;完成时复用原工具 formatter 更新 system text、写入轻量媒体引用并标记 `completed`,任务本身失败时写入 `error` 并标记 `failed`。任务结果读取或 completed payload 解析 / formatter 首次失败后,在同一次 GET 内最多重试 3 次,每次等待 100ms 并重新读取主任务;仍失败才把工具消息写为 `failed`,不依赖前端再次刷新。排队 / 执行保持 `not_completed`。worker 的 `result_payload_json` 只保留 formatter 与媒体引用所需的轻量生成回包,不向通用 summary 状态接口投影。 - 影响范围:画布 Agent 共享契约、确认/取消接口、编辑器生成入队 helper、对话状态展示与恢复。 - 验证方式:`cargo check -p api-server -p shared-contracts --manifest-path server-rs/Cargo.toml`、画布 Agent 定向前端测试、`npm run typecheck`、`npm run check:encoding`、`git diff --check`。 diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index ae292e591..7ed660513 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -85,6 +85,7 @@ npm run check:server-rs-ddd - 对话附件只允许引用当前工程 `editor_project_resource` 或当前账号 `editor_asset` 的图片;前端可提交展示用 `imageSrc` / `thumbnailSrc`,后端必须按 `resourceId` / `assetId` 重新归一、校验 owner / project 和 `objectKey`,再给 LLM 或生成工具使用。 - 画布 Agent 工具复用既有编辑器图片生成 / 修改 / 图标 spritesheet BFF,并继续使用后端模型定价和 `execute_billable_asset_operation_with_cost`;前端不提交 `priceMudPoints`。 - `/messages/{messageId}/confirm` 与 `/messages/{messageId}/cancel` 只返回成功确认;前端成功后立即重新读取整个会话,以会话详情中的权威消息状态和 `externalJobId` 驱动气泡展示与任务轮询。 +- 会话详情的终态懒回填必须在单次 GET 和同一 conversation lock 内完成有界重试:任务结果读取、completed payload 解析或工具 formatter 首次失败后最多重试 3 次,每次等待 100ms 并重新读取主任务;仍失败后把 OSS 工具消息原子写为 `failed`,保存“重试 3 次后仍失败”的最后错误。不得依赖浏览器再次轮询才能耗尽重试,也不得让消息永久停在 `not_completed`。 - 画布 Agent 是“正式任务 payload 不进入通用用户 read model”规则的窄例外消费者:`GET /api/editor/agent-conversations/{conversationId}` 只按会话中已有的 `externalJobId` 定向读取主任务,完成后由对应工具 formatter 从 `result_payload_json` 提取并归一有界的图片 / 视频 / 音频引用,写入 OSS 工具消息后返回。前端仍不得通过通用任务列表 / 状态接口读取或解析 `request_payload_json` / `result_payload_json`;OSS 轻量媒体引用只是会话展示与后续 Agent 上下文,不替代 `editor_project_resource`、`editor_asset`、`editor_canvas.layers_json` 或 `external_generation_job` 的业务真相。 ### 创作 / 游玩统一流程主干 diff --git a/docs/【编辑器】画布Agent对话面板-2026-07-03.md b/docs/【编辑器】画布Agent对话面板-2026-07-03.md index 1f5489165..aa00c6f9e 100644 --- a/docs/【编辑器】画布Agent对话面板-2026-07-03.md +++ b/docs/【编辑器】画布Agent对话面板-2026-07-03.md @@ -78,7 +78,7 @@ - 工具消息只保存 `status`(`not_completed` / `completed` / `failed` / `cancelled`)和可选 `externalJobId`。`external_generation_job` 仍是队列、执行、lease 与计费结算真相;OSS status 仅表示该条对话消息是否已经回填完成结果或失败,不复制 queued / running。 - 确认接口必须先把工具参数转换为既有编辑器 worker payload,再使用 `editor-agent:{conversationId}:{messageId}:{toolName}` 稳定 dedupe key 入队;同一确认的请求重试只能得到同一个 external job。入队成功后把返回的 job id 写回同一条 OSS 工具消息,不新增 Agent 工具执行关联表。 - 前端根据 `externalJobId` 查询通用 external-generation job 状态;worker 继续通过 `canvasCompletion` 把生成结果写回工程与素材库。浏览器断线、刷新或 api-server 重启不得导致确认接口重新扣费或重新提交 provider。 -- `GET /conversation` 会在同一个 conversation lock 内扫描 `status=not_completed` 且已有 `externalJobId` 的工具消息:只对这些消息按 job id 定向读取主任务;任务完成后复用对应工具的 `format_execute_message` 替换 system text、回填轻量图片 / 视频 / 音频引用并写为 `completed`,任务失败则回填 `error` 并写为 `failed`。排队和执行中都保持 `not_completed`,整轮扫描结果一次性写回 OSS。前端轮询发现任务终态后调用该 GET,使用回填后的正常工具卡展示结果,再刷新画布和任务列表。 +- `GET /conversation` 会在同一个 conversation lock 内扫描 `status=not_completed` 且已有 `externalJobId` 的工具消息:只对这些消息按 job id 定向读取主任务;任务完成后复用对应工具的 `format_execute_message` 替换 system text、回填轻量图片 / 视频 / 音频引用并写为 `completed`,任务失败则回填 `error` 并写为 `failed`。任务结果读取或 completed payload 解析 / formatter 回填失败时,必须在同一次 GET 内完成首次尝试及最多 3 次重试,三次重试各间隔 100ms 并重新读取任务结果;仍失败才把该工具消息写为 `failed` 并保存最后错误。该重试不依赖前端再次刷新。排队和执行中都保持 `not_completed`,整轮扫描结果一次性写回 OSS。 - `EditorAgentToolCall.args` 保留为工具返回的原始 JSON,是确认接口重新反序列化并执行工具的唯一参数真相。图片参数继续只保存由真实 data key 计算出的 opaque SHA-256 `imageId`;不得为了前端预览把 `args` 中的图片 ID 改写成 `objectKey`、URL 或展示对象,也不得由前端重组或回传一份新的执行参数。 - `EditorAgentToolCall.displayArgs` 是必填、只读的用户确认展示投影,与 `args` 分离: - `stringArgs` 保存提示词、比例、清晰度、模型、时长等可展示参数的稳定名称、用户可见标题和值; diff --git a/server-rs/crates/api-server/src/editor_agent/reconcile.rs b/server-rs/crates/api-server/src/editor_agent/reconcile.rs index e0447aa1f..d8de3a2fa 100644 --- a/server-rs/crates/api-server/src/editor_agent/reconcile.rs +++ b/server-rs/crates/api-server/src/editor_agent/reconcile.rs @@ -29,6 +29,10 @@ use shared_contracts::editor_agent::{ EditorAgentConversationMessagesDocument, EditorAgentMessage, EditorAgentToolCallStatus, }; use spacetime_client::{EditorAgentConversationRecord, ExternalGenerationJobGetRecordInput}; +use tokio::time::{Duration, sleep}; + +pub(crate) const EDITOR_AGENT_RESULT_RECONCILE_MAX_RETRIES: u8 = 3; +const EDITOR_AGENT_RESULT_RECONCILE_RETRY_DELAY: Duration = Duration::from_millis(100); pub async fn reconcile_editor_agent_tool_calls( state: &AppState, @@ -50,59 +54,98 @@ pub async fn reconcile_editor_agent_tool_calls( let mut reconciled = Vec::new(); for (message_index, job_id) in candidates { - let job = match state - .spacetime_client() - .get_external_generation_job_result(ExternalGenerationJobGetRecordInput { - job_id, - owner_user_id: conversation.owner_user_id.clone(), - }) - .await - { - Ok(job) => job, - Err(_) => continue, - }; + let mut last_reconcile_error = None; - let changed = match job.status.as_str() { - "completed" => reconcile_completed_editor_agent_tool_call( - &mut document.messages[message_index], - job.result_payload_json.as_deref(), - ), - "failed" => { - let tool_call = document.messages[message_index] - .tool_call - .as_mut() - .expect("candidate contains a tool call"); - let error = job - .last_error_message - .unwrap_or_else(|| "生成失败".to_string()); - tool_call.error = Some(error.clone()); - tool_call.status = EditorAgentToolCallStatus::Failed; - document.messages[message_index].text = - format!("[tool_call:{}] output: {error}", tool_call.tool_name); - Ok(()) - } - _ => continue, - }; + for retry_count in 0..=EDITOR_AGENT_RESULT_RECONCILE_MAX_RETRIES { + let job = match state + .spacetime_client() + .get_external_generation_job_result(ExternalGenerationJobGetRecordInput { + job_id: job_id.clone(), + owner_user_id: conversation.owner_user_id.clone(), + }) + .await + { + Ok(job) => job, + Err(error) => { + last_reconcile_error = Some(format!("读取生成任务结果失败:{error}")); + if should_retry_result_reconcile(retry_count) { + sleep(EDITOR_AGENT_RESULT_RECONCILE_RETRY_DELAY).await; + } + continue; + } + }; - match changed { - Ok(()) => reconciled.push(document.messages[message_index].clone()), - Err(error) => { - let tool_call = document.messages[message_index] - .tool_call - .as_mut() - .expect("candidate contains a tool call"); - tool_call.error = Some(error.clone()); - tool_call.status = EditorAgentToolCallStatus::Failed; - document.messages[message_index].text = - format!("[tool_call:{}] output: {error}", tool_call.tool_name); - reconciled.push(document.messages[message_index].clone()); + match job.status.as_str() { + "completed" => { + let original_message = document.messages[message_index].clone(); + match reconcile_completed_editor_agent_tool_call( + &mut document.messages[message_index], + job.result_payload_json.as_deref(), + ) { + Ok(()) => { + last_reconcile_error = None; + reconciled.push(document.messages[message_index].clone()); + break; + } + Err(error) => { + document.messages[message_index] = original_message; + last_reconcile_error = Some(error); + } + } + } + "failed" => { + mark_job_failed( + &mut document.messages[message_index], + job.last_error_message + .unwrap_or_else(|| "生成失败".to_string()), + ); + last_reconcile_error = None; + reconciled.push(document.messages[message_index].clone()); + break; + } + _ => { + last_reconcile_error = None; + break; + } } + + if should_retry_result_reconcile(retry_count) { + sleep(EDITOR_AGENT_RESULT_RECONCILE_RETRY_DELAY).await; + } + } + + if let Some(error) = last_reconcile_error { + mark_result_reconcile_failed(&mut document.messages[message_index], error); + reconciled.push(document.messages[message_index].clone()); } } Ok(reconciled) } +fn should_retry_result_reconcile(retry_count: u8) -> bool { + retry_count < EDITOR_AGENT_RESULT_RECONCILE_MAX_RETRIES +} + +fn mark_job_failed(message: &mut EditorAgentMessage, error: String) { + let tool_call = message + .tool_call + .as_mut() + .expect("reconcile candidate contains a tool call"); + tool_call.error = Some(error.clone()); + tool_call.status = EditorAgentToolCallStatus::Failed; + message.text = format!("[tool_call:{}] output: {error}", tool_call.tool_name); +} + +fn mark_result_reconcile_failed(message: &mut EditorAgentMessage, error: String) { + mark_job_failed( + message, + format!( + "任务终态结果回填重试 {EDITOR_AGENT_RESULT_RECONCILE_MAX_RETRIES} 次后仍失败:{error}" + ), + ); +} + fn reconcile_completed_editor_agent_tool_call( message: &mut EditorAgentMessage, result_payload_json: Option<&str>, @@ -189,3 +232,55 @@ fn reconcile_completed_editor_agent_tool_call( fn parse_reconciled_value(value: &Value) -> Result { serde_json::from_value(value.clone()).map_err(|error| error.to_string()) } + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + fn pending_tool_message() -> EditorAgentMessage { + serde_json::from_value(json!({ + "id": 1, + "role": "system", + "text": "waiting", + "attachments": [], + "toolCall": { + "toolName": "generate_image", + "status": "not_completed", + "args": {}, + "displayArgs": { + "stringArgs": [], + "imageArgs": [], + "extras": { "priceMudPoints": 0 } + }, + "externalJobId": "job-1", + "images": [] + }, + "createdAt": "2026-07-16T00:00:00Z" + })) + .expect("pending tool message should deserialize") + } + + #[test] + fn terminal_result_reconcile_retries_three_times_after_the_initial_attempt() { + assert!(should_retry_result_reconcile(0)); + assert!(should_retry_result_reconcile(1)); + assert!(should_retry_result_reconcile(2)); + assert!(!should_retry_result_reconcile(3)); + } + + #[test] + fn exhausted_terminal_result_reconcile_marks_the_tool_call_failed() { + let mut message = pending_tool_message(); + + mark_result_reconcile_failed(&mut message, "invalid result payload".to_string()); + + let tool_call = message.tool_call.expect("tool call should remain present"); + assert_eq!(tool_call.status, EditorAgentToolCallStatus::Failed); + assert_eq!( + tool_call.error.as_deref(), + Some("任务终态结果回填重试 3 次后仍失败:invalid result payload") + ); + assert!(message.text.contains("重试 3 次后仍失败")); + } +}