SSE 流式协议
wesgine 的 Run 使用 Server-Sent Events(SSE)流式返回事件。
基本格式
每个事件是一行 data: + JSON:
data: {"type":"stream_delta","content":"你好"}
data: {"type":"done","reason":"end_turn","run_id":"run-abc123"}
事件之间以空行分隔。
事件类型一览
流式输出
| type | 说明 | 关键字段 |
|---|---|---|
stream_delta | 模型文本增量 | content |
thinking_delta | 推理内容增量(Thinking 模型) | content |
生命周期
| type | 说明 | 关键字段 |
|---|---|---|
run_start | Run 开始 | run_id, session_id, model |
turn_start | 新轮次开始 | turn_num |
turn_end | 轮次结束 | turn_num |
done | Run 正常完成 | reason, run_id |
error | Run 错误 | error, kind, severity |
interrupted | Run 被中断 | reason |
工具调用
| type | 说明 | 关键字段 |
|---|---|---|
tool_start | 工具调用开始 | tool_name, call_id, arguments |
tool_progress | 执行进度 | call_id, output(增量) |
tool_end | 工具调用完成 | call_id, result, is_error |
Plan
| type | 说明 | 关键字段 |
|---|---|---|
plan_created | 计划创建 | plan_id, steps[] |
plan_updated | 步骤更新 | task_id, status, result |
plan_completed | 计划完成 | total_steps, done_steps |
plan_interrupted | 计划中断 | reason |
HITL
| type | 说明 | 关键字段 |
|---|---|---|
approval_request | 需要用户审批 | request_id, tool_name, arguments |
终端事件
每个 Run 最终必定产出一个终端事件:
| 终端事件 | 说明 |
|---|---|
done | 正常结束 |
error | 错误终止 |
interrupted | 被中断 |
INV-TERM-04:缺少终端事件时,引擎自动合成
error(missing_terminal_event)。
done.reason 值
| reason | 含义 |
|---|---|
end_turn | 模型正常结束 |
context_cancelled | 用户中断 |
cycle_detected | 循环检测终止 |
budget_exhausted | 预算耗尽 |
error_streak | 连续错误终止 |
SSE 重连
断开后可通过 Replay 端点续播:
# 从第 42 个事件开始重放
curl -N http://localhost:9091/cells/my-cell/runs/run-abc123/events?from_seq=42 \
-H "Authorization: Bearer $TOKEN"
Run 生命周期与 SSE 完全解耦:SSE 断开不会取消 Run。Run 在后台继续执行,所有事件通过 RunEventBuffer 保存,重连后可续播。
客户端实现建议
JavaScript
// 注意:浏览器 EventSource 不支持 POST,使用 fetch + ReadableStream
const resp = await fetch('http://localhost:9091/cells/my-cell/run', {
method: 'POST',
headers: {
'Authorization': 'Bearer ' + token,
'Content-Type': 'application/json',
},
body: JSON.stringify({
actor: 'user-001',
model: 'gpt-4o',
messages: [{role: 'user', content: [{type: 'text', text: '你好'}]}],
}),
})
const reader = resp.body.getReader()
const decoder = new TextDecoder()
let buffer = ''
while (true) {
const { done, value } = await reader.read()
if (done) break
buffer += decoder.decode(value, { stream: true })
const lines = buffer.split('\n')
buffer = lines.pop() // 保留不完整的最后一行
for (const line of lines) {
if (!line.startsWith('data: ')) continue
const event = JSON.parse(line.slice(6))
switch (event.type) {
case 'stream_delta':
document.getElementById('output').textContent += event.content
break
case 'done':
console.log('\n--- Run 完成:', event.reason)
break
case 'error':
console.error('错误:', event.error)
break
}
}
}
Python
import requests
import json
resp = requests.post(
'http://localhost:9091/cells/my-cell/run',
headers={
'Authorization': f'Bearer {token}',
'Content-Type': 'application/json',
},
json={
'actor': 'user-001',
'model': 'gpt-4o',
'messages': [{'role': 'user', 'content': [{'type': 'text', 'text': '你好'}]}],
},
stream=True,
)
for line in resp.iter_lines():
if not line:
continue
line = line.decode('utf-8')
if not line.startswith('data: '):
continue
event = json.loads(line[6:])
if event['type'] == 'stream_delta':
print(event['content'], end='', flush=True)
elif event['type'] == 'done':
print(f"\n--- Run 完成: {event['reason']}")
elif event['type'] == 'error':
print(f"错误: {event['error']}")