Files
th-hotel-simple/docs/project/requirements/M007-agentbus-superagent-auto-dispatch-v1.md

13 KiB
Raw Blame History

M007 AgentBus SourceMessage 自动分发 SuperAgent 链路 V1

文档信息

项目 内容
文档版本 1.0
日期 2026-07-12
状态 后端 V1 已实现,等待测试机联调
适用范围 AgentBus 实时邮件入库后,异步调用 SuperAgent Open API 的生产链路
主要读者 产品、后端、测试、运维、SuperAgent 对接方、后续协作 agent

1. 背景和目标

M001 已完成 AgentBus 邮件来源事实入库M004 已完成 Debug EML 人工上传并调用 SuperAgent 的调试链路。当前缺口是:真实 AgentBus WebSocket 收到邮件并写入 SourceMessage Inbox 后,还没有生产链路自动把该邮件推送给 SuperAgent。

本 checkpoint 目标是补齐这条链路:

AgentBus WebSocket 收到邮件
→ 写入 SourceMessage Inbox
→ 创建 SuperAgent dispatch / outbox 记录
→ 异步 worker 调用 SuperAgent Open API
→ 保存 session_id / run_id / raw_answer / parsed_json / 状态
→ 等待 SuperAgent 通过现有任务结果通知接口或 MCP 提交业务结果

中文说明AgentBus 仍然是消息入口SuperAgent 仍然是外部 AI / Agent 能力提供方。自动分发只负责把已入库 SourceMessage 交给 SuperAgent不直接创建订单、任务、OPERA 操作或客户回复。

2. 设计原则

  • AgentBus WebSocket 回调内只做快速入库和轻量投递标记,不同步等待 SuperAgent。
  • SourceMessage 捕获事务和 SuperAgent 外部调用事务分离,避免外部慢请求影响邮件入库。
  • 使用 outbox / dispatch run 表记录状态、重试、幂等键和 SuperAgent run 信息。
  • Debug EML 链路和生产自动分发链路复用 SuperAgent Open API client但运行记录表和 provider 必须区分。
  • SuperAgent 返回内容只能保存为外部能力输出,不直接改变业务最终状态。
  • 业务任务创建继续依赖 SuperAgent 后续调用本系统任务结果通知接口或 MCP 工具。

3. 触发规则

第一版只对以下 SourceMessage 创建自动分发记录:

条件 规则
来源 provider 必须是 AGENTBUS
捕获状态 必须是 RECEIVED
幂等结果 RECEIVED 状态会尝试幂等创建 dispatch数据库唯一键保证不会重复创建
重复投递 会再次尝试 insertIfAbsent 补偿可能缺失的 dispatch如 payload changed只记录 SourceMessage 重复诊断
Debug EML 不参与生产自动分发
FAILED Inbox 不分发,保留入库失败原因供排查

中文说明:后续如果需要对历史 SourceMessage 补发 SuperAgent应单独设计人工 replay / backfill 接口,不在本 checkpoint 内隐式补发。

4. SuperAgent Open API 稳定性要求

2026-07-12 导入的 OPEN_AGENT_API_JAVA_SSE_CLIENT.md 对 Java SSE 调用提出新的强约束。本项目后续 Open API client 必须满足:

  • 请求 messages/stream?include_trace=true 时生成稳定 X-Request-ID
  • Open API 请求必须携带 CSRF double-submitX-CSRF-TokenCookie: csrf_token=<same-token> 使用后端临时随机值,不写入配置文件或环境变量。
  • 同一业务 SourceMessage 的 idempotency_key 在所有尝试中保持不变。
  • 初始 POST 成功后保存响应头 Content-Location,解析并保存 SuperAgent run_id
  • SSE 必须按帧解析 event:data:id: 和 heartbeat comment。
  • 成功条件必须同时满足:最终 AI 内容、run.completed status=success、顶层 event: end,且没有顶层 errorrun.failed
  • EOF、Premature EOF、incomplete chunked response 不能当成功。
  • 断流后如果已有 run_id,不得重新 POST 初始消息;应先查询 GET /runs/{run_id},再通过 GET /runs/{run_id}/events 携带 Last-Event-ID 恢复。
  • 恢复失败应记录为可诊断失败,不返回部分回答。

中文说明:这条规则同时适用于 Debug EML 和 AgentBus 自动分发。实现时应优先改造共享 SuperAgent Open API client避免两条链路行为不一致。

5. 数据模型建议

新增表建议命名为 platform_superagent_dispatch_run,作为 SourceMessage 到 SuperAgent Open API 的生产分发运行记录。

核心字段建议:

