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_startRun 开始run_id, session_id, model
turn_start新轮次开始turn_num
turn_end轮次结束turn_num
doneRun 正常完成reason, run_id
errorRun 错误error, kind, severity
interruptedRun 被中断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']}")