5.4 KiB
5.4 KiB
Run 流式架构改造说明
1. 改造目标
本次改造将对话输出从 HTTP 直连流 升级为 异步任务执行 + SSE 增量拉取,目标是:
- 页面离开/刷新不影响后台执行。
- 前端支持断线重连与续流。
- 提升系统稳定性、并发能力与可观测性。
- 控制存储成本:过程数据短期保存,结果数据长期保存。
2. 改造前后对比(仅与 HTTP 直连流对比)
| 维度 | 改造前(HTTP 直连流) | 改造后(Run + SSE) |
|---|---|---|
| 触发方式 | POST /agent/{id} 后长连接直接流式输出 |
POST /runs 创建任务,worker 异步执行 |
| 任务生命周期 | 绑定前端连接 | 与前端连接解耦 |
| 页面离开/刷新 | 常导致任务中断或前端丢上下文 | 任务继续执行,前端可续流 |
| 前端消费方式 | 同一个请求内读取 chunk | GET /runs/{id}/events?after_seq=... 增量拉取 |
| 恢复能力 | 弱,重连后难恢复 | 强,依赖 seq 游标恢复 |
| 取消语义 | 中断连接即可能影响任务 | 仅显式 cancel 才取消任务 |
3. 架构方案
3.1 组件职责
- FastAPI:负责创建 run、查询 run、SSE 输出、cancel 接口。
- ARQ Worker:负责真正执行模型流式任务。
- Redis:
- ARQ 队列
- run 过程事件流(Redis Stream)
- 取消信号(key + pub/sub)
- Postgres:
- run 执行状态(
agent_runs) - 最终业务消息(
messages/tool_calls) - checkpointer(会话运行状态)
- run 执行状态(
3.2 架构图
flowchart LR
FE["Frontend"] -->|"POST /runs"| API["FastAPI"]
API -->|"create run"| PG[("Postgres")]
API -->|"enqueue"| R[("Redis")]
W["ARQ Worker"] -->|"dequeue"| R
W -->|"update run status"| PG
W -->|"write stream events"| R
FE -->|"GET /runs/:id/events?after_seq=..."| API
API -->|"read incremental events"| R
API -->|"SSE events"| FE
FE -->|"POST /runs/:id/cancel"| API
API -->|"cancel mark"| PG
API -->|"publish cancel"| R
W -->|"persist messages/tool_calls"| PG
4. 端到端流程(事件流转)
sequenceDiagram
participant FE as Frontend
participant API as FastAPI
participant R as Redis
participant W as ARQ Worker
participant PG as Postgres
FE->>API: POST /api/chat/agent/{agent_id}/runs
API->>PG: create agent_runs(status=pending)
API->>R: enqueue process_agent_run(run_id)
API-->>FE: run_id
FE->>API: GET /api/chat/runs/{run_id}/events?after_seq=0
W->>R: dequeue run job
W->>PG: mark running
W->>R: append loading/tool/state events (stream)
API->>R: read events after_seq
API-->>FE: SSE incremental events
FE->>API: POST /api/chat/runs/{run_id}/cancel (optional)
API->>PG: mark cancel_requested
API->>R: publish cancel signal
W->>W: cancel current task
W->>PG: mark terminal status
W->>PG: persist messages/tool_calls
API-->>FE: close event
5. 接口与协议变更
5.1 对外路径(保持稳定)
POST /api/chat/agent/{agent_id}/runsGET /api/chat/runs/{run_id}GET /api/chat/runs/{run_id}/events?after_seq=...POST /api/chat/runs/{run_id}/cancel
5.2 after_seq 语义
- 主格式为字符串游标(Redis Stream ID,如
1700000000000-3)。 - 兼容旧整数参数输入。
- SSE 返回
seq字段统一为字符串,前端按单调递增去重。
5.3 SSE 事件格式
{
"run_id": "...",
"seq": "1700000000000-3",
"event_type": "loading",
"payload": {"items": [...]},
"ts": 1700000000000
}
控制事件:heartbeat / error / close。
6. 前端行为变化
- 本地记录活跃 run 快照:
active_run:{threadId}。 - 刷新/切回页面时按
run_id + last_seq自动续流。 - 接收事件先做 seq 去重,再更新 UI。
- 保留打字机效果(
requestAnimationFrame + throttle)。 - 保留首条消息自动更新会话标题逻辑。
7. 稳定性设计
- 幂等:
request_id避免重复创建 run。 - 重试:仅可恢复错误触发 ARQ 重试(
max_tries=2)。 - 取消:DB 状态 + Redis 信号双通道。
- SSE 生命周期:心跳、超时、终态关闭、断线重连。
- 状态单一真相:执行态在
agent_runs,业务态在messages/tool_calls + checkpointer。
8. 本次代码变更范围(未提交部分)
后端
/Yuxi-Know/backend/package/yuxi/services/run_queue_service.py/Yuxi-Know/backend/package/yuxi/services/run_worker.py/Yuxi-Know/backend/package/yuxi/services/agent_run_service.py/Yuxi-Know/backend/package/yuxi/repositories/agent_run_repository.py/Yuxi-Know/server/routers/chat_router.py/Yuxi-Know/server/worker_main.py/Yuxi-Know/backend/package/yuxi/storage/postgres/manager.py/Yuxi-Know/backend/package/yuxi/storage/postgres/models_business.py
前端
/Yuxi-Know/web/src/apis/agent_api.js/Yuxi-Know/web/src/components/AgentChatComponent.vue
测试
/Yuxi-Know/test/test_run_queue_service.py/Yuxi-Know/test/test_agent_run_service.py/Yuxi-Know/test/test_run_worker.py
配置
/Yuxi-Know/docker-compose.yml/Yuxi-Know/docker-compose.prod.yml/Yuxi-Know/.env.template
9. 验收标准
- 发送消息后,SSE 能持续收到增量事件。
- 页面刷新后,可按
after_seq恢复输出。 - 页面离开不影响后台执行。
- 取消后 run 状态正确收敛,输出停止。
- 最终消息与工具调用正常入库。