字段 中文说明
id 内部主键 ID
hotel_id SourceMessage 所属酒店 ID
source_message_id 本系统内部 SourceMessage Inbox ID
source_provider 来源 provider第一版主要为 AGENTBUS
source_channel 来源 channel例如 OUTLOOKEMAIL
external_message_id AgentBus 邮件外部消息 ID
external_conversation_id AgentBus 邮件会话 ID
dispatch_source 分发来源代码,第一版为 AGENTBUS_REALTIME
idempotency_key 发送给 SuperAgent 的业务幂等键,同一 SourceMessage 固定
request_id 本系统生成的调用关联 ID同时写入 X-Request-ID
dispatch_status PENDINGRUNNINGSUCCEEDEDRETRYABLE_FAILEDFAILEDIGNORED
attempt_count 已尝试次数
max_attempts 最大尝试次数
next_attempt_at 下次可重试 UTC 时间
locked_until worker 抢占锁过期 UTC 时间
superagent_session_id SuperAgent Open API session ID
superagent_run_id SuperAgent run ID
superagent_run_uri SuperAgent run 查询路径或 URI
last_event_id 最后完整处理的 SSE event id
superagent_raw_answer SuperAgent 最终原始回答
superagent_parsed_json 能解析为 JSON 时保存的结构化结果
superagent_trace_json 可安全保存的公开 trace 摘要
safe_error_code 安全错误代码
safe_error_summary 安全错误摘要不包含正文、HTML、附件 URL 或 Secret
created_at / updated_at 记录创建和更新时间UTC

约束建议:

UNIQUE KEY uk_dispatch_source_message (source_message_id, dispatch_source)
KEY idx_dispatch_status_next_attempt (dispatch_status, next_attempt_at)
KEY idx_dispatch_external_message (hotel_id, external_message_id)

中文说明:source_message_id + dispatch_source 是本系统内的分发幂等边界。外部 idempotency_key 是给 SuperAgent 的幂等边界,两者都需要稳定。

6. 配置建议

第一版新增配置建议:

配置 默认值 中文说明
agentbus.superagent-dispatch.enabled false 是否在 AgentBus 新邮件入库后创建自动分发记录
agentbus.superagent-dispatch.worker-enabled false 是否启动异步 worker 处理 dispatch
agentbus.superagent-dispatch.max-attempts 3 单条 dispatch 最大尝试次数
agentbus.superagent-dispatch.batch-size 10 worker 每轮领取数量
agentbus.superagent-dispatch.lock-ttl 5m worker 处理锁有效期
agentbus.superagent-dispatch.initial-backoff 30s 首次失败后的重试等待时间
agentbus.superagent-dispatch.max-backoff 15m 最大重试等待时间
superagent.open-api.agentbus-external-subject-id th-hotel-agentbus-source-message AgentBus 自动分发创建 session 时使用的 external subject id
superagent.open-api.sse-recovery-max-attempts 5 SSE 断流恢复最大次数

中文说明:生产首次上线建议先开启 AgentBus 入库,确认稳定后再单独开启自动分发和 worker。Debug EML 使用的访问 key 和生产自动分发无关。

7. 消息组装

AgentBus 自动分发发送给 SuperAgent 的 message 第一版应基于 SourceMessage 原始 payload 构造。Debug EML 发送给 SuperAgent 的 agentbus_like_payload 也按 AgentBus Outlook-like 主结构组装,并包含 reply_policy.modereply_policy.final_onlyDebug EML V1 暂用 mode=manualfinal_only=true,不能再使用旧的 Debug 专属 debug_only。Debug 来源区分仍保留在 SourceMessage provider=DEBUG_EML_UPLOAD、payload 表 schemaVersion=debug-eml-upload-v1 和 Open API metadata 中,不混入主 payload。后续 AgentBus 明确真实 reply_policy.mode 枚举后Debug EML 和实时链路需要一起对齐。

发送 metadata 建议包含:

{
  "source": "th-hotel-agentbus-realtime",
  "dispatch_run_id": "内部 dispatch run ID",
  "source_message_id": "内部 SourceMessage Inbox ID",
  "external_message_id": "AgentBus external_message_id",
  "external_conversation_id": "AgentBus external_conversation_id",
  "hotel_id": "系统酒店 ID"
}

注意metadata 里的内部 SourceMessage ID 只用于本系统排查SuperAgent 后续任务结果通知中的 source_message_id 仍应使用 AgentBus source.external_message_id

8. 状态流转

PENDING / RETRYABLE_FAILED / 过期 RUNNING
→ RUNNING
→ SUCCEEDED

RUNNING
→ RETRYABLE_FAILED
→ RUNNING

RUNNING / RETRYABLE_FAILED
→ FAILED

任意未启用或不应处理场景
→ IGNORED

失败分类建议:

