13 KiB
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-submit:
X-CSRF-Token与Cookie: csrf_token=<same-token>使用后端临时随机值,不写入配置文件或环境变量。 - 同一业务 SourceMessage 的
idempotency_key在所有尝试中保持不变。 - 初始 POST 成功后保存响应头
Content-Location,解析并保存 SuperAgentrun_id。 - SSE 必须按帧解析
event:、data:、id:和 heartbeat comment。 - 成功条件必须同时满足:最终 AI 内容、
run.completed status=success、顶层event: end,且没有顶层error或run.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,例如 OUTLOOK 或 EMAIL |
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 |
PENDING、RUNNING、SUCCEEDED、RETRYABLE_FAILED、FAILED、IGNORED |
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 的 AgentBus-like payload 语义一致。
发送 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 崩溃后,锁过期的
RUNNINGdispatch 可以被重新领取处理。 - Debug EML 不会进入生产 dispatch run。
FAILEDSourceMessage 不会触发 dispatch。- worker 能领取
PENDING、RETRYABLE_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-Location、run_id、SSEid、Last-Event-ID、GET /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,运行测试,中文提交。