108 lines
4.0 KiB
Markdown
108 lines
4.0 KiB
Markdown
# 工作流 SSE / 事件协议 v1.0(冻结)
|
||
|
||
> 实现源:`packages/harness/deerflow/workflows/events.py`
|
||
> 错误体:`packages/harness/deerflow/workflows/errors.py`
|
||
> 事件清单与开发方案对齐:`WORKFLOW_STUDIO_BACKEND_DEV_ZH.md` §9
|
||
|
||
## 不变式
|
||
|
||
1. **先持久化,后发布** — 用户可见事件先写入 `workflow_run_events`,提交成功后再 fan-out。
|
||
2. **`seq` 严格单调** — 唯一键 `(run_id, seq)`。
|
||
3. **可重放** — `Last-Event-ID` 优先于 `after_seq`,否则从 `0`。
|
||
4. **增量合并** — 模型 token 建议 30–100ms 或约 1KB 合并落库。
|
||
5. **终态停连** — `run.completed` / `run.failed` / `run.cancelled` 且库已补齐后关闭 SSE。
|
||
6. **不泄密** — 无凭据、无完整认证头、无 DB 连接串、无隐藏思维链。
|
||
|
||
## 端点(阶段 2 已实现)
|
||
|
||
```http
|
||
GET /api/workflows/runs/{run_id}/stream?after=128
|
||
Accept: text/event-stream
|
||
Last-Event-ID: 128
|
||
```
|
||
|
||
实现:`app/gateway/routers/workflow_runs.py`。
|
||
|
||
- **库是排序的权威来源**:连接先按 `Last-Event-ID`(优先)或 `after` 从 `workflow_run_events` 回放,再挂 `WorkflowLiveHub` 订阅实时帧;`seq` 不大于已回放游标的实时帧被丢弃,下一轮 DB 轮询(0.8s)会带上它。断线重连因此不会看到空洞或乱序。
|
||
- **终态必闭**:读到 `run.completed` / `run.failed` / `run.cancelled` 即关闭。若运行已是终态却没有终态事件(例如崩在写事件之前),补发一条 `synthesised: true` 的终态帧再关闭,客户端不会永久挂住。
|
||
- **响应头** `X-Workflow-Run-Status` 给出连接建立瞬间的状态。
|
||
- Heartbeat:SSE comment `: heartbeat`(**不写库**),15s 静默后发送。
|
||
- 同一份事件也可用 `GET /api/workflows/runs/{run_id}/events?after=&limit=` 分页拉取(无长连接场景)。
|
||
|
||
## Envelope
|
||
|
||
```json
|
||
{
|
||
"schemaVersion": "1.0",
|
||
"runId": "run_01",
|
||
"workflowId": "wf_01",
|
||
"versionId": "wv_03",
|
||
"seq": 42,
|
||
"event": "node.output.delta",
|
||
"nodeId": "agent_writer",
|
||
"nodeRunId": "nr_09",
|
||
"timestamp": "2026-08-26T10:00:00+08:00",
|
||
"data": {}
|
||
}
|
||
```
|
||
|
||
SSE 帧:
|
||
|
||
```text
|
||
id: 42
|
||
event: node.output.delta
|
||
data: { ... envelope ... }
|
||
```
|
||
|
||
Python:`WorkflowEventEnvelope.to_sse_frame()`。
|
||
|
||
## 事件类型(封闭集合)
|
||
|
||
| 事件 | `data` 要点 |
|
||
|------|-------------|
|
||
| `run.created` | input 摘要、version |
|
||
| `run.queued` | queue 信息 |
|
||
| `run.started` | startedAt |
|
||
| `node.queued` | node 信息 |
|
||
| `node.started` | attempt、输入摘要 |
|
||
| `node.progress` | phase、message、percent |
|
||
| `node.output.delta` | channel、delta |
|
||
| `node.tool.started` | toolCallId、name、input 摘要 |
|
||
| `node.tool.finished` | toolCallId、status、output 摘要 |
|
||
| `artifact.created` | artifact 元数据 |
|
||
| `node.completed` | output 摘要、duration |
|
||
| `node.failed` | error code/message/retryable |
|
||
| `run.awaiting_input` | formSchema、actions、resumeToken |
|
||
| `run.resumed` | action |
|
||
| `run.cancel_requested` | requestedBy |
|
||
| `run.cancelled` | finishedAt |
|
||
| `run.completed` | output、artifacts |
|
||
| `run.failed` | error、failedNodeId |
|
||
|
||
## 统一错误体
|
||
|
||
```json
|
||
{
|
||
"code": "WORKFLOW_SQL_POLICY_DENIED",
|
||
"message": "SQL 节点只允许单条只读查询",
|
||
"retryable": false,
|
||
"nodeId": "sql_1",
|
||
"details": { "rule": "single_read_statement" }
|
||
}
|
||
```
|
||
|
||
完整错误码列表:`ALL_WORKFLOW_ERROR_CODES`(`deerflow.workflows.errors`)。
|
||
|
||
## `resumeToken` 只走事件流
|
||
|
||
`run.awaiting_input` 的 `data.resumeToken` 是恢复运行的唯一凭证,**只在事件流里出现**——`GET /api/workflows/runs/{run_id}` 的运行视图不包含它。恢复时 `POST /api/workflows/runs/{run_id}/resume` 携带该 token,服务端以「消费即失效」的方式原子翻转状态(同一个 token 不能恢复两次)。
|
||
|
||
## 前后端契约测试
|
||
|
||
前后端共用 fixture:
|
||
|
||
- `workflows/samples/event_envelope_example.json`
|
||
- 本文件事件名表
|
||
|
||
实现测试:`tests/test_workflow_runs.py`(回放游标、终态关闭、暂停/恢复的事件序列)。
|