当前错误代码 中文说明 是否可重试
SUPERAGENT_DISPATCH_FAILED SuperAgent Open API 调用失败,具体 HTTP / SSE / run 失败原因写入安全摘要 是,直到超过 max attempts
SOURCE_MESSAGE_PAYLOAD_NOT_FOUND SourceMessage 原始 payload 缺失 是,直到超过 max attempts通常需要人工排查数据
SUPERAGENT_DISPATCH_UNEXPECTED dispatch worker 内部非预期错误 是,直到超过 max attempts

中文说明:如果已经获得 superagent_run_id,后续恢复不得重新发送初始 POST避免同一业务消息被执行两次。

9. 非目标范围

本 checkpoint 不做:

  • 不创建 Reservation 订单。
  • 不创建 Reservation 任务。
  • 不调用 OPERA / OHIP。
  • 不自动 ACK AgentBus。
  • 不自动发送 task.result
  • 不自动回复客户。
  • 不做历史 SourceMessage 批量补发。
  • 不做前端页面。
  • 不改变 SuperAgent 任务结果通知接口的业务处理规则。

10. 验收标准

  • AgentBus 新邮件入库成功后,在开启配置时幂等创建一条 dispatch run。
  • 重复 AgentBus 邮件投递不会创建第二条 dispatch run。
  • 如首次入库后 dispatch 创建失败,后续重复 RECEIVED 投递可以补偿创建缺失的 dispatch run。
  • worker 崩溃后,锁过期的 RUNNING dispatch 可以被重新领取处理。
  • Debug EML 不会进入生产 dispatch run。
  • FAILED SourceMessage 不会触发 dispatch。
  • worker 能领取 PENDINGRETRYABLE_FAILED 和锁已过期的 RUNNING 记录并调用 SuperAgent Open API。
  • SuperAgent 调用成功时保存 session、run、raw answer、parsed json、trace 摘要和 SUCCEEDED 状态。
  • SuperAgent SSE 缺少 end、缺少最终内容或 run.completed 时不能成功。
  • SSE EOF 后如果已有 run_id,使用 /events 恢复,不重发初始 POST。
  • 恢复耗尽后记录安全错误摘要和可重试状态。
  • 任务创建仍只由 SuperAgent 后续任务结果通知接口或 MCP 提交触发。
  • 测试覆盖 AgentBus 触发、幂等、失败、不处理 Debug EML、SSE 严格成功条件和断流恢复。

11. 已实现代码范围

本次后端 V1 已落地:

  • 新增 platform_superagent_dispatch_run 表和 V20__create_superagent_dispatch_run.sql
  • 新增 SuperAgentDispatchRunEntity、Mapper、Repository、Service 和配置类。
  • AgentBus 捕获 AGENTBUS + RECEIVED 后,在配置开启时幂等创建 dispatch run重复投递可补偿缺失 outbox。
  • worker 处理 PENDING / RETRYABLE_FAILED / 锁已过期 RUNNING,调用共享 SuperAgent Open API client。
  • Open API client 支持 Content-Locationrun_id、SSE idLast-Event-IDGET /runs/{run_id}GET /runs/{run_id}/events 恢复。
  • Debug EML 继续复用共享 Open API client。
  • dispatch 成功保存 session、run、last event id、raw answer、parsed json、trace 摘要和状态。

12. 后续目标模式建议

进入目标模式,目标:实现 M007 AgentBus SourceMessage 自动分发 SuperAgent 链路 V1。

范围:
1. 按 docs/project/requirements/M007-agentbus-superagent-auto-dispatch-v1.md 实现 AgentBus 入库后的异步 SuperAgent dispatch / outbox。
2. 新增 platform_superagent_dispatch_run 表、Entity、Mapper、Repository、Service、ServiceImpl 和必要枚举 / DTO。
3. AgentBus 捕获新 SourceMessage 成功后,在配置开启时创建 dispatch run重复投递不重复创建。
4. 新增 worker 处理 PENDING / RETRYABLE_FAILED dispatch调用 SuperAgent Open API。
5. 改造 SuperAgent Open API client满足 20260712 SSE 稳定性要求严格成功条件、Content-Location/run_id、SSE id、Last-Event-ID、/runs 查询和 /events 恢复EOF 后不重发初始 POST。
6. Debug EML 继续复用改造后的 SuperAgent Open API client行为不回退。
7. 保存 SuperAgent session_id、run_id、last_event_id、raw answer、parsed json、trace 摘要、状态和安全错误摘要。
8. 更新集成文档、上线注意事项和配置说明。
9. 补充测试覆盖触发、幂等、失败、Debug EML 不触发、SSE 严格成功条件和断流恢复。

不做:
1. 不创建订单。
2. 不创建任务。
3. 不调用 OPERA / OHIP。
4. 不自动 ACK AgentBus。
5. 不自动发送 task.result 或客户回复。
6. 不做历史 SourceMessage 批量补发。
7. 不做前端页面。

完成后:
code review运行测试中文提交。