249 lines
13 KiB
Markdown
249 lines
13 KiB
Markdown
# 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 目标是补齐这条链路:
|
||
|
||
```text
|
||
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`,解析并保存 SuperAgent `run_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 |
|
||
|
||
约束建议:
|
||
|
||
```text
|
||
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 建议包含:
|
||
|
||
```json
|
||
{
|
||
"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. 状态流转
|
||
|
||
```text
|
||
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 能领取 `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`、SSE `id`、`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. 后续目标模式建议
|
||
|
||
```text
|
||
进入目标模式,目标:实现 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,运行测试,中文提交。
|
||
```
|