Agent 工程体系 · 第 91/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。

Agent 流式事件协议:Delta、Tool、Usage、Finish、断线和恢复

流式 Agent API 不是“把模型输出按 token 打印出来”这么简单。

一次 Agent 运行可能包含多次模型调用、工具调用、工具返回、Agent handoff、人工审批、错误、重试、历史压缩和最终持久化。用户看到的文字只是其中一种事件。若协议只定义 text,前端无法可靠区分“模型正在生成”“工具正在执行”“等待人工确认”和“运行已经完成”。

本文建立一套面向生产系统的 Agent 流式事件协议,重点解决六个问题:

  1. Delta 如何表达增量文本、结构化参数和其他可拼接内容;
  2. Tool 如何表达调用、参数、执行、结果和错误;
  3. Usage 如何在多轮、多工具、多模型调用中统计;
  4. Finish 如何区分当前响应结束、当前 Agent 结束和整个 Run 完成;
  5. 断线时如何判断客户端究竟丢了什么;
  6. 恢复时如何避免重复输出、重复执行工具和重复扣费。

OpenAI Agents SDK 当前把流式事件分成三类:直接透传模型的原始事件、表示完整运行项的高层事件,以及表示当前 Agent 变化的事件。SDK 文档同时明确指出:最后一个可见 token 到达后,运行仍可能继续进行会话持久化、审批状态处理或历史压缩,只有流迭代器结束后,运行才真正完成。(openai.github.io)

因此,生产协议必须把“可见输出”与“运行生命周期”分开建模。


一、先建立正确的对象模型:消息、事件、Run 和响应

1.1 Delta 不是消息

Delta 是一次流式传输中的增量片段。它通常不能独立解释,必须与同一个逻辑输出项的前序 Delta 按顺序拼接。

例如完整文本:

杭州今天多云,最高气温 28°C。

可能被拆成:

"杭州"
"今天"
"多云,"
"最高气温"
" 28°C。"

这些片段不是五条消息,而是同一条消息的五个增量。

定义一个输出项 item,其 Delta 序列为:

D=(d1,d2,,dn)D = (d_1, d_2, \ldots, d_n)

若文本拼接运算为 ,则最终文本为:

T=d1d2dnT = d_1 \oplus d_2 \oplus \cdots \oplus d_n

这个公式成立的前提是:

  1. 所有 Delta 属于同一个 item_id
  2. Delta 按协议序号递增;
  3. 每个 Delta 只发送尚未发送的内容;
  4. 客户端不会把不同输出项的 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 类型相关的数据

seqevent_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.createdresponse.output_text.deltaresponse.completederror。(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

这段代码中的关键不是字符串拼接,而是:

  1. index < next_index 表示重复投递;
  2. index == next_index 表示可以应用;
  3. 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 解析。正确做法是:

  1. tool_call_id 建立参数缓冲区;
  2. index 校验顺序;
  3. 追加 Delta;
  4. 收到 tool.arguments.finished 后再解析;
  5. 解析失败则生成协议错误,而不是执行工具。

反例:

# 错误:每个 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

这里的约束是:

e,execute(e)1\forall e,\quad \text{execute}(e) \leq 1

更准确地说,是“具有副作用的逻辑执行最多一次”。如果底层网络请求在服务端已经发出但进程在写入 succeeded 前崩溃,数据库事务本身无法证明下游副作用是否发生。因此,对于支付、发货、删除等操作,还需要下游系统支持幂等键或查询确认。


五、Usage:用量是运行事实,不是 UI 进度条

5.1 Usage 的三种口径

Usage 至少有三个层次:

请求级

一次模型 API 请求的用量:

{
  "request_id": "req_01",
  "input_tokens": 1200,
  "output_tokens": 80,
  "reasoning_tokens": 0
}

响应级

一次模型响应可能对应一次请求,但在重试、流式重连或多供应商适配时,不应想当然地认为二者永远一一对应。

Run 级

整个 Agent Run 的累计用量:

Urun=i=1nUiU_{\text{run}} = \sum_{i=1}^{n} U_i

其中每个 UiU_i 是一次模型调用的用量,包括:

  • 初始回答;
  • 生成工具调用;
  • 工具返回后再次回答;
  • 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 完成

至少存在以下结束边界:

  1. message.finished:一条消息不再产生新的文本 Delta;
  2. response.finished:一次模型响应结束;
  3. run.completed:整个 Agent Run 完成;
  4. run.paused:运行暂停,等待审批或外部输入;
  5. run.failed:运行失败;
  6. 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 必须满足终态条件

建议定义:

run.completed{所有必需工具调用已终结没有未处理审批最终输出已确定Usage 已封存或明确不可用事件日志已持久化\text{run.completed} \Rightarrow \begin{cases} \text{所有必需工具调用已终结}\\ \text{没有未处理审批}\\ \text{最终输出已确定}\\ \text{Usage 已封存或明确不可用}\\ \text{事件日志已持久化} \end{cases}

一个简单状态机如下:

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 有消息边界,就省略 seqevent_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
}

服务端必须保证以下性质:

after_seq=k返回所有 seq>k 的可恢复事件\text{after\_seq}=k \Rightarrow \text{返回所有 } seq > k \text{ 的可恢复事件}

如果事件已被清理,必须返回明确错误:

{
  "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
  }
}

恢复流程:

  1. 客户端请求快照;
  2. 服务端返回 snapshot_seq = 800
  3. 客户端将状态替换为快照;
  4. 客户端请求 after_seq=800
  5. 从 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

需要注意三个边界:

  1. raw_response_event 面向模型级低延迟输出;
  2. run_item_stream_event 面向消息、工具和 handoff 等高层语义;
  3. result.final_outputresult.context_wrapper.usage 应在流迭代结束后读取。

SDK 的高层事件名称包含 tool_calledtool_outputmessage_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.completedrun.failedrun.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、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。