Agent 工程体系 · 第 91/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。
Agent 流式事件协议:Delta、Tool、Usage、Finish、断线和恢复
流式 Agent API 不是“把模型输出按 token 打印出来”这么简单。
一次 Agent 运行可能包含多次模型调用、工具调用、工具返回、Agent handoff、人工审批、错误、重试、历史压缩和最终持久化。用户看到的文字只是其中一种事件。若协议只定义 text,前端无法可靠区分“模型正在生成”“工具正在执行”“等待人工确认”和“运行已经完成”。
本文建立一套面向生产系统的 Agent 流式事件协议,重点解决六个问题:
Delta如何表达增量文本、结构化参数和其他可拼接内容;Tool如何表达调用、参数、执行、结果和错误;Usage如何在多轮、多工具、多模型调用中统计;Finish如何区分当前响应结束、当前 Agent 结束和整个 Run 完成;- 断线时如何判断客户端究竟丢了什么;
- 恢复时如何避免重复输出、重复执行工具和重复扣费。
OpenAI Agents SDK 当前把流式事件分成三类:直接透传模型的原始事件、表示完整运行项的高层事件,以及表示当前 Agent 变化的事件。SDK 文档同时明确指出:最后一个可见 token 到达后,运行仍可能继续进行会话持久化、审批状态处理或历史压缩,只有流迭代器结束后,运行才真正完成。(openai.github.io)
因此,生产协议必须把“可见输出”与“运行生命周期”分开建模。
一、先建立正确的对象模型:消息、事件、Run 和响应
1.1 Delta 不是消息
Delta 是一次流式传输中的增量片段。它通常不能独立解释,必须与同一个逻辑输出项的前序 Delta 按顺序拼接。
例如完整文本:
杭州今天多云,最高气温 28°C。
可能被拆成:
"杭州"
"今天"
"多云,"
"最高气温"
" 28°C。"
这些片段不是五条消息,而是同一条消息的五个增量。
定义一个输出项 item,其 Delta 序列为:
若文本拼接运算为 ⊕,则最终文本为:
这个公式成立的前提是:
- 所有 Delta 属于同一个
item_id; - Delta 按协议序号递增;
- 每个 Delta 只发送尚未发送的内容;
- 客户端不会把不同输出项的 Delta 混合拼接。
因此,客户端不应使用“收到事件就追加到当前文本框”的无状态逻辑,而应维护:
run_id
item_id
next_seq
assembled_content
1.2 Event 不是消息,也不是日志
事件是对运行状态变化的外部通知。一个事件至少需要回答:
- 它属于哪次运行?
- 它属于哪个逻辑输出项?
- 它在整个运行中的位置是什么?
- 它能否被重复消费?
- 客户端可以根据它恢复到什么状态?
建议采用统一事件信封:
{
"protocol": "agent-events.v1",
"event_id": "evt_01J...",
"run_id": "run_01J...",
"seq": 42,
"type": "message.delta",
"occurred_at": "2026-09-01T10:00:01.245Z",
"payload": {
"item_id": "item_msg_01",
"delta": "杭州"
}
}
其中:
| 字段 | 含义 |
|---|---|
protocol |
协议版本,不等于模型版本 |
event_id |
事件唯一标识,用于去重和审计 |
run_id |
一次 Agent 运行的唯一标识 |
seq |
同一 run_id 内单调递增的事件序号 |
type |
事件类型 |
occurred_at |
服务端产生事件的时间 |
payload |
类型相关的数据 |
seq 和 event_id 解决的是两个不同问题:
seq判断顺序、缺口和恢复位置;event_id判断重复投递。
不能只使用 event_id。随机事件 ID 能去重,但不能判断中间是否丢失了事件。也不能只使用时间戳,因为多个事件可能具有相同时间戳,跨机器时还可能出现时钟偏差。
1.3 Run 是比响应更大的生命周期
一次 Agent Run 可以包含多个模型响应:
用户输入
↓
模型调用 1:决定调用天气工具
↓
工具执行
↓
模型调用 2:根据工具结果生成答案
↓
最终输出
如果使用 Agents SDK,Runner.run_streamed() 返回的是 RunResultStreaming,其 stream_events() 会产出原始模型事件、运行项事件和 Agent 更新事件;高层运行项事件可以表示消息、工具调用、工具输出和 handoff 等完整项。(openai.github.io)
所以协议中至少要区分:
run_id 一次完整 Agent 运行
response_id 一次模型响应
item_id 一个消息、工具调用或工具结果
event seq 运行内事件位置
二、事件协议的核心分层
一个可用的协议通常分为四层:
flowchart TD
A[Agent Run] --> B[生命周期事件]
A --> C[模型响应事件]
A --> D[运行项事件]
A --> E[运营事件]
B --> B1[run.started]
B --> B2[run.paused]
B --> B3[run.completed]
B --> B4[run.failed]
B --> B5[run.cancelled]
C --> C1[response.created]
C --> C2[message.delta]
C --> C3[tool.arguments.delta]
C --> C4[response.finished]
D --> D1[tool.started]
D --> D2[tool.finished]
D --> D3[handoff.started]
D --> D4[approval.required]
E --> E1[usage.updated]
E --> E2[heartbeat]
E --> E3[checkpoint.created]
这四层不能混为一谈:
- 生命周期事件回答“Run 现在处于什么阶段”;
- 模型响应事件回答“模型输出了什么增量”;
- 运行项事件回答“一个消息或工具项是否已经形成”;
- 运营事件回答“消耗、心跳、检查点和诊断信息是什么”。
OpenAI Responses API 的流式事件同样同时包含响应生命周期事件、输出项事件和文本 Delta;官方示例中常见的事件包括 response.created、response.output_text.delta、response.completed 和 error。(developers.openai.com)
三、Delta:增量必须可定位、可校验、可恢复
3.1 文本 Delta
建议的文本 Delta:
{
"protocol": "agent-events.v1",
"event_id": "evt_0042",
"run_id": "run_1001",
"seq": 42,
"type": "message.delta",
"payload": {
"response_id": "resp_02",
"item_id": "msg_01",
"index": 7,
"delta": "多云,"
}
}
其中 index 是该 item_id 内的增量编号,seq 是整个 Run 的事件编号。
为什么需要两个编号?
假设事件顺序为:
seq=40, item=msg_01, index=0
seq=41, item=tool_01, index=0
seq=42, item=msg_01, index=1
全局 seq 能描述事件之间的关系;局部 index 能帮助客户端校验某个消息是否缺片段。
客户端处理逻辑:
def apply_message_delta(state: dict, event: dict) -> None:
payload = event["payload"]
item_id = payload["item_id"]
index = payload["index"]
delta = payload["delta"]
item = state.setdefault(item_id, {
"next_index": 0,
"text": "",
"status": "streaming",
})
if index < item["next_index"]:
# 重复事件:忽略
return
if index > item["next_index"]:
raise ValueError(
f"delta gap: item={item_id}, "
f"expected={item['next_index']}, got={index}"
)
item["text"] += delta
item["next_index"] += 1
这段代码中的关键不是字符串拼接,而是:
index < next_index表示重复投递;index == next_index表示可以应用;index > next_index表示出现缺口,不能继续盲目渲染。
3.2 Delta 不一定是完整 UTF-8 字符
传输层必须以完整事件或完整字符串为边界,而不能把任意字节块直接当作字符片段。
错误做法:
TCP chunk 1: E6 B5
TCP chunk 2: 8B E5 B7 B...
UTF-8 字符可能被拆在字节层,客户端若先把字节转换为字符串,再做业务级 Delta 处理,可能遇到解码错误或替换字符。
因此:
- SSE 的事件边界应由服务端生成;
- WebSocket 的消息边界应由协议实现保证;
- 网关不得任意截断或重新拼接 JSON 字节;
- 业务 Delta 应在 JSON 字符串层表达,不应暴露底层网络分片。
3.3 结构化 Delta
工具参数经常也是流式生成的:
{
"type": "tool.arguments.delta",
"payload": {
"tool_call_id": "call_01",
"index": 0,
"delta": "{\"city\":\"杭"
}
}
下一事件:
{
"type": "tool.arguments.delta",
"payload": {
"tool_call_id": "call_01",
"index": 1,
"delta": "州\",\"unit\":\"celsius\"}"
}
}
拼接后才是完整 JSON:
{"city":"杭州","unit":"celsius"}
客户端不能在每个 Delta 到达后都执行 JSON 解析。正确做法是:
- 按
tool_call_id建立参数缓冲区; - 按
index校验顺序; - 追加 Delta;
- 收到
tool.arguments.finished后再解析; - 解析失败则生成协议错误,而不是执行工具。
反例:
# 错误:每个 delta 都尝试执行
args = json.loads(event["payload"]["delta"])
run_tool(args)
这个实现会在第一个片段 {"city":"杭 到达时失败;如果为了“容错”而补括号执行,则可能把不完整或错误的参数提交给真实工具。
四、Tool:工具调用是一个状态机,不是一条消息
4.1 工具调用的完整生命周期
一个工具调用至少包含以下阶段:
planned
↓
arguments_streaming
↓
arguments_ready
↓
approval_required ── reject ──> rejected
↓ approve
executing
↓
succeeded / failed / timed_out / cancelled
建议事件序列:
tool.call.created
tool.arguments.delta
tool.arguments.finished
tool.approval.required # 可选
tool.started
tool.progress # 可选
tool.finished
tool.output # 可选
tool.failed # 可选
工具事件示例:
{
"protocol": "agent-events.v1",
"event_id": "evt_tool_001",
"run_id": "run_1001",
"seq": 51,
"type": "tool.call.created",
"payload": {
"response_id": "resp_02",
"item_id": "tool_item_01",
"tool_call_id": "call_01",
"name": "get_weather",
"arguments": null,
"status": "arguments_streaming"
}
}
参数完成:
{
"protocol": "agent-events.v1",
"event_id": "evt_tool_003",
"run_id": "run_1001",
"seq": 53,
"type": "tool.arguments.finished",
"payload": {
"tool_call_id": "call_01",
"arguments": {
"city": "杭州",
"unit": "celsius"
},
"arguments_sha256": "..."
}
}
执行完成:
{
"protocol": "agent-events.v1",
"event_id": "evt_tool_005",
"run_id": "run_1001",
"seq": 55,
"type": "tool.finished",
"payload": {
"tool_call_id": "call_01",
"status": "succeeded",
"output": {
"temperature": 28,
"condition": "cloudy"
},
"duration_ms": 83
}
}
4.2 tool.call.created 不代表工具已经执行
模型产生工具调用意图,与服务端真正执行工具之间存在边界。
以下两种事件不能等价:
tool.call.created 模型要求调用工具
tool.started 执行器已经接受并开始执行
如果工具需要人工确认,可能出现:
tool.call.created
tool.arguments.finished
approval.required
run.paused
这时不能发送 tool.started,更不能因为前端显示了“正在删除文件”就认为文件已经删除。
OpenAI Agents SDK 支持流式运行在工具审批处暂停:流事件迭代结束后,待审批项通过 interruptions 暴露;应用将结果转换为 RunState,批准或拒绝后再恢复运行。(openai.github.io)
4.3 工具结果必须携带幂等键
断线恢复最危险的情况是:
客户端断线
服务端已经扣款
客户端未收到 tool.finished
客户端重连
服务端重新执行扣款
因此每次工具执行都必须有稳定的执行键:
tool_execution_id = hash(run_id, tool_call_id, attempt)
但仅有客户端传入的幂等键还不够,工具执行器或下游服务必须真正保存结果:
CREATE TABLE tool_execution (
execution_id TEXT PRIMARY KEY,
run_id TEXT NOT NULL,
tool_call_id TEXT NOT NULL,
tool_name TEXT NOT NULL,
arguments_hash TEXT NOT NULL,
status TEXT NOT NULL,
output_json TEXT,
error_json TEXT,
created_at TIMESTAMP NOT NULL,
finished_at TIMESTAMP
);
执行流程:
def execute_once(execution_id, tool_name, arguments):
existing = store.get(execution_id)
if existing is not None:
if existing.arguments_hash != sha256_json(arguments):
raise RuntimeError("same execution id with different arguments")
return existing.output_or_error()
store.insert_running(
execution_id=execution_id,
tool_name=tool_name,
arguments_hash=sha256_json(arguments),
)
try:
output = TOOL_REGISTRY[tool_name](arguments)
store.mark_succeeded(execution_id, output)
return output
except Exception as exc:
store.mark_failed(execution_id, serialize_error(exc))
raise
这里的约束是:
更准确地说,是“具有副作用的逻辑执行最多一次”。如果底层网络请求在服务端已经发出但进程在写入 succeeded 前崩溃,数据库事务本身无法证明下游副作用是否发生。因此,对于支付、发货、删除等操作,还需要下游系统支持幂等键或查询确认。
五、Usage:用量是运行事实,不是 UI 进度条
5.1 Usage 的三种口径
Usage 至少有三个层次:
请求级
一次模型 API 请求的用量:
{
"request_id": "req_01",
"input_tokens": 1200,
"output_tokens": 80,
"reasoning_tokens": 0
}
响应级
一次模型响应可能对应一次请求,但在重试、流式重连或多供应商适配时,不应想当然地认为二者永远一一对应。
Run 级
整个 Agent Run 的累计用量:
其中每个 是一次模型调用的用量,包括:
- 初始回答;
- 生成工具调用;
- 工具返回后再次回答;
- handoff 后的新 Agent 调用;
- 运行结束前的历史压缩请求。
Agents SDK 会跟踪一次 Run 中的请求数、输入 token、输出 token、总 token,以及每请求明细;这些统计包含产生工具调用或 handoff 的模型调用。(openai.github.io)
5.2 为什么不能用 Delta 数量估算 token
假设输出 Delta 为:
"Hello"
" world"
它们是字符串片段,不等价于两个 token。token 数取决于模型的 tokenizer,中文、英文、空格、标点和代码的切分方式都可能不同。
错误估算:
estimated_output_tokens = number_of_deltas
这只能作为前端动画进度的粗略指标,不能用于:
- 计费;
- 限额;
- 成本分析;
- 模型对比;
- 上下文窗口保护。
5.3 Usage 事件应该如何发送
推荐两种模式。
模式 A:结束时发送最终 Usage
{
"type": "usage.final",
"payload": {
"requests": 2,
"input_tokens": 1830,
"output_tokens": 147,
"total_tokens": 1977,
"details": {
"reasoning_tokens": 0,
"cached_tokens": 900
}
}
}
适合普通聊天,协议简单,但运行期间无法实施动态预算。
模式 B:阶段性发送累计 Usage
{
"type": "usage.updated",
"payload": {
"scope": "run",
"requests": 2,
"input_tokens": 1830,
"output_tokens": 147,
"total_tokens": 1977,
"as_of_seq": 61
}
}
usage.updated 应表达累计值,而不是本次增量值。这样重复投递不会导致客户端重复累加。
若发送的是增量 Usage,则必须提供独立的 usage_seq:
{
"input_tokens_delta": 300,
"output_tokens_delta": 20,
"usage_seq": 4
}
客户端只能对尚未应用的 usage_seq 累加。
5.4 SDK 中的 Usage 与事件协议的关系
Agents SDK 的用量主要从运行结果和运行上下文访问,而不是简单依赖“最后一条文本事件”。运行完成后可以从 result.context_wrapper.usage 读取聚合值;还可以读取 request_usage_entries 获取每次请求的明细。(openai.github.io)
当使用第三方模型适配器时,用量是否准确取决于上游是否返回 Usage,以及适配器是否保留它。某些 Chat Completions 流式后端需要显式启用 usage chunk;而 preserve_raw_usage 只能保留已经到达适配器的原始 Usage,并不会主动向提供商请求 Usage。(openai.github.io)
因此,生产协议中应区分:
usage.status = reported
usage.status = unavailable
usage.status = estimated
不能把“没有收到 Usage”序列化成:
{"total_tokens": 0}
因为“真实为 0”和“没有上报”是两个不同事实。
六、Finish:至少要有三种结束
6.1 文本完成不等于 Run 完成
至少存在以下结束边界:
message.finished:一条消息不再产生新的文本 Delta;response.finished:一次模型响应结束;run.completed:整个 Agent Run 完成;run.paused:运行暂停,等待审批或外部输入;run.failed:运行失败;run.cancelled:运行被取消。
典型工具流程:
response.created
tool.call.created
tool.arguments.delta
tool.arguments.finished
tool.started
tool.finished
response.finished
response.created
message.delta
message.finished
response.finished
run.completed
此时第一个 response.finished 只表示模型完成了工具调用响应,不能告诉前端“最终答案已经完成”。
6.2 finish_reason 是原因,不是生命周期事件
finish_reason 通常表示某个模型响应为什么停止,例如:
stop
length
tool_calls
content_filter
error
它是响应级元数据:
{
"type": "response.finished",
"payload": {
"response_id": "resp_02",
"finish_reason": "tool_calls",
"usage": null
}
}
但 finish_reason = "stop" 也不一定意味着整个 Agent Run 完成,因为运行时可能还要保存 Session、执行后处理或更新审批状态。Agents SDK 文档明确说明,流迭代器结束前,final_output 可能仍然是 None;只有流处理结束后才可视为最终结果已经形成。(openai.github.io)
6.3 run.completed 必须满足终态条件
建议定义:
一个简单状态机如下:
stateDiagram-v2
[*] --> queued
queued --> running
running --> paused: approval.required
paused --> running: approval.approved
paused --> cancelled: approval.rejected
running --> completed: final output committed
running --> failed: unrecoverable error
running --> cancelled: cancel requested
completed --> [*]
failed --> [*]
cancelled --> [*]
注意:run.completed 不是“服务端准备发送一条完成消息”,而是“服务端已经提交了运行终态”。发送动作可以重试,但终态不能因为网络断开而回退。
七、推荐的完整事件序列
下面是一个带工具调用的完整算例。
用户输入:
杭州今天的天气怎么样?
事件:
{"seq":1,"type":"run.started","payload":{"run_id":"run_1001"}}
{"seq":2,"type":"response.created","payload":{"response_id":"resp_01","agent":"weather-agent"}}
{"seq":3,"type":"tool.call.created","payload":{"item_id":"tool_01","tool_call_id":"call_01","name":"get_weather"}}
{"seq":4,"type":"tool.arguments.delta","payload":{"tool_call_id":"call_01","index":0,"delta":"{\"city\":\"杭"}}
{"seq":5,"type":"tool.arguments.delta","payload":{"tool_call_id":"call_01","index":1,"delta":"州\"}"}}
{"seq":6,"type":"tool.arguments.finished","payload":{"tool_call_id":"call_01","arguments":{"city":"杭州"}}}
{"seq":7,"type":"tool.started","payload":{"tool_call_id":"call_01","execution_id":"exec_01"}}
{"seq":8,"type":"tool.finished","payload":{"tool_call_id":"call_01","execution_id":"exec_01","status":"succeeded","output":{"condition":"多云","temperature":28}}}
{"seq":9,"type":"response.finished","payload":{"response_id":"resp_01","finish_reason":"tool_calls"}}
{"seq":10,"type":"response.created","payload":{"response_id":"resp_02","agent":"weather-agent"}}
{"seq":11,"type":"message.delta","payload":{"item_id":"msg_01","index":0,"delta":"杭州今天"}}
{"seq":12,"type":"message.delta","payload":{"item_id":"msg_01","index":1,"delta":"多云,"}}
{"seq":13,"type":"message.delta","payload":{"item_id":"msg_01","index":2,"delta":"气温约 28°C。"}}
{"seq":14,"type":"message.finished","payload":{"item_id":"msg_01","text":"杭州今天多云,气温约 28°C。"}}
{"seq":15,"type":"response.finished","payload":{"response_id":"resp_02","finish_reason":"stop"}}
{"seq":16,"type":"usage.final","payload":{"requests":2,"input_tokens":620,"output_tokens":48,"total_tokens":668}}
{"seq":17,"type":"run.completed","payload":{"final_output":"杭州今天多云,气温约 28°C。"}}
前端行为应当是:
message.delta:更新正在显示的文本;tool.call.created:显示“正在查询天气”,但不要显示为已完成;tool.started:显示工具确实开始执行;tool.finished:显示工具结果或更新工具卡片;response.finished:结束当前模型响应;run.completed:解除运行中状态,允许发送下一条消息。
如果第 15 条事件后连接断开,前端不能直接认为 Run 成功结束,因为还没有看到 run.completed。它应通过恢复接口查询或继续消费事件。
八、传输层:SSE 和 WebSocket 只解决搬运,不解决语义
8.1 SSE
SSE 适合服务端单向推送,事件可以编码为:
id: 42
event: message.delta
data: {"protocol":"agent-events.v1","run_id":"run_1001","seq":42,...}
id 可以使用事件序号或事件 ID,客户端重连时携带 Last-Event-ID。但是否真正支持恢复,取决于服务端是否保留事件日志,不能只依赖 HTTP 连接本身。
8.2 WebSocket
WebSocket 适合双向交互,特别是:
- 人工审批;
- 客户端取消;
- 动态追加输入;
- 心跳;
- 长时间运行控制。
但 WebSocket 消息也必须使用同样的业务事件信封:
{
"type": "client.approval",
"run_id": "run_1001",
"approval_id": "approval_01",
"decision": "approve",
"idempotency_key": "approve-run_1001-approval_01"
}
不能因为 WebSocket 有消息边界,就省略 seq、event_id 和幂等键。
九、断线:先判断断在哪里,再决定是否重试
断线恢复不能简单写成:
请求失败 → 重新调用 Agent
因为重新调用可能造成:
- 文本重复;
- 工具重复执行;
- 运行费用增加;
- 人工审批状态丢失;
- Session 写入两次;
- 两个并行 Run 同时修改同一会话。
服务端应持久化至少三类状态:
Run 状态
事件日志
工具执行记录
推荐数据关系:
flowchart LR
R[Run 状态表] --> E[事件日志]
R --> C[Checkpoint]
E --> S[客户端订阅游标]
R --> T[工具执行记录]
T --> O[下游幂等请求]
9.1 断线分类
情况一:客户端断线,Run 仍在运行
服务端继续执行并写入事件日志。客户端恢复时从最后确认的 seq 之后读取。
情况二:Run 已完成,但最后完成事件未送达
恢复接口返回已提交的终态和缺失事件。客户端补齐 run.completed。
情况三:Run 在工具执行期间断线
客户端不能根据“没有收到 tool.finished”推断工具没有执行。必须查询:
GET /runs/{run_id}
GET /runs/{run_id}/events?after_seq=...
GET /tool-executions/{execution_id}
情况四:Run 暂停等待审批
恢复后应显示审批请求,而不是重新发送用户输入。审批是原 Run 的延续,不是新的用户轮次。
9.2 恢复游标
客户端每处理完一个连续事件,就更新本地游标:
last_applied_seq = 42
恢复请求:
GET /v1/runs/run_1001/events?after_seq=42
服务端返回:
{
"run_id": "run_1001",
"from_seq": 43,
"events": [
{
"seq": 43,
"type": "tool.finished",
"payload": {}
}
],
"terminal": false
}
服务端必须保证以下性质:
如果事件已被清理,必须返回明确错误:
{
"code": "cursor_expired",
"message": "events before seq=800 are no longer retained",
"snapshot_url": "/v1/runs/run_1001/snapshot"
}
不能返回空数组并伪装成“没有新事件”,否则客户端会错误地认为自己已经追平。
9.3 Snapshot 与事件重放
事件日志可以无限增长,因此通常要提供快照:
{
"run_id": "run_1001",
"snapshot_seq": 800,
"status": "running",
"items": {
"msg_01": {
"type": "message",
"text": "前面已经生成的内容",
"status": "streaming"
}
},
"pending_tools": [],
"usage": {
"total_tokens": 1500
}
}
恢复流程:
- 客户端请求快照;
- 服务端返回
snapshot_seq = 800; - 客户端将状态替换为快照;
- 客户端请求
after_seq=800; - 从 801 开始重放事件。
快照不是最终结果的替代品。它必须包含足够的信息,使事件重放从该序号继续成立。
十、恢复的伪代码
服务端恢复处理:
async def resume_run(run_id: str, after_seq: int):
run = await run_store.get(run_id)
if run is None:
raise NotFound("run_not_found")
oldest_seq = await event_store.oldest_seq(run_id)
if oldest_seq is not None and after_seq < oldest_seq - 1:
return {
"kind": "snapshot_required",
"snapshot": await snapshot_store.get(run_id),
}
events = await event_store.list_after(run_id, after_seq)
return {
"kind": "events",
"events": events,
"terminal": run.status in {
"completed",
"failed",
"cancelled",
},
"run_status": run.status,
}
客户端恢复处理:
async def reconnect(run_id: str, state: dict):
result = await api.resume(run_id, after_seq=state["last_seq"])
if result["kind"] == "snapshot_required":
state = apply_snapshot(result["snapshot"])
result = await api.resume(
run_id,
after_seq=state["snapshot_seq"],
)
for event in result["events"]:
if event["seq"] <= state["last_seq"]:
continue
if event["seq"] != state["last_seq"] + 1:
raise RuntimeError("event gap during replay")
apply_event(state, event)
state["last_seq"] = event["seq"]
if result["terminal"]:
state["mode"] = "finished"
return state
这个流程默认事件序号连续。如果协议允许服务端在某些不可见事件上跳号,则必须明确提供:
seq_range
omitted_reason
否则客户端无法区分“服务端故意省略”和“网络丢包”。
十一、基于 OpenAI Agents SDK 的端到端流式消费
下面的示例使用 Agents SDK 的公开流式接口,展示三类处理:
- 原始模型事件:读取文本 Delta;
- 高层运行项事件:读取工具调用和工具输出;
- 迭代结束后:读取最终结果和 Usage。
Agents SDK 官方文档使用 Runner.run_streamed() 创建流式运行,并通过 result.stream_events() 读取事件;原始事件中可以识别 Responses API 的文本 Delta,高层事件则包装工具调用、工具输出和消息输出。(openai.github.io)
import asyncio
from agents import Agent, Runner, ItemHelpers
from agents import function_tool
from openai.types.responses import ResponseTextDeltaEvent
@function_tool
def get_weather(city: str) -> dict:
"""Return a deterministic weather result for demonstration."""
if city != "杭州":
return {"condition": "unknown", "temperature": None}
return {
"condition": "多云",
"temperature": 28,
"unit": "celsius",
}
async def main() -> None:
agent = Agent(
name="Weather agent",
instructions=(
"回答天气问题时先调用 get_weather。"
"拿到工具结果后,用中文给出简洁回答。"
),
tools=[get_weather],
)
result = Runner.run_streamed(
agent,
input="杭州今天的天气怎么样?",
)
async for event in result.stream_events():
if event.type == "raw_response_event":
if isinstance(event.data, ResponseTextDeltaEvent):
print(event.data.delta, end="", flush=True)
elif event.type == "run_item_stream_event":
if event.name == "tool_called":
print("\n[tool called]")
elif event.name == "tool_output":
print(f"\n[tool output] {event.item}")
elif event.name == "message_output_created":
# 这是完整消息级事件,不是 token 级 Delta。
text = ItemHelpers.text_message_output(event.item)
print(f"\n[message completed] {text}")
elif event.type == "agent_updated_stream_event":
print(f"\n[agent changed] {event.new_agent.name}")
# 必须等 stream_events() 完整结束后再读取最终结果。
print("\n[run complete]")
print("final_output =", result.final_output)
usage = result.context_wrapper.usage
print("requests =", usage.requests)
print("input_tokens =", usage.input_tokens)
print("output_tokens =", usage.output_tokens)
print("total_tokens =", usage.total_tokens)
if __name__ == "__main__":
asyncio.run(main())
前置条件:
pip install openai-agents
export OPENAI_API_KEY='你的 API key'
python example.py
需要注意三个边界:
raw_response_event面向模型级低延迟输出;run_item_stream_event面向消息、工具和 handoff 等高层语义;result.final_output和result.context_wrapper.usage应在流迭代结束后读取。
SDK 的高层事件名称包含 tool_called、tool_output、message_output_created 等;其中 handoff_occured 的拼写是历史兼容行为,不能在自定义协议中照搬这个拼写。(openai.github.io)
如果工具需要审批,流会在审批点结束当前迭代,应用需将结果转换为状态、处理 interruptions,再继续运行:
result = Runner.run_streamed(
agent,
"如果临时文件不再需要,请删除它们。",
)
async for _event in result.stream_events():
pass
if result.interruptions:
state = result.to_state()
for interruption in result.interruptions:
state.approve(interruption)
resumed = Runner.run_streamed(agent, state)
async for _event in resumed.stream_events():
pass
审批暂停不是新一轮用户消息。它是同一个 Run 的可恢复状态。官方文档也要求先排空流、检查 interruptions,再从 RunState 恢复。(openai.github.io)
十二、取消、暂停和断线恢复不是同一件事
12.1 取消
取消表示系统明确要求终止运行:
run.cancel_requested
run.cancelled
Agents SDK 支持立即取消,也支持在当前 turn 完成后取消。后者用于避免在工具调用或当前响应中间直接截断状态。(openai.github.io)
12.2 暂停
暂停表示运行仍然有效,但下一步需要外部条件:
approval.required
run.paused
恢复后继续原来的状态。
12.3 断线
断线只是传输连接消失:
transport disconnected
它不改变 Run 状态。客户端必须重新查询 Run,而不能向服务端发送:
“请重新回答刚才的问题”
否则会创建另一个逻辑运行。
十三、常见错误与诊断方法
错误一:把最后一个 Delta 当作完成
表现:
- UI 显示答案已完成;
- 但工具结果还未持久化;
- 重新发送下一条消息时,上一轮上下文缺失。
诊断:
检查是否收到:
message.finished
response.finished
run.completed
真正允许关闭 Run 的条件应是 run.completed、run.failed 或 run.cancelled。
错误二:只保存最终文本,不保存事件游标
表现:
- 断线后只能重新生成;
- 重新生成的文本与第一次不同;
- 工具被重复执行。
诊断:
检查数据库中是否存在:
run_id
last_event_seq
tool_execution_id
run_status
只保存 final_output 无法支持中途恢复。
错误三:把重复事件当成新 Delta
表现:
杭州今天多云,杭州今天多云,气温约 28°C。
原因:
客户端没有按 event_id 或局部 index 去重。
修复:
对每个 item_id 保存下一个期待的 Delta 序号;对于已应用序号直接忽略,对于跳跃序号触发恢复。
错误四:工具参数尚未完整就执行
表现:
- JSON 解析失败;
- 工具收到缺失字段;
- 默认值覆盖了模型尚未生成的字段。
修复:
只在 tool.arguments.finished 后执行,且对完整参数做 schema 校验。
错误五:没有区分工具调用意图和工具副作用
表现:
前端显示“已删除文件”,但服务端只收到了模型的工具调用意图。
修复:
分别渲染:
模型请求调用工具
工具执行中
工具执行成功
工具执行失败
错误六:把缺失 Usage 记为零
表现:
成本报表显示某些请求完全免费。
原因:
上游没有返回 Usage,适配器也没有保留原始字段,但业务层将 None 转成了 0。
修复:
使用明确状态:
{
"total_tokens": null,
"status": "unavailable"
}
错误七:恢复时重新追加用户消息
表现:
恢复后的上下文变成:
用户:查询杭州天气
用户:查询杭州天气
原因:
把暂停或断线误当作新轮次。
OpenAI Agents SDK 的结果对象提供 to_state()、to_input_list()、interruptions 等不同恢复表面;其中审批暂停应使用状态恢复,而不是追加新的用户输入。(openai.github.io)
十四、协议契约的最小要求
如果要把这套协议落到 Agent API 中,最小可接受契约应包含:
运行标识
run_id
response_id
item_id
tool_call_id
execution_id
顺序与去重
event_id
seq
item_index
cursor
生命周期
run.started
run.paused
run.completed
run.failed
run.cancelled
输出
message.delta
message.finished
工具
tool.call.created
tool.arguments.delta
tool.arguments.finished
tool.started
tool.finished
tool.failed
approval.required
用量
usage.updated
usage.final
usage.status
request_usage_entries
恢复
after_seq
snapshot_seq
cursor_expired
terminal
错误
{
"code": "tool_timeout",
"retryable": true,
"message": "weather service timed out",
"phase": "tool_execution",
"tool_call_id": "call_01"
}
retryable 只能表达协议层建议,不能替代幂等保证。网络超时并不等于工具未执行;如果工具具有副作用,重试前必须查询执行记录或使用下游幂等键。
十五、协议保证、实现选择与经验建议
最后需要把三类结论分开。
协议保证
这些应写入契约并由服务端保证:
- 同一 Run 内
seq单调递增; - 事件携带稳定的
event_id; - Delta 可按
item_id和局部序号重建; run.completed表示 Run 已进入终态;- 工具执行具有可查询的幂等标识;
- 客户端可以从游标或快照恢复;
- 缺失 Usage 不伪装为零。
常见实现
这些是合理实现,但不是所有框架都必须采用:
- SSE 作为文本流传输;
- WebSocket 作为双向控制传输;
- Redis Streams、Kafka 或数据库保存事件日志;
- JSONL 作为内部事件格式;
snapshot + after_seq作为恢复机制;usage.updated发送累计值。
经验建议
这些取决于产品场景:
- 普通问答可只在结束时发送 Usage;
- 长任务应发送心跳和阶段性 Usage;
- 高风险工具应强制人工审批;
- 工具结果应按权限脱敏;
- 面向 UI 的协议和面向审计的事件可以通过适配层分开;
- 不要把模型供应商的原始事件类型直接当成业务协议。
Agent 流式协议真正要解决的不是“如何更快地显示 token”,而是如何把一次可能跨越多个模型调用、工具副作用和网络故障的运行,变成一组有顺序、有状态、有边界、可重放、可审计的事实。
Delta 负责表达增量,Tool 负责表达外部动作,Usage 负责表达消耗,Finish 负责表达边界,事件序号负责表达连续性,快照和幂等执行记录负责让断线恢复不会变成重复执行。
系列导航与关联阅读
- 系列入口:Agent 工程完整路线:从运行循环、记忆与协议到安全、评测和生产交付
- 上一篇:Agent 限流与容量:并发、Token 速率、队列、GPU 配额和过载
- 下一篇:Agent API 契约:会话、消息、附件、事件、错误和幂等键
- 延伸:AG-UI 协议:Agent 事件、前端状态、流式交互和人工确认
官方资料
本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论