add retry in result reconciliation logic

This commit is contained in:
2026-07-16 17:41:06 +08:00
parent ca4bf6db5e
commit ce2806dea8
4 changed files with 142 additions and 46 deletions
@@ -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`
@@ -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` 的业务真相。
### 创作 / 游玩统一流程主干
@@ -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` 保存提示词、比例、清晰度、模型、时长等可展示参数的稳定名称、用户可见标题和值;
@@ -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 mut last_reconcile_error = None;
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: job_id.clone(),
owner_user_id: conversation.owner_user_id.clone(),
})
.await
{
Ok(job) => job,
Err(_) => continue,
Err(error) => {
last_reconcile_error = Some(format!("读取生成任务结果失败:{error}"));
if should_retry_result_reconcile(retry_count) {
sleep(EDITOR_AGENT_RESULT_RECONCILE_RETRY_DELAY).await;
}
continue;
}
};
let changed = match job.status.as_str() {
"completed" => reconcile_completed_editor_agent_tool_call(
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(),
),
"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,
};
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);
) {
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<T: serde::de::DeserializeOwned>(value: &Value) -> Result<T, String> {
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 次后仍失败"));
}
}