云原生 Agent 开发(一):从 LLM 响应到 Agent 会话

· Updated Agent#Agent#LLM#云原生#源码阅读

模型与应用的交互首先由 LLM API 规定。应用提交输入和工具声明,模型返回文本或工具调用;应用执行工具后,再把结果送回模型。框架把对象转换、流处理和执行循环封装成可复用能力,Agent Runtime 进一步组织会话、工具执行与客户端接口。

1. LLM API 如何定义模型与应用的交互

1.1 从一次普通文本生成开始

应用请求模型生成文本时,需要说明使用哪个模型、输入是什么。例如,让模型用一句话解释 HTTP,可以构造下面的 Chat Completions 请求。MODEL_ID 是待替换的模型标识,HTTP 头省略:

数据结构JSON 数据结构 · 1.1 从一次普通文本生成开始
object · 5 个子节点
${2 个字段}
model"MODEL_ID"
messages[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "model": "MODEL_ID",
  "messages": [
    {
      "role": "user",
      "content": "用一句话解释 HTTP。"
    }
  ]
}

messages 是提交给模型的输入序列。这里只有一条 user 消息,它的 content 是用户的问题。应用发送请求后,得到的响应也有结构。下面是一份构造的响应,只保留与本次解释有关的字段:

数据结构JSON 数据结构 · 1.1 从一次普通文本生成开始
object · 8 个子节点
${2 个字段}
id"chatcmpl_text_1"
choices[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "id": "chatcmpl_text_1",
  "choices": [
    {
      "index": 0,
      "message": {
        "role": "assistant",
        "content": "HTTP 是客户端与服务器交换请求和响应的应用层协议。"
      },
      "finish_reason": "stop"
    }
  ]
}

choices[0] 表示第一个生成候选,它的 message 是模型产生的 assistant 消息。应用取出 message.content,就能展示“HTTP 是客户端与服务器交换请求和响应的应用层协议”。响应中的 id 标识这次生成,index 标识候选位置,finish_reason 描述生成为什么停止。

这些字段会参与不同的处理:界面显示文本,日志关联生成身份,后续程序依据停止原因决定下一步。本例的 stop 表示生成在正常停止条件下结束。加入工具之后,应用还会遇到请求执行函数的返回对象,需要先完成行动,再继续生成。1

1.2 从文本输出到一次工具往返

将用户输入改为“读取 README.md,并用一句话概括项目用途”。README 位于应用的工作环境中,应用需要提供读取能力。它向模型声明 read_file(path):工具名称是 read_file,参数包含文件路径 path,实际实现负责访问文件。

工具声明描述应用可以做什么。模型选择这个工具时,返回本次调用的名称、参数和关联身份,例如 read_file、{"path":"README.md"}、call_read_1。执行器据此查找函数并读取文件。名称用来找到实现,关联身份用来说明稍后返回的文本属于哪次调用。

假设读取取得以下文本:

# Example
这是一个演示 HTTP API 的示例项目。

应用把这段内容与原调用一起加入下一次模型输入。这样模型既知道自己请求了哪个文件,也知道应用取得了什么。第二次生成便可以依据文件正文回答:“这个项目演示如何提供 HTTP API。”下面的时序使用构造数据展示这次往返:

sequenceDiagram
    participant U as 用户
    participant R as 应用或 Runtime
    participant M as 模型 API
    participant T as 工具执行器
    U->>R: 读取 README.md,概括项目用途
    R->>M: 请求 1:用户输入与 read_file 工具声明
    M-->>R: read_file 调用 call_read_1
    R->>T: 执行 read_file(README.md)
    T-->>R: README 文件文本
    R->>M: 请求 2:保留调用,追加关联结果
    M-->>R: 这个项目演示如何提供 HTTP API
    R-->>U: 展示模型生成的概括

这里发生了两次模型生成和一次工具执行。第一次模型返回读取请求,工具执行器取得正文,第二次模型使用正文形成概括。工具声明、调用和结果分别参与这三个位置,不能只保存最后的一句话。2

四类模型 API 都能表达这样的函数往返,但用于保存输入、调用和结果的对象不同。要把它们接入同一个执行循环,需要先看清这些对象如何关联。

2. 主流 LLM Response API

2.1 工具声明、调用与结果的关系

前一节中,模型的第一次响应给出读取请求,第二次响应才给出项目用途的概括。两次生成之间,应用必须把工具声明、调用和结果接起来。四类 API 的差别,首先体现在这些对象的组织方式上。

声明告诉模型“应用提供了什么能力”,包括名称、用途和参数结构;调用告诉应用“这次要使用哪项能力、传入什么参数”;结果告诉模型“这次调用得到了什么”。read_file 声明中只有参数规则,没有文件正文,也没有函数实现。文件何时被读取,取决于应用何时接收并执行那次调用。

HTTP 将这些对象送到两端。开启流式输出后,模型服务可以保持响应连接,用 SSE(Server-Sent Events)陆续发送事件。SSE 的响应类型是 text/event-stream,内容采用 UTF-8 文本:字段按行书写,空行结束一个事件。下面只展示两个参数事件,HTTP 头与内容采用教学数据:

HTTP/1.1 200 OK
Content-Type: text/event-stream
event: response.function_call_arguments.delta
data: {"type":"response.function_call_arguments.delta","item_id":"fc_read_1","output_index":0,"delta":"{\"path\":","sequence_number":2}
event: response.function_call_arguments.delta
data: {"type":"response.function_call_arguments.delta","item_id":"fc_read_1","output_index":0,"delta":"\"README.md\"}","sequence_number":3}

event 给出事件名,data 承载事件数据;LLM API 通常在 data 里放 JSON。一个事件内的多行 data 会用换行拼接,没有 event 时事件名默认为 message。SSE 还定义了 id 和 retry 字段,分别供 EventSource 记录事件标识和设置重连等待时间;以冒号开头的注释行可用作心跳。它们属于传输格式,模型的调用身份与结束状态仍由 API 的 JSON 定义。3

网络读取的一个 chunk 可能只包含半行,也可能同时包含多个事件。因此,接收端先累积字节并解码,按行读取到空行后组出 SSE 事件,再解析其中的 JSON;不能对每个网络 chunk 直接执行 JSON.parse。SSE 事件完整之后,工具参数仍可能只有 {"path": 这样的一段,应用还要按模型 API 聚合并判断调用是否定稿。4

浏览器的原生 EventSource 使用 GET;需要以 POST 提交模型输入、设置鉴权头的调用,通常由 SDK 或 fetch 读取响应流并解析同一格式。自动重连也不等于模型请求可从任意位置恢复,续接能力要看具体接口的契约。

下面用同一次文件读取对照四类接口。用户输入“读取 README.md,并用一句话概括项目用途。”,工具取得的正文统一为:

# Example
这是一个演示 HTTP API 的示例项目。

四组 JSON 都是按原生结构构造的教学数据。每组依次给出第一次请求、裁剪后的调用响应、第二次请求需要追加的历史;其中调用 ID 也是构造值,实际应用须使用响应返回的值。MODEL_ID 需替换为支持对应接口和函数工具的模型标识,HTTP 头省略。示例裁剪用于说明对象关系,实际续接应保留接口要求的完整材料。

2.2 Chat Completions:围绕消息组织工具往返

Chat Completions 用 messages 保存这次生成需要的上下文。原有的 developer 消息给出处理要求,user 消息给出任务;应用另在 tools 中声明 read_file。声明中的 parameters 描述参数形状:这里需要一个带字符串 path 的对象。

数据结构JSON 数据结构 · 2.2 Chat Completions:围绕消息组织工具往返
object · 22 个子节点
${3 个字段}
model"MODEL_ID"
messages[2 项]
tools[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "model": "MODEL_ID",
  "messages": [
    {
      "role": "developer",
      "content": "根据文件内容概括项目用途。"
    },
    {
      "role": "user",
      "content": "读取 README.md,并用一句话概括项目用途。"
    }
  ],
  "tools": [
    {
      "type": "function",
      "function": {
        "name": "read_file",
        "description": "读取文本文件。",
        "parameters": {
          "type": "object",
          "properties": {
            "path": {
              "type": "string"
            }
          },
          "required": [
            "path"
          ],
          "additionalProperties": false
        }
      }
    }
  ]
}

模型此时还没有文件正文。假设它先选择调用工具,响应会把行动请求放进 assistant message 的 tool_calls。下面裁剪掉时间、模型、usage 等响应字段,只留下这次交接需要的部分。

数据结构JSON 数据结构 · 2.2 Chat Completions:围绕消息组织工具往返
object · 15 个子节点
${2 个字段}
id"chatcmpl_demo_1"
choices[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "id": "chatcmpl_demo_1",
  "choices": [
    {
      "index": 0,
      "message": {
        "role": "assistant",
        "content": null,
        "tool_calls": [
          {
            "id": "call_read_1",
            "type": "function",
            "function": {
              "name": "read_file",
              "arguments": "{\"path\":\"README.md\"}"
            }
          }
        ]
      },
      "finish_reason": "tool_calls"
    }
  ]
}

从外到内,completion 是整次生成的响应,choices 是候选结果,message 是一个候选的消息;本例只有 index: 0 的一个候选。该消息没有文字内容,却包含一次工具调用,因而不能用 message.content 是否非空判断模型有没有返回有效结果。

应用沿 tool_calls 取得 read_file 和参数。这里的 arguments 是 JSON 编码的字符串,需要先解析成对象、校验参数,再把 path 交给读取实现。read_file 用于找到实现,call_read_1 用于标识这一次调用;同一函数被调用两次时,名称相同,结果仍须对应各自的调用 ID。

finish_reason: "tool_calls" 表示模型将工作交给工具。此时模型的这次生成已经停止,文件读取尚待应用完成。应用取得正文后,向原 messages 追加下面两条消息。这是新增的历史片段,下一次请求还需要原消息、模型和工具声明。

数据结构JSON 数据结构 · 2.2 Chat Completions:围绕消息组织工具往返
array · 14 个子节点
$[2 项]
[0]{3 个字段}
role"assistant"
contentnull
[1]{3 个字段}
role"tool"
tool_call_id"call_read_1"
content"# Example\n这是一个演示 HTTP API 的示例项目。"

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
array
完整值
[
  {
    "role": "assistant",
    "content": null,
    "tool_calls": [
      {
        "id": "call_read_1",
        "type": "function",
        "function": {
          "name": "read_file",
          "arguments": "{\"path\":\"README.md\"}"
        }
      }
    ]
  },
  {
    "role": "tool",
    "tool_call_id": "call_read_1",
    "content": "# Example\n这是一个演示 HTTP API 的示例项目。"
  }
]

第一条保存模型原先的调用,第二条用 tool_call_id 把正文接到该调用上。于是第二次生成能看到一个完整关系:assistant 请求了 read_file("README.md"),tool 返回了该文件的正文。模型可以据此概括项目用途。若将正文作为普通 user 文字追加,历史中就不再表达“这是哪次工具调用的结果”。

流式输出延续这套结构,只是消息逐步形成。chunk 中的 choice index 定位候选,delta.tool_calls[].index 再定位该候选内部的调用;参数片段追加到对应调用的字符串中。调用 ID 和名称未必在每个 chunk 重复出现,聚合器需要保存先前收到的值。它完成的工作,是将分散更新还原为上面那条 assistant message。

停止原因和计量也要分别处理。length 说明生成受到长度限制,不能与本例的工具交接混为一谈。usage 描述 token 消耗;启用流式 include_usage 后,末尾可能有 choices 为空的 usage chunk。该 chunk 不产生新消息,断流时也可能拿不到它。应用据停止原因安排下一步,据 usage 记录消耗,两者承担不同职责。

这套结构把工具往返保存在消息序列里。对封装而言,维护历史比较直观,但还需区分整次响应、候选、消息和消息里的调用;把它们压成一条“模型回答”会失去执行和续接需要的信息。14

2.3 Responses:用 items 和事件组织生成过程

Responses 把不同输出组织成独立的 item。message item 承载文字等内容,function-call item 承载调用,reasoning item 承载推理相关材料。一份 output 因而可以包含多种对象;便捷字段 output_text 只取可展示的文字,不能代替完整输出。

同一次读取在第一次请求中放进 input,处理要求放进 instructions。工具定义的 name 和 parameters 直接位于 function 工具对象中。

数据结构JSON 数据结构 · 2.3 Responses:用 items 和事件组织生成过程
object · 20 个子节点
${4 个字段}
model"MODEL_ID"
instructions"根据文件内容概括项目用途。"
input[1 项]
tools[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "model": "MODEL_ID",
  "instructions": "根据文件内容概括项目用途。",
  "input": [
    {
      "role": "user",
      "content": "读取 README.md,并用一句话概括项目用途。"
    }
  ],
  "tools": [
    {
      "type": "function",
      "name": "read_file",
      "description": "读取文本文件。",
      "parameters": {
        "type": "object",
        "properties": {
          "path": {
            "type": "string"
          }
        },
        "required": [
          "path"
        ],
        "additionalProperties": false
      },
      "strict": true
    }
  ]
}

假设这次生成只返回一个 function call,裁剪后的响应如下:

数据结构JSON 数据结构 · 2.3 Responses:用 items 和事件组织生成过程
object · 10 个子节点
${3 个字段}
id"resp_demo_1"
status"completed"
output[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "id": "resp_demo_1",
  "status": "completed",
  "output": [
    {
      "type": "function_call",
      "id": "fc_read_1",
      "call_id": "call_read_1",
      "name": "read_file",
      "arguments": "{\"path\":\"README.md\"}",
      "status": "completed"
    }
  ]
}

这里有两个不能互换的身份。fc_read_1 是输出 item 的 ID,流式事件用它定位正在生成的对象;call_read_1 是工具调用的 ID,应用用它关联执行结果。前者回答“更新哪个输出对象”,后者回答“这是谁的返回值”。

应用读取文件后,以 function_call_output 表示结果。手动维护上下文时,将实际 output 追加到旧 input,再追加执行结果。本例新增的两个 item 如下,结果中的 call_id 与前面的调用一致。

数据结构JSON 数据结构 · 2.3 Responses:用 items 和事件组织生成过程
array · 11 个子节点
$[2 项]
[0]{6 个字段}
type"function_call"
id"fc_read_1"
call_id"call_read_1"
name"read_file"
arguments"{\"path\":\"README.md\"}"
status"completed"
[1]{3 个字段}
type"function_call_output"
call_id"call_read_1"
output"# Example\n这是一个演示 HTTP API 的示例项目。"

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
array
完整值
[
  {
    "type": "function_call",
    "id": "fc_read_1",
    "call_id": "call_read_1",
    "name": "read_file",
    "arguments": "{\"path\":\"README.md\"}",
    "status": "completed"
  },
  {
    "type": "function_call_output",
    "call_id": "call_read_1",
    "output": "# Example\n这是一个演示 HTTP API 的示例项目。"
  }
]

第二次生成于是同时看到函数请求和文件正文。status: "completed" 描述的是前一次 Response 的生成状态,以及示例中 function-call item 的生成状态;它们表示调用对象已经形成,文件读取仍由应用执行。Response 完成后继续一次工具往返,是这段程序的正常路径。

实际 output 还可能包含 assistant message 和 reasoning 等 item。手动续接不能只抽取调用再丢弃其他必要材料。例如 reasoning 的 encrypted_content 是 provider 产生的 opaque 数据:应用不解释它的内部内容,但要按契约保存并回传。将输出压成文字,再重建看似相同的消息,无法恢复这些数据;手写占位字符串也不能替代它们。

Responses 也可以让服务端持有先前上下文。使用 previous_response_id 时,应用指向已有 Response,再提交新增输入;使用 Conversation 时,items 由一个持续积累的会话对象持有。两者不能在同一请求中同时指定。前一次 instructions 不会随 previous_response_id 自动继承,因此应用仍需安排后续指令。服务端负责续接模型上下文后,应用自己的用户任务、文件读取和执行状态仍需要另行管理。

流式事件把这些对象的形成过程拆开:Response 有整体生命周期,item 有创建与定稿,message 的内容 part 还有各自位置。参数 delta 用 item_id、output_index 定位,文字 delta 还包含 content_index。response.output_item.done 说明一个 item 已定稿,response.completed、response.failed、response.incomplete 才说明整次生成的结果。保存 reasoning 等续接材料时,也需要保存定稿后的对象。

独立 item 让封装可以分别处理文字、调用和续接材料;相应地,聚合器必须维护它们的身份与完成状态。下一节会用这里的 function-call item 展开参数从片段到完整对象的过程。54

2.4 Anthropic Messages:以 content blocks 表达多种内容

Anthropic Messages 仍以 message 组织历史,但一个 message 的 content 是一组有类型的 block。text block 表示文字,tool-use block 表示调用,thinking block 则携带推理相关材料。常规历史使用 user 和 assistant 角色,系统要求放在顶层 system;多轮请求由应用重新提供历史。

下面是第一次请求。工具参数规则位于 input_schema,应用仍只是在声明能力。

数据结构JSON 数据结构 · 2.4 Anthropic Messages:以 content blocks 表达多种内容
object · 18 个子节点
${5 个字段}
model"MODEL_ID"
max_tokens256
system"根据文件内容概括项目用途。"
messages[1 项]
tools[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "model": "MODEL_ID",
  "max_tokens": 256,
  "system": "根据文件内容概括项目用途。",
  "messages": [
    {
      "role": "user",
      "content": "读取 README.md,并用一句话概括项目用途。"
    }
  ],
  "tools": [
    {
      "name": "read_file",
      "description": "读取文本文件。",
      "input_schema": {
        "type": "object",
        "properties": {
          "path": {
            "type": "string"
          }
        },
        "required": [
          "path"
        ]
      }
    }
  ]
}

模型选择读取文件后,返回的 assistant message 包含一个 tool_use block。下面裁剪掉其余响应字段,保留调用与交接信息。

数据结构JSON 数据结构 · 2.4 Anthropic Messages:以 content blocks 表达多种内容
object · 11 个子节点
${5 个字段}
id"msg_demo_1"
type"message"
role"assistant"
content[1 项]
stop_reason"tool_use"

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "id": "msg_demo_1",
  "type": "message",
  "role": "assistant",
  "content": [
    {
      "type": "tool_use",
      "id": "toolu_read_1",
      "name": "read_file",
      "input": {
        "path": "README.md"
      }
    }
  ],
  "stop_reason": "tool_use"
}

此处 input 已是对象,应用可以取得 path 并校验后执行读取。toolu_read_1 标识这次调用,stop_reason: "tool_use" 表示等待客户端工具结果。应用取得正文后,把原 assistant content 和下面的结果消息追加到旧 messages:

数据结构JSON 数据结构 · 2.4 Anthropic Messages:以 content blocks 表达多种内容
array · 16 个子节点
$[2 项]
[0]{2 个字段}
role"assistant"
[1]{2 个字段}
role"user"

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
array
完整值
[
  {
    "role": "assistant",
    "content": [
      {
        "type": "tool_use",
        "id": "toolu_read_1",
        "name": "read_file",
        "input": {
          "path": "README.md"
        }
      }
    ]
  },
  {
    "role": "user",
    "content": [
      {
        "type": "tool_result",
        "tool_use_id": "toolu_read_1",
        "content": "# Example\n这是一个演示 HTTP API 的示例项目。"
      }
    ]
  }
]

结果虽然处于 user 消息中,仍由 tool_result block 明确标识为工具输出,tool_use_id 将它关联到 toolu_read_1。因此,角色为 user 不等于“人又说了一句话”;应用需要同时读取 message 角色与 block 类型,才知道这段正文的来源和用途。

回填还受历史顺序约束:结果消息要紧跟对应的 assistant 调用;同一 user 消息中若还包含普通文字,tool_result blocks 要放在前部。读取失败时可以在结果上标 is_error: true,让模型据失败信息继续处理。是否使用该标志取决于工具执行契约:读取成功后得到的文件内容本身也可能描述一个错误,不能仅凭正文含有“错误”就认为工具失败。

非流响应的 input 是完整对象,流式传输则仍需聚合参数。message_start 之后,按 index 定位的 block 依次收到 start、delta、stop;工具参数通过 input_json_delta 传入,文字、thinking、signature 使用各自的 delta。最后的 message_delta 和 message_stop 描述消息结束。usage 更新是累计值,若把每次更新直接相加,会重复计量同一次生成。

停止原因同样驱动后续动作。max_tokens 表示生成受限;服务端工具可能以 pause_turn 停在本次请求的迭代边界,调用方需要把响应续送以继续工作。工具定义来自 Anthropic,也不能据此推断执行位置:bash、text editor 是客户端工具,web search 等是服务端工具。

block 的类型与顺序还参与模型续接。thinking block 的 signature、redacted-thinking 的 data 都需要依照契约随原 block 原样保留;可见 thinking 文字并不能替代它们。把 message 的各个 block 拼成一段字符串,会同时丢掉工具关联、内容顺序和续接材料。以 message 保存历史时,完整保存 content blocks 才能保住这些关系。6

2.5 Gemini GenerateContent:围绕 Content 与 Part 组织候选结果

GenerateContent 的输入 contents 是 Content 序列。每个 Content 有角色和一组 Part;Part 可以承载文字、媒体、functionCall 或 functionResponse。响应的 candidates 是候选结果,各自持有模型生成的 Content。应用需要沿 candidate → Content → Part 读取输出,才能取得其中的调用。

模型标识位于请求 URL,下面给出请求主体。functionDeclarations 声明函数,参数采用 REST Schema 的 Type 枚举值。

数据结构JSON 数据结构 · 2.5 Gemini GenerateContent:围绕 Content 与 Part 组织候选结果
object · 23 个子节点
${3 个字段}
systemInstruction{1 个字段}
contents[1 项]
tools[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "systemInstruction": {
    "parts": [
      {
        "text": "根据文件内容概括项目用途。"
      }
    ]
  },
  "contents": [
    {
      "role": "user",
      "parts": [
        {
          "text": "读取 README.md,并用一句话概括项目用途。"
        }
      ]
    }
  ],
  "tools": [
    {
      "functionDeclarations": [
        {
          "name": "read_file",
          "description": "读取文本文件。",
          "parameters": {
            "type": "OBJECT",
            "properties": {
              "path": {
                "type": "STRING"
              }
            },
            "required": [
              "path"
            ]
          }
        }
      ]
    }
  ]
}

裁剪后的响应在一个 candidate 中返回 functionCall。这里的 args 已是对象;省略的内容包括 usage、modelVersion,以及 Part 可能携带的签名材料。

数据结构JSON 数据结构 · 2.5 Gemini GenerateContent:围绕 Content 与 Part 组织候选结果
object · 14 个子节点
${2 个字段}
responseId"response_demo_1"
candidates[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "responseId": "response_demo_1",
  "candidates": [
    {
      "index": 0,
      "content": {
        "role": "model",
        "parts": [
          {
            "functionCall": {
              "id": "call_read_1",
              "name": "read_file",
              "args": {
                "path": "README.md"
              }
            }
          }
        ]
      },
      "finishReason": "STOP"
    }
  ]
}

应用取得 path,执行读取,再把结果放进 functionResponse。下一次请求的 contents 包含旧历史、实际返回的 model Content,以及新增的结果 Content。下面只展示最后追加的结果:

数据结构JSON 数据结构 · 2.5 Gemini GenerateContent:围绕 Content 与 Part 组织候选结果
object · 8 个子节点
${2 个字段}
role"user"
parts[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "role": "user",
  "parts": [
    {
      "functionResponse": {
        "id": "call_read_1",
        "name": "read_file",
        "response": {
          "content": "# Example\n这是一个演示 HTTP API 的示例项目。"
        }
      }
    }
  ]
}

函数请求和结果各占一个 Part,调用 ID 将它们联系起来。REST schema 将 function call 的 id 定义为可选字段,Gemini 3 的调用契约则要求每次返回 ID,并在结果中匹配它。本例展示带 ID 的往返,实际应用应使用原响应里的值,遵循所用模型的契约。

保存实际 model Content,是为了把模型产生的 Part 连同元数据一并送回。thoughtSignature 是附着在 Part 上的 opaque 续接材料,需要随原 Part 回传;Gemini 3 的函数调用有强制要求,Gemini 2.5 的规则不同。非函数流式响应也可能在空 text Part 中携带签名,因此过滤空文字、合并 Part 都可能丢掉后续请求需要的数据。上面的裁剪响应不提供可回放的签名,实际程序应保存原响应对象。

REST 多轮调用由应用重传这份 contents,SDK chats 可以代为维护客户端历史。Python SDK 还可以自动执行自定义函数,把读取和回填循环一起封装;函数仍在客户端执行。这两种便利封装减少了应用手写代码,并没有把 read_file 移到 Google 的服务端。

流式 streamGenerateContent 返回一系列 GenerateContentResponse,更新仍处于 candidate/Part 结构中,没有 Responses 的那套独立 item-done 事件。finishReason 描述 candidate 的生成结束原因;本例为 STOP 时,Part 里已经有一个完整函数请求,应用接下来仍要执行它。便利 text 接口只能展示文字,执行与续接则需要这些完整 Part。7

2.6 统一封装需要保留哪些信息

至此,同一次读取已经分别形成了四种历史结构。它们都让第二次生成看到 read_file("README.md") 的请求与文件正文,但用于保存这段关系的基本对象不同。

对照维度Chat CompletionsResponsesAnthropic MessagesGemini GenerateContent
基本输出结构choice 中的 messageoutput items;message 还含 partsmessage 中的 typed blockscandidate 中的 Content/Parts
函数请求assistant tool_callsfunction_call itemtool_use blockfunctionCall Part
结果关联tool message 的 tool_call_idoutput item 的 call_idtool_result.tool_use_idfunctionResponse 的 name 和匹配 ID;依模型契约
增量定位choice index、tool indexoutput index、item ID、content indexcontent block indexcandidate/Part 结构;按具体接口与模型处理
生成结束依据finish_reasonResponse status 与终态事件stop_reason、message 边界candidate finishReason
多轮上下文应用维护 messages手动 items、previous response 或 Conversation重传 messages重传 contents;SDK chats 可代管

这些选择直接影响封装的职责。以消息为中心,工具往返随历史排列;将行动拆成独立 item 或 block,则给它们更明确的类型与位置。框架可以向上统一提供“文字、函数调用、结果”,但向下构造请求时仍需保留调用身份、对象顺序和原生续接材料。上层接口是否方便,与这些信息是否足够,必须一起判断。

执行位置又是另一项职责。本章的 read_file 是应用函数,第一次模型调用后由应用读取文件,取得结果再发起第二次模型调用。provider 的服务端工具可以在服务端执行并继续生成,例如 Responses 内置工具、Gemini code execution;SDK 自动函数调用则在客户端封装执行循环。权限、资源和取消由谁处理,要沿实际执行位置判断,不能只看工具名称或是谁提供了 schema。

响应里的各类材料有各自用途:文字供展示,usage 供计量,调用 ID 关联行动,签名及 opaque 元数据供模型续接。reasoning token 数量不是 reasoning item,可见 thinking 文字也不是其签名。只保存最终文字,足以显示一句答案,却不足以恢复这次工具往返;保存完整结构,展示、执行和后续输入才能各自取得需要的部分。846

接下来可以把这段交互变成执行循环:接收模型输出,取得完整调用,执行 read_file,保存结果,再构造下一次请求。前面的 JSON 展示了完整调用的样子;流式响应中,调用还要从增量片段逐步形成。下一节就从这些片段开始。

3. 从流式响应到工具调用循环

3.1 沿 Responses 追踪一个工具调用的形成

前面的非流式样例一次返回完整的 function_call。流式返回将这个对象分成若干事件,应用需要把属于同一调用的信息重新放在一起。以 read_file 为例,最终要形成的仍然是:

字段本次值用途
item 身份fc_read_1定位正在形成的输出对象。
输出位置output_index=0区分本次响应中的不同输出。
调用身份call_read_1将后续文件内容关联到这次调用。
名称read_file在应用工具表中找到实现。
参数{"path":"README.md"}提供本次读取的文件路径。

下面是一组构造事件,只追踪这个函数对象。响应中若还有文字或 reasoning,聚合器也需要分别保存:

到达的事件本次更新聚合后的状态与用途
response.createdresponse id、in_progress建立本次模型响应状态。
response.output_item.addeditem id、call_id、name、output_index建立 fc_read_1 的函数槽,参数缓冲区为空。
response.function_call_arguments.deltadelta={"path":"REA槽内参数成为部分 JSON。
下一条参数 deltadelta=DME.md"}拼接后得到 {"path":"README.md"}。
response.function_call_arguments.done完整 arguments 字符串使用完整值定稿参数。
response.output_item.done完整 function_call item保存这个 item 的最终表示。
response.completed完整 response、status、usage本次模型响应完成,循环可以检查下一步。

一个参数片段和参数定稿事件分别如下:

数据结构JSON 数据结构 · 3.1 沿 Responses 追踪一个工具调用的形成
object · 5 个子节点
${5 个字段}
type"response.function_call_arguments.delta"
item_id"fc_read_1"
output_index0
delta"{\"path\":\"REA"
sequence_number2

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "type": "response.function_call_arguments.delta",
  "item_id": "fc_read_1",
  "output_index": 0,
  "delta": "{\"path\":\"REA",
  "sequence_number": 2
}
数据结构JSON 数据结构 · 3.1 沿 Responses 追踪一个工具调用的形成
object · 5 个子节点
${5 个字段}
type"response.function_call_arguments.done"
item_id"fc_read_1"
output_index0
arguments"{\"path\":\"README.md\"}"
sequence_number4

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "type": "response.function_call_arguments.done",
  "item_id": "fc_read_1",
  "output_index": 0,
  "arguments": "{\"path\":\"README.md\"}",
  "sequence_number": 4
}

两条事件都指向 fc_read_1。delta 的内容需要追加,done 的内容则是完整参数,应当覆盖临时缓冲区。假如把两者都当成增量,已经闭合的 JSON 后面会再接上一份 JSON,工具执行器就无法按预期解析它。

多个函数调用可以交错到达。因此,聚合器要为每个输出对象维护独立的槽,再把 item_id、output_index 和调用身份关联起来。sequence_number 表示事件序号,不能仅凭这个字段推导出任意断点续传能力。

flowchart LR
    A[创建函数 item<br/>fc_read_1 / call_read_1] --> B[参数片段<br/>部分 JSON]
    B --> C[继续累积<br/>JSON 已闭合]
    C --> D[arguments.done<br/>使用完整参数]
    D --> E[output_item.done<br/>保存最终 item]
    E --> F[response 终态<br/>循环判定下一步]

图中的顺序描述本例从片段到完整响应的形成过程。参数完成说明获得了完整参数,item 完成说明这个输出对象定稿,response 完成说明本次生成进入终态。它们的作用范围不同。本章的教学循环等待整次模型结果,再处理工具;框架可以按自己的契约选择调度时点,稍后 Eino 和 Pi 的源码会展示具体差别。9

可以逐帧推进下面的教学演示,观察 SSE 数据进入参数槽后发生的变化。切换到截断情况,再判断“参数可以解析”和“调用可以执行”是否在同一时点成立。

01 / STREAM ASSEMBLY

从 SSE 帧到完整函数调用

同一次 read_file 调用,逐帧观察参数形成。可以解析 JSON 时,调用是否已经能交给执行器?

教学推演
本次生成的结果
事件 1 / 7

response.created

先建立 Response。此时没有函数对象,不能从响应已开始推断有工具可执行。

数据结构解析出的事件
object · 6 个子节点
${3 个字段}
type"response.created"
response{3 个字段}
id"resp_demo_1"
status"in_progress"
output[0 项]
sequence_number0

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "type": "response.created",
  "response": {
    "id": "resp_demo_1",
    "status": "in_progress",
    "output": []
  },
  "sequence_number": 0
}

继续接收。本演示选择在正常 Response 终态后交接;当前 JSON 即使可解析,也尚未满足本次交接的完成条件。

原始 SSE 帧 · 空行结束一帧
event: response.created
data: {"type":"response.created","response":{"id":"resp_demo_1","status":"in_progress","output":[]},"sequence_number":0}

这段推演采用了哪些简化?

每次只喂入一个完整 SSE 帧,省略其他事件与响应字段。实际网络块可能截断一帧,也可能包含多帧;传输层先按空行分帧,再把多行 data 用换行连接。本例只有一个函数 item,完整应用还要按 item 身份聚合其他输出。截断场景是一种构造序列,不表示所有 incomplete 响应都缺少 done 事件;实际处理还要检查终态 output 的 status 和参数。正常终态后执行是这里选择的调度策略,并非所有 Runtime 都必须采用的唯一时机。不发起模型请求或读取真实文件。

3.2 从完整消息驱动下一步执行

聚合结束后,执行循环接到两组信息:完整模型输出,以及其中可以交给应用执行的函数调用。完整输出需要进入上下文,函数调用则交给工具执行器。保存两组信息之间的关系,下一次模型请求才能恢复“模型请求了什么,应用返回了什么”。105

下面用 TypeScript 风格的伪代码表示这个循环。接口用来说明职责,不对应某个库的实际 API:model.collect 收齐结果,Context 构造和更新模型上下文,tools.executeValidated 执行已经定稿的应用函数调用。

工具循环:伪代码
type ToolCall = {
callId: string;
name: string;
argumentsText: string;
};
type ModelResult = {
outcome: "completed" | "incomplete" | "failed";
output: ModelOutput; // 完整结构与 provider 续接材料
toolCalls: ToolCall[]; // 已定稿的应用函数调用
};
async function runTask(context: Context, maxSteps: number) {
for (let step = 0; step < maxSteps; step++) {
const result: ModelResult = await model.collect(
context.buildRequest(),
);
if (result.outcome !== "completed") {
return {
kind: "generation_stopped",
outcome: result.outcome,
output: result.output,
};
}
context.appendModelOutput(result.output);
if (result.toolCalls.length === 0) {
return { kind: "finished", output: result.output };
}
for (const call of result.toolCalls) {
const toolResult = await tools.executeValidated(call);
// toolResult 保留 callId;可报告的工具错误也是一种结果
context.appendToolResult(toolResult);
}
}
return { kind: "budget_exhausted" };
}

第一轮,model.collect 返回包含 read_file 的完成结果。appendModelOutput 先保存原生输出及续接材料,工具阶段再取得 call_read_1、函数名和参数字符串。executeValidated 根据名称找到读取实现,解析 JSON、检查参数约束,最后读取 README。得到的结果保留 call_read_1,由 appendToolResult 加入上下文。

对于 Responses 手动续接,这一步保存的是原始 output items,再追加 function_call_output。第二轮的 buildRequest 读取这份更新后的上下文,把原调用和文件文本一起送入模型。若改用服务端续接,Context 就保存 response id 并提交新增输入;循环的工具路由和关联仍然需要保留。

第二轮返回最终概括,没有新的应用函数调用。循环先保存模型输出,再进入 toolCalls.length === 0 分支,返回 finished。这样,停止条件由本轮的结构化结果决定,不需要从回答文字里猜测模型是否打算继续行动。

代码还有两个边界。生成失败或截断时,它把结果返回宿主,没有将这些输出当成正常完成后继续执行工具。达到 maxSteps 后,它返回预算耗尽,交给宿主处理。文件不存在或无权读取这类可以反馈给模型的执行错误,则由工具阶段编码成关联结果,让下一轮知道尚未取得文件内容。

本例的工具由应用执行。provider 已执行的服务端工具应由适配层识别,不能看到一个带“tool”的对象就再次运行本地函数。调用参数可解析也只说明 JSON 语法成立,执行器还要检查函数 schema 和工具自己的约束。

3.3 一次工具往返怎样更新上下文

把这次交互按上下文的变化展开,可以看到每一步交给下一步的具体对象:

位置已有内容接下来的动作
第一次模型请求用户输入、指令、read_file 声明模型生成读取调用。
第一次模型响应完成原生输出中的 read_file 调用,身份为 call_read_1保存输出,解析参数并执行文件读取。
工具执行完成README 文本与 call_read_1保存关联结果,构造第二次请求。
第二次模型请求原输入、原调用与关联的文件文本模型依据 README 生成概括。
第二次模型响应完成最终概括,没有新的应用函数调用保存输出,结束本次循环。

整个用户任务跨越两次生成,工具执行位于两次生成之间。模型响应、工具执行和用户任务因此有各自的状态。框架会给它们定义消息类型与运行事件,Runtime 还会给本次工作以及所属会话分配身份。

下一步封装的起点已经明确:原生事件如何转为完整消息,工具结果如何回到上下文,以及循环如何发出可观察的进度。Eino 和 Pi 会沿这些关系提供可复用的实现。

下面将这次往返展开为可推进的上下文变化。选择一种 API,查看原调用与工具结果怎样进入下一次输入;文件读取失败时,也要用同一调用身份返回错误结果。

02 / TOOL ROUND TRIP

串起两次请求,追踪同一次调用

模型提出调用,应用执行并回填,再发起下一次生成。点击链上的节点,查看这一刻的数据与历史。

工具执行结果
请求 1 · 应用 → 模型

应用保存用户输入,声明 read_file 并发起请求 1。工具声明只描述参数规则,模型此时还没有文件正文。下方是裁剪后的历史字段。

查看响应 1 的 SSE 聚合
数据结构请求 1 · input 历史字段
object · 4 个子节点
${1 个字段}
input[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "input": [
    {
      "role": "user",
      "content": "读取 README.md,并用一句话概括项目用途。"
    }
  ]
}
展开累计 input 历史1 个条目
数据结构累计原生历史 · input
object · 4 个子节点
${1 个字段}
input[1 项]

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
object
完整值
{
  "input": [
    {
      "role": "user",
      "content": "读取 README.md,并用一句话概括项目用途。"
    }
  ]
}

教学快照,不发起模型请求或读取真实文件。请求节点仅展示原生历史字段,省略 system、工具 schema 等配置,并非完整 HTTP 请求;真实续接还必须保存协议要求的签名、reasoning 和其他 opaque 材料。调用 ID 为教学构造值。

4. Eino 与 Pi 的二次封装

上一节的最小循环把 read_file("README.md") 的往返串了起来:原生事件组成调用,文件读取返回结果,更新后的上下文进入下一次请求。框架把这些步骤变成可复用的模型组件、消息类型和执行循环。本节沿用这份 README,追踪调用和结果经过封装后存在哪里、由谁继续使用。

4.1 Eino:从统一模型接口走到 ADK 工具图

先看三个对象。模型组件接受上下文,返回模型消息;工具节点持有可以执行的 read_file;ADK 的 state.Messages 保存已经形成的消息,供下一次模型调用读取。工具声明描述名称和参数结构,交给模型;工具实例执行读取,交给工具节点。模型生成调用时,文件还没有被读取。

Eino 的 AgenticModel 连接上下文与模型输出,输入和输出都使用 AgenticMessage。消息把内容分为 ContentBlocks:文本、推理、应用函数调用与函数结果各有自己的块类型。FunctionToolCall 保存名称、参数和调用身份,FunctionToolResult 保存对应结果。

interface.go:BaseModel 与 AgenticModel 的完整声明
type BaseModel[M messageType] interface {
Generate(ctx context.Context, input []M, opts ...Option) (M, error)
Stream(ctx context.Context, input []M, opts ...Option) (*schema.StreamReader[M], error)
}
type AgenticModel = BaseModel[*schema.AgenticMessage]

AgenticModel 将泛型接口中的消息类型固定为 *schema.AgenticMessage:Generate 接受 []*schema.AgenticMessage,返回完整的 *schema.AgenticMessage;Stream 接受同一历史,返回可逐段读取的 StreamReader[*schema.AgenticMessage]。ADK 模型节点据此调用和消费模型组件,HTTP 请求与 SSE 事件的处理留在适配器内部。接口返回的是消息,工具执行由后续节点负责。

工具声明通过请求级 model.WithTools(...) option 传入。ADK 从当前状态的 ToolInfos 派生这一选项,连同 state.Messages 一起交给模型。这个选项描述本次请求允许模型提出哪些调用;执行这些调用所需的工具实例仍由工具节点持有。

现在进入第一次模型调用。Responses 适配器读取 function_call_arguments.delta,把原生参数片段写成 FunctionToolCall 的 chunk。转换器为块分配 StreamingMeta.Index:同一原生函数 item 的片段得到相同索引,后面的聚合器才能把它们放回同一调用。

responses_event_convertor.go:完整的 delta 转换函数
func (r *streamReceiver) functionCallArgumentsDeltaEventToContentBlock(ev responses.ResponseFunctionCallArgumentsDeltaEvent) *schema.ContentBlock {
meta := &schema.StreamingMeta{
Index: r.getBlockIndex(makeFunctionToolCallIndexKey(ev.OutputIndex)),
}
block := schema.NewContentBlockChunk(&schema.FunctionToolCall{
Arguments: ev.Delta,
}, meta)
setItemID(block, ev.ItemID)
return block
}

输入的 ev.Delta 仍是字符串片段,输出的 Arguments 也只包含这一段。output_index 用来找到框架块索引,item_id 被保留为块的扩展信息;框架索引由转换器分配,并不直接等于原生 output_index。此时对象的职责是记录生成进度,工具节点还没有收到完整参数。

例如,这次 read_file 被分到框架索引 0。参数先到 {"path":,再到 "README.md"};output_item.done 带来 CallID=call_read_1、Name=read_file 及 item ID、状态。ConcatAgenticMessages 按块索引归组、排序,再对同组的函数调用执行下面的合并。

agentic_message.go:完整的函数调用聚合
func concatFunctionToolCalls(calls []*FunctionToolCall) (*FunctionToolCall, error) {
if len(calls) == 0 {
return nil, fmt.Errorf("no function tool call found")
}
if len(calls) == 1 {
return calls[0], nil
}
ret := &FunctionToolCall{}
for _, c := range calls {
if c == nil {
continue
}
if ret.CallID == "" {
ret.CallID = c.CallID
} else if c.CallID != "" && c.CallID != ret.CallID {
return nil, fmt.Errorf("expected call ID '%s' for function tool call, but got '%s'", ret.CallID, c.CallID)
}
if ret.Name == "" {
ret.Name = c.Name
} else if c.Name != "" && c.Name != ret.Name {
return nil, fmt.Errorf("expected tool name '%s' for function tool call, but got '%s'", ret.Name, c.Name)
}
ret.Arguments += c.Arguments
}
return ret, nil
}

这里同时完成两件事:累积参数,确认身份。空的 CallID 或 Name 可以由后续 chunk 补齐;已有非空值时,另一段非空值必须一致,否则合并返回错误。每段 Arguments 按输入顺序相接,得到 {"path":"README.md"}。合并后的块因而同时包含执行所需的名称、参数和调用标识。

这解释了当前 Responses 适配器为什么依赖 delta:它的 item.done 转换补充身份与状态,却没有再写完整参数;事件分派也没有单独处理 function_call_arguments.done。原生协议提供了完整值,并不意味着每个适配器都把它用于聚合。这里的调用参数来自前面累积的字符串。11

接下来,完整 assistant 消息要进入 state.Messages。ADK 有两条流消费路径:事件发送 wrapper 对模型流执行 Copy(2),一份用于观察生成进度,另一份交给状态 wrapper 聚合。观察者可以逐段显示参数,工具图则等待完整消息。

wrappers.go:完整的 typedStateModelWrapper.Stream
func (w *typedStateModelWrapper[M]) Stream(ctx context.Context, _ []M, opts ...model.Option) (*schema.StreamReader[M], error) {
var (
stateMessages []M
stateToolInfos []*schema.ToolInfo
stateDeferredToolInfos []*schema.ToolInfo
)
_ = compose.ProcessState(ctx, func(_ context.Context, st *typedState[M]) error {
stateMessages = st.Messages
stateToolInfos = st.ToolInfos
stateDeferredToolInfos = st.DeferredToolInfos
return nil
})
// Backfill: old checkpoints or fresh starts have nil ToolInfos.
// Use compose-level tools from opts (which always reflects the latest bc.toolInfos)
// rather than w.toolInfos (which may be stale if the graph was reused).
if stateToolInfos == nil {
composeLevelOpts := model.GetCommonOptions(&model.Options{}, opts...)
if composeLevelOpts.Tools != nil {
stateToolInfos = composeLevelOpts.Tools
} else {
stateToolInfos = w.toolInfos
}
}
state := &TypedChatModelAgentState[M]{
Messages: stateMessages,
ToolInfos: stateToolInfos,
DeferredToolInfos: stateDeferredToolInfos,
}
if msgState, ok := any(state).(*ChatModelAgentState); ok {
for _, m := range w.middlewares {
if m.BeforeChatModel != nil {
if err := m.BeforeChatModel(ctx, msgState); err != nil {
return nil, err
}
}
}
}
baseOpts := &model.Options{Tools: w.toolInfos}
commonOpts := model.GetCommonOptions(baseOpts, opts...)
mc := &TypedModelContext[M]{Tools: commonOpts.Tools, ModelRetryConfig: w.modelRetryConfig, cancelContext: w.cancelContext}
for _, handler := range w.handlers {
var err error
ctx, state, err = handler.BeforeModelRewriteState(ctx, state, mc)
if err != nil {
return nil, err
}
}
// Persist state (including tool infos) after BeforeModelRewriteState.
_ = compose.ProcessState(ctx, func(_ context.Context, st *typedState[M]) error {
st.Messages = state.Messages
st.ToolInfos = state.ToolInfos
st.DeferredToolInfos = state.DeferredToolInfos
return nil
})
// Derive model options from state. Append after caller opts so state takes precedence
// (model.GetCommonOptions applies left-to-right, last wins).
// Use explicit copy to avoid mutating the caller's opts slice.
derivedOpts := make([]model.Option, len(opts), len(opts)+2)
copy(derivedOpts, opts)
derivedOpts = append(derivedOpts, model.WithTools(state.ToolInfos))
if state.DeferredToolInfos != nil {
derivedOpts = append(derivedOpts, model.WithDeferredTools(state.DeferredToolInfos))
}
wrappedEndpoint := w.wrapStreamEndpoint(w.inner.Stream)
stream, err := wrappedEndpoint(ctx, state.Messages, derivedOpts...)
if err != nil {
return nil, err
}
result, err := concatMessageStream(stream)
if err != nil {
return nil, err
}
// Re-read State.Messages after Stream completes: same rationale as in Generate above.
if w.modelRetryConfig != nil && w.modelRetryConfig.ShouldRetry != nil {
_ = compose.ProcessState(ctx, func(_ context.Context, st *typedState[M]) error {
state.Messages = st.Messages
return nil
})
}
state.Messages = append(state.Messages, result)
for _, handler := range w.handlers {
ctx, state, err = handler.AfterModelRewriteState(ctx, state, mc)
if err != nil {
return nil, err
}
}
if msgState, ok := any(state).(*ChatModelAgentState); ok {
for _, m := range w.middlewares {
if m.AfterChatModel != nil {
if err := m.AfterChatModel(ctx, msgState); err != nil {
return nil, err
}
}
}
}
// Persist state (including tool infos) after AfterModelRewriteState.
_ = compose.ProcessState(ctx, func(_ context.Context, st *typedState[M]) error {
st.Messages = state.Messages
st.ToolInfos = state.ToolInfos
st.DeferredToolInfos = state.DeferredToolInfos
return nil
})
if len(state.Messages) == 0 {
return nil, errors.New("no messages left in state after model call")
}
return schema.StreamReaderFromArray([]M{state.Messages[len(state.Messages)-1]}), nil
}

模型请求读取 state.Messages。concatMessageStream 消费完流后产生 result,错误在这里返回;成功后,完整消息才追加到状态。重试钩子存在时先重读共享状态,再追加结果,避免继续使用调用前保存的旧消息列表。后置处理完成后,该函数将消息和工具信息写回图状态,并把最后一条完整消息作为单值流交给后续节点。展示端消费的是生成过程,工具节点消费的是聚合后的调用。

工具图随后查看最新 assistant 消息。AgenticMessage 可以包含文本、推理、应用函数和 provider 服务端工具等不同块;AgenticToolsNode 只提取应用函数,将其转成已有 ToolsNode 能处理的调用表示。

agentic_tools_node.go:完整的 AgenticToolsNode.Invoke
func (a *AgenticToolsNode) Invoke(ctx context.Context, input *schema.AgenticMessage, opts ...ToolsNodeOption) ([]*schema.AgenticMessage, error) {
result, err := a.inner.Invoke(ctx, agenticMessageToToolCallMessage(input), opts...)
if err != nil {
return nil, err
}
return toolMessageToAgenticMessage(result), nil
}

调用与结果转换保留了 CallID、名称、参数和 user role 的结果块:

agentic_tools_node.go:完整的调用与结果转换函数
func agenticMessageToToolCallMessage(input *schema.AgenticMessage) *schema.Message {
var tc []schema.ToolCall
for _, block := range input.ContentBlocks {
if block.Type != schema.ContentBlockTypeFunctionToolCall || block.FunctionToolCall == nil {
continue
}
tc = append(tc, schema.ToolCall{
ID: block.FunctionToolCall.CallID,
Function: schema.FunctionCall{
Name: block.FunctionToolCall.Name,
Arguments: block.FunctionToolCall.Arguments,
},
Extra: block.Extra,
})
}
return &schema.Message{
Role: schema.Assistant,
ToolCalls: tc,
}
}
func toolMessageToAgenticMessage(input []*schema.Message) []*schema.AgenticMessage {
results := make([]*schema.AgenticMessage, len(input))
for i, m := range input {
if msg, ok := toolSearchResultMessageToAgenticMessage(m, nil); ok {
results[i] = msg
continue
}
ftr := &schema.FunctionToolResult{
CallID: m.ToolCallID,
Name: m.ToolName,
}
if len(m.UserInputMultiContent) > 0 {
ftr.Content = messageInputPartsToFunctionToolBlocks(m.UserInputMultiContent)
} else if m.Content != "" {
ftr.Content = []*schema.FunctionToolResultContentBlock{
newFuncToolResultContentBlock(&schema.UserInputText{Text: m.Content}),
}
}
results[i] = &schema.AgenticMessage{
Role: schema.AgenticRoleTypeUser,
ContentBlocks: []*schema.ContentBlock{{
Type: schema.ContentBlockTypeFunctionToolResult,
FunctionToolResult: ftr,
Extra: m.Extra,
}},
Extra: m.Extra,
}
}
return results
}

输入和结果都使用 AgenticMessage。agenticMessageToToolCallMessage 先提取类型为 FunctionToolCall 且内容非空的块,将调用身份、名称和参数交给内部 ToolsNode 执行;toolMessageToAgenticMessage 再把执行结果转回函数结果块。这里的内部转换复用了已有执行器,ADK 传递的消息和 state.Messages 仍沿 AgenticMessage 组织。对本例,工具节点据此选择 read_file,解析 {"path":"README.md"} 并执行文件读取;错误则从 Invoke 返回,不能当成读取成功。

读取返回 # Example 和“这是一个演示 HTTP API 的示例项目。”两行文本。toolMessageToAgenticMessage 将执行器结果转成 user role 的 FunctionToolResult,其中 CallID 仍为 call_read_1;ADK 的 afterToolCalls 再把它追加到 state.Messages。这次循环的状态变化可以放在一起看:

时刻state.Messages 中新增的内容下游读取者
第一次模型调用前用户请求“读取 README.md,并用一句话概括项目用途。”模型组件
模型流完整聚合后assistant 的 read_file 调用,参数与 call_read_1 已形成工具节点
文件读取返回后对应 call_read_1 的文件内容第二次模型调用
第二次模型调用后assistant 的“这个项目演示如何提供 HTTP API。”结果消费者

图中的回边让模型节点再次读取这份状态。第二次请求包含“模型要求读哪个文件”和“这次读取实际返回什么”,模型才能据此生成概括。整个往返仍是两次模型调用、一次工具执行;框架承担的是对象转换、流消费和状态传递。

Eino 还提供两类续接。Responses 适配器的 EnableAutoCache 可以提取已标记消息中的 response ID,以 previous_response_id 续接 provider 上下文并只发送后续输入;ADK Runner 的 checkpoint 则保存继续框架执行所需的状态。前者连接模型服务里的上下文,后者连接执行图。Eino 也提供 TurnLoop 等更外层的控制能力,选用组件时仍应按实际职责理解它们。12

4.2 Pi:从 provider events 走到 agent loop

Pi 也先建立模型层的统一表示。pi-ai 负责 provider adapter 与模型消息,pi-agent-core 接收这些消息并执行工具循环,pi-coding-agent 再把循环纳入产品会话。继续追踪同一个 read_file,可以看到每个包交给下一层什么对象。

pi-ai 用 AssistantMessage.content 保存 text、thinking、toolCall;其中 ToolCall.arguments 是解析后的 JSON 对象。生成中的消息通过 AssistantMessageEvent 更新,文本和工具调用都有各自的 start、delta、end;整个模型流最后以 done 或 error 交付最终消息。模型来源、response ID、原生 stop reason 和签名另有字段,后续回放仍可使用这些信息。

Responses 适配器先建立 Map<output_index, slot>。slot 保存原生输出对应的统一内容块及其位置:函数 item fc_read_1 创建 toolCall,ID 为 call_read_1|fc_read_1,用一个字段保存原协议的调用 ID 与 item ID;参数字符串暂存于 partialJson。后续 delta 通过 output_index 找到同一个 slot。

这个 slot 内同时存在两种参数表示:partialJson 保存收到的原始片段,arguments 保存当前能够解析出来的对象。参数继续生成时,后者便于展示;完整值到达时,前者被完整字符串替换。

openai-responses-shared.ts:完整的 processResponsesStream
export async function processResponsesStream<TApi extends Api>(
openaiStream: AsyncIterable<ResponseStreamEvent>,
output: AssistantMessage,
stream: AssistantMessageEventStream,
model: Model<TApi>,
options?: OpenAIResponsesStreamOptions,
): Promise<void> {
let sawTerminalResponseEvent = false;
const outputSlots = new Map<number, ResponsesOutputSlot>();
const reasoningBlocksById = new Map<string, ThinkingContent>();
const applyMessagePhaseStopReason = (item: ResponseOutputItem): void => {
if (item.type === "message" && item.phase === "final_answer") {
output.stopReason = "stop";
}
};
const getSlot = <TType extends ResponsesOutputSlot["type"]>(
outputIndex: number,
type: TType,
): Extract<ResponsesOutputSlot, { type: TType }> | undefined => {
const slot = outputSlots.get(outputIndex);
return slot?.type === type ? (slot as Extract<ResponsesOutputSlot, { type: TType }>) : undefined;
};
const pushToolCallDelta = (slot: ToolCallOutputSlot, delta: string | undefined): void => {
if (delta === undefined) return;
stream.push({
type: "toolcall_delta",
contentIndex: slot.contentIndex,
delta,
partial: output,
});
};
const createSlot = (outputIndex: number, item: ResponseOutputItem): ResponsesOutputSlot | undefined => {
if (item.type === "reasoning") {
const block: ThinkingContent = { type: "thinking", thinking: "" };
output.content.push(block);
const slot = {
type: "thinking",
block,
contentIndex: output.content.length - 1,
} satisfies ResponsesOutputSlot;
outputSlots.set(outputIndex, slot);
stream.push({ type: "thinking_start", contentIndex: slot.contentIndex, partial: output });
return slot;
}
if (item.type === "message") {
applyMessagePhaseStopReason(item);
const block: TextContent = { type: "text", text: "" };
output.content.push(block);
const slot = { type: "text", block, contentIndex: output.content.length - 1 } satisfies ResponsesOutputSlot;
outputSlots.set(outputIndex, slot);
stream.push({ type: "text_start", contentIndex: slot.contentIndex, partial: output });
return slot;
}
if (item.type === "function_call") {
const block: StreamingToolCall = {
type: "toolCall",
id: `${item.call_id}|${item.id}`,
name: item.name,
arguments: {},
...(item.namespace !== undefined ? { namespace: item.namespace } : {}),
partialJson: item.arguments || "",
};
output.content.push(block);
const slot = {
type: "toolCall",
block,
contentIndex: output.content.length - 1,
} satisfies ResponsesOutputSlot;
outputSlots.set(outputIndex, slot);
stream.push({ type: "toolcall_start", contentIndex: slot.contentIndex, partial: output });
return slot;
}
if (item.type === "custom_tool_call") {
const inputProperty = options?.grammarToolInputProperties?.get(item.name) ?? "input";
const input = item.input || "";
const block: StreamingToolCall = {
type: "toolCall",
id: `${item.call_id}|${item.id}`,
name: item.name,
arguments: { [inputProperty]: input },
...(item.namespace !== undefined ? { namespace: item.namespace } : {}),
customInput: {
property: inputProperty,
jsonBuffer: { input: "", started: false, closed: false },
},
};
output.content.push(block);
const slot = {
type: "toolCall",
block,
contentIndex: output.content.length - 1,
} satisfies ResponsesOutputSlot;
outputSlots.set(outputIndex, slot);
stream.push({ type: "toolcall_start", contentIndex: slot.contentIndex, partial: output });
return slot;
}
return undefined;
};
const getOrCreateSlot = (outputIndex: number, item: ResponseOutputItem): ResponsesOutputSlot | undefined => {
return outputSlots.get(outputIndex) ?? createSlot(outputIndex, item);
};
// Azure OpenAI can omit reasoning.encrypted_content from response.output_item.done
// and provide it only in response.completed.response.output. Backfill the
// persisted reasoning signature from the terminal response to keep store:false
// multi-turn replay stateless. See https://github.com/earendil-works/pi/issues/6409.
const backfillReasoningSignatures = (responseOutput: ResponseOutputItem[]): void => {
for (const item of responseOutput) {
if (item.type !== "reasoning" || !item.encrypted_content) continue;
const block = reasoningBlocksById.get(item.id);
if (!block?.thinkingSignature) continue;
const storedItem = JSON.parse(block.thinkingSignature) as ResponseReasoningItem;
if (storedItem.encrypted_content) continue;
block.thinkingSignature = JSON.stringify({
...storedItem,
encrypted_content: item.encrypted_content,
});
}
};
const finalizeResponse = (
response: Extract<ResponseStreamEvent, { type: "response.completed" | "response.incomplete" }>["response"],
): void => {
sawTerminalResponseEvent = true;
backfillReasoningSignatures(response.output ?? []);
if (response?.id) {
output.responseId = response.id;
}
if (response?.usage) {
const inputDetails = response.usage.input_tokens_details as
| { cached_tokens?: number; cache_write_tokens?: number }
| undefined;
const cachedTokens = inputDetails?.cached_tokens || 0;
const cacheWriteTokens = inputDetails?.cache_write_tokens || 0;
output.usage = {
// OpenAI includes cached and cache-write tokens in input_tokens, so subtract both.
input: Math.max(0, (response.usage.input_tokens || 0) - cachedTokens - cacheWriteTokens),
output: response.usage.output_tokens || 0,
cacheRead: cachedTokens,
cacheWrite: cacheWriteTokens,
reasoning: response.usage.output_tokens_details?.reasoning_tokens || 0,
totalTokens: response.usage.total_tokens || 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
}
calculateCost(model, output.usage);
if (options?.applyServiceTierPricing) {
const serviceTier = options.resolveServiceTier
? options.resolveServiceTier(response?.service_tier, options.serviceTier)
: (response?.service_tier ?? options.serviceTier);
options.applyServiceTierPricing(output.usage, serviceTier);
}
// Map status to stop reason. For incomplete responses, retain the provider's
// specific reason so max-output truncation and content filtering stay distinct.
const status = response?.status;
const incompleteDetails = response?.incomplete_details as { reason?: unknown } | null | undefined;
const incompleteReason = typeof incompleteDetails?.reason === "string" ? incompleteDetails.reason : undefined;
output.rawStopReason = incompleteReason ? `${status}.${incompleteReason}` : status;
const mappedStop = mapStopReason(status, incompleteReason);
output.stopReason = mappedStop.stopReason;
if (mappedStop.errorMessage === undefined) delete output.errorMessage;
else output.errorMessage = mappedStop.errorMessage;
if (output.content.some((b) => b.type === "toolCall") && output.stopReason === "stop") {
output.stopReason = "toolUse";
}
};
for await (const event of openaiStream) {
await options?.onProviderStreamEvent?.(event, model);
if (event.type === "response.created") {
output.responseId = event.response.id;
} else if (event.type === "response.output_item.added") {
createSlot(event.output_index, event.item);
} else if (event.type === "response.reasoning_summary_text.delta") {
const slot = getSlot(event.output_index, "thinking");
if (!slot) continue;
slot.block.thinking += event.delta;
stream.push({
type: "thinking_delta",
contentIndex: slot.contentIndex,
delta: event.delta,
partial: output,
});
} else if (event.type === "response.reasoning_summary_part.done") {
const slot = getSlot(event.output_index, "thinking");
if (!slot) continue;
slot.block.thinking += "\n\n";
stream.push({
type: "thinking_delta",
contentIndex: slot.contentIndex,
delta: "\n\n",
partial: output,
});
} else if (event.type === "response.reasoning_text.delta") {
const slot = getSlot(event.output_index, "thinking");
if (!slot) continue;
slot.block.thinking += event.delta;
stream.push({
type: "thinking_delta",
contentIndex: slot.contentIndex,
delta: event.delta,
partial: output,
});
} else if (event.type === "response.output_text.delta") {
const slot = getSlot(event.output_index, "text");
if (!slot) continue;
slot.block.text += event.delta;
stream.push({
type: "text_delta",
contentIndex: slot.contentIndex,
delta: event.delta,
partial: output,
});
} else if (event.type === "response.refusal.delta") {
const slot = getSlot(event.output_index, "text");
if (!slot) continue;
slot.block.text += event.delta;
stream.push({
type: "text_delta",
contentIndex: slot.contentIndex,
delta: event.delta,
partial: output,
});
} else if (event.type === "response.function_call_arguments.delta") {
const slot = getSlot(event.output_index, "toolCall");
if (!slot || slot.block.partialJson === undefined) continue;
slot.block.partialJson += event.delta;
slot.block.arguments = parseStreamingJson(slot.block.partialJson);
pushToolCallDelta(slot, event.delta);
} else if (event.type === "response.function_call_arguments.done") {
const slot = getSlot(event.output_index, "toolCall");
if (!slot || slot.block.partialJson === undefined) continue;
const previousPartialJson = slot.block.partialJson;
slot.block.partialJson = event.arguments;
slot.block.arguments = parseStreamingJson(slot.block.partialJson);
if (event.arguments.startsWith(previousPartialJson)) {
const delta = event.arguments.slice(previousPartialJson.length);
if (delta.length > 0) pushToolCallDelta(slot, delta);
}
} else if (event.type === "response.custom_tool_call_input.delta") {
const slot = getSlot(event.output_index, "toolCall");
if (!slot || !slot.block.customInput) continue;
pushToolCallDelta(
slot,
appendCustomToolCallInput(slot.block, getCustomToolCallInput(slot.block) + event.delta, false),
);
} else if (event.type === "response.custom_tool_call_input.done") {
const slot = getSlot(event.output_index, "toolCall");
if (!slot || !slot.block.customInput) continue;
pushToolCallDelta(slot, appendCustomToolCallInput(slot.block, event.input, true));
} else if (event.type === "response.output_item.done") {
const item = event.item;
applyMessagePhaseStopReason(item);
const slot = getOrCreateSlot(event.output_index, item);
if (item.type === "reasoning" && slot?.type === "thinking") {
const summaryText = item.summary?.map((s) => s.text).join("\n\n") || "";
const contentText = item.content?.map((c) => c.text).join("\n\n") || "";
slot.block.thinking = summaryText || contentText || slot.block.thinking;
slot.block.thinkingSignature = JSON.stringify(item);
reasoningBlocksById.set(item.id, slot.block);
stream.push({
type: "thinking_end",
contentIndex: slot.contentIndex,
content: slot.block.thinking,
partial: output,
});
outputSlots.delete(event.output_index);
} else if (item.type === "message" && slot?.type === "text") {
slot.block.text = item.content?.map((c) => (c.type === "output_text" ? c.text : c.refusal)).join("") || "";
slot.block.textSignature = encodeTextSignatureV1(item.id, item.phase ?? undefined);
stream.push({
type: "text_end",
contentIndex: slot.contentIndex,
content: slot.block.text,
partial: output,
});
outputSlots.delete(event.output_index);
} else if (
item.type === "function_call" &&
slot?.type === "toolCall" &&
slot.block.partialJson !== undefined
) {
slot.block.arguments = parseStreamingJson(item.arguments || slot.block.partialJson || "{}");
if (item.namespace !== undefined) slot.block.namespace = item.namespace;
// Finalize in-place and strip the scratch buffer so replay only
// carries parsed arguments.
delete slot.block.partialJson;
stream.push({
type: "toolcall_end",
contentIndex: slot.contentIndex,
toolCall: slot.block,
partial: output,
});
outputSlots.delete(event.output_index);
} else if (item.type === "custom_tool_call" && slot?.type === "toolCall" && slot.block.customInput) {
pushToolCallDelta(
slot,
appendCustomToolCallInput(slot.block, item.input ?? getCustomToolCallInput(slot.block), true),
);
if (item.namespace !== undefined) slot.block.namespace = item.namespace;
delete slot.block.customInput;
stream.push({
type: "toolcall_end",
contentIndex: slot.contentIndex,
toolCall: slot.block,
partial: output,
});
outputSlots.delete(event.output_index);
}
} else if (event.type === "response.completed" || event.type === "response.incomplete") {
finalizeResponse(event.response);
} else if (event.type === "error") {
throw new Error(`Error Code ${event.code}: ${event.message}` || "Unknown error");
} else if (event.type === "response.failed") {
sawTerminalResponseEvent = true;
output.rawStopReason = event.response?.status;
const error = event.response?.error;
const details = event.response?.incomplete_details;
const msg = error
? `${error.code || "unknown"}: ${error.message || "no message"}`
: details?.reason
? `incomplete: ${details.reason}`
: "Unknown error (no error details in response)";
throw new Error(msg);
}
}
if (!sawTerminalResponseEvent) {
throw new Error("OpenAI Responses stream ended before a terminal response event");
}
// The agent runs every tool call in the final message. Refuse to hand over calls whose
// output_item.done never arrived: their arguments may be cut off or mixed up, e.g. when a
// non-compliant server omits output_index. Finished calls have their scratch buffers removed.
if (output.stopReason === "toolUse") {
for (const block of output.content) {
if (block.type !== "toolCall") continue;
const toolCall = block as StreamingToolCall;
if (toolCall.partialJson !== undefined || toolCall.customInput !== undefined) {
throw new Error(
`OpenAI Responses stream completed with an unfinished tool call: ${toolCall.name} (${toolCall.id})`,
);
}
}
}
}

delta 分支把片段追加到 partialJson,再更新解析对象并发出展示增量。done 分支先检查 slot 仍在聚合,然后直接用 event.arguments 覆盖临时串。只有完整串包含此前前缀时,才把新增后缀再发成 delta。这能避免把完整参数错误地追加第二遍,也保留生成末尾尚未展示的内容。

例如,收到 {"path":"REA 时,slot 仍持有临时串;随后到达完整值 {"path":"README.md"},arguments 更新为 {path: "README.md"}。接着,output_item.done 使用 item 的完整参数定稿 toolCall,删除 partialJson,发出 toolcall_end 并移除该 output_index 的 slot。最终消息只留下解析后的参数,生成过程中的临时串到这里完成职责。

一个调用已经定稿,整个模型响应也要结束,adapter 才能把最终消息交给 agent loop。processResponsesStream 在退出事件循环后检查这两层状态。

在上方的完整 processResponsesStream 中,高亮的第 328–344 行就是流结束检查:先要求收到终态事件,再拒绝交接仍保留参数缓冲的工具调用。两组检查与参数聚合发生在同一次流处理内。

response.completed、response.incomplete 等分支先记录 response 终态。没有看到终态时,上面的第一处判断返回错误。模型要交付工具调用时,第二处判断检查每个工具块是否还带着参数 scratch;删除 scratch 发生在 item 定稿时,所以残留意味着调用尚未完成。最终消息能够解析参数、调用已经定稿、response 已经结束,分别由不同步骤确立。

在 pi-agent-core 中,streamAssistantResponse 把这条模型流映射成 message_start/update/end。请求前,它先用可选的 transformContext 处理 Agent 上下文,再用 convertToLlm 和 normalizeContext 形成模型输入;消费时先把 partial 放进 context.messages 的最后一个位置,逐次更新,最后用终态消息替换。返回给工具循环的 message 因而已经是完整 assistant 消息。

循环从这条消息选出 toolCall,执行工具,再把结果追加到下一次请求要使用的上下文。

agent-loop.ts:完整的 runLoop
async function runLoop(
initialContext: AgentContext,
newMessages: AgentMessage[],
initialConfig: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
streamFunction: StreamFn,
): Promise<void> {
let currentContext = initialContext;
let config = initialConfig;
let lastCompletedTurn: PrepareNextTurnContext | undefined;
let explicitContinuation = false;
// Check for steering messages at start (user may have typed while waiting)
let pendingMessages: AgentMessage[] = (await config.getSteeringMessages?.()) || [];
// Outer loop: continues when queued follow-up messages arrive after agent would stop
while (true) {
let hasMoreToolCalls = true;
// Inner loop: process tool calls and steering messages
while (hasMoreToolCalls || pendingMessages.length > 0) {
let preparedMessages: AgentMessage[] = [];
if (lastCompletedTurn) {
const nextTurnSnapshot = await config.prepareNextTurn?.(lastCompletedTurn);
if (nextTurnSnapshot) {
currentContext = nextTurnSnapshot.context ?? currentContext;
preparedMessages = nextTurnSnapshot.messages ?? [];
config = {
...config,
model: nextTurnSnapshot.model ?? config.model,
reasoning:
nextTurnSnapshot.thinkingLevel === undefined
? config.reasoning
: nextTurnSnapshot.thinkingLevel === "off"
? undefined
: nextTurnSnapshot.thinkingLevel,
};
}
// Preparation can be long-running (for example, compaction). Pick up steering
// queued while it ran. Only poll again if the earlier poll returned nothing;
// otherwise one-at-a-time mode would deliver two messages in this turn.
if (pendingMessages.length === 0) {
pendingMessages = (await config.getSteeringMessages?.()) || [];
}
await emit({ type: "turn_start" });
}
// Process prepared and queued messages before the next assistant response.
for (const message of declareToolChanges(currentContext, [...preparedMessages, ...pendingMessages])) {
await emit({ type: "message_start", message });
await emit({ type: "message_end", message });
currentContext.messages.push(message);
newMessages.push(message);
}
pendingMessages = [];
const requestUpdate = await config.prepareRequest?.(
{
context: currentContext,
model: config.model,
thinkingLevel: config.reasoning ?? "off",
},
signal,
);
if (requestUpdate) {
currentContext = requestUpdate.context ?? currentContext;
config = {
...config,
model: requestUpdate.model ?? config.model,
reasoning:
requestUpdate.thinkingLevel === undefined
? config.reasoning
: requestUpdate.thinkingLevel === "off"
? undefined
: requestUpdate.thinkingLevel,
};
}
// Stream assistant response
const message = await streamAssistantResponse(currentContext, config, signal, emit, streamFunction);
newMessages.push(message);
if (message.stopReason === "error" || message.stopReason === "aborted") {
lastCompletedTurn = {
message,
toolResults: [],
context: currentContext,
newMessages,
};
await config.finishTurn?.(lastCompletedTurn, signal);
await emit({ type: "turn_end", message, toolResults: [] });
await emit({ type: "agent_end", messages: newMessages });
return;
}
// Check for tool calls
const toolCalls = message.content.filter((c) => c.type === "toolCall");
const toolResults: ToolResultMessage[] = [];
hasMoreToolCalls = false;
if (toolCalls.length > 0) {
// A "length" stop means the output was cut off by the token limit, so
// every tool call in the message may carry truncated arguments. Fail
// them all instead of executing potentially borked calls.
const executedToolBatch =
message.stopReason === "length"
? await failToolCallsFromTruncatedMessage(toolCalls, emit)
: await executeToolCalls(currentContext, message, config, signal, emit);
toolResults.push(...executedToolBatch.messages);
hasMoreToolCalls = !executedToolBatch.terminate;
for (const result of toolResults) {
currentContext.messages.push(result);
newMessages.push(result);
}
}
lastCompletedTurn = {
message,
toolResults,
context: currentContext,
newMessages,
};
const decision = await config.finishTurn?.(lastCompletedTurn, signal);
await emit({ type: "turn_end", message, toolResults });
if (decision?.action === "end") {
await emit({ type: "agent_end", messages: newMessages });
return;
}
explicitContinuation = decision?.action === "continue";
pendingMessages = (await config.getSteeringMessages?.()) || [];
if (hasMoreToolCalls || pendingMessages.length > 0) {
explicitContinuation = false;
}
}
// Agent would stop here. Check for follow-up messages.
const followUpMessages = (await config.getFollowUpMessages?.()) || [];
if (followUpMessages.length > 0) {
// Set as pending so inner loop processes them
explicitContinuation = false;
pendingMessages = followUpMessages;
continue;
}
// No natural request was selected, so fulfill the continuation decision with one context-only turn.
if (explicitContinuation) {
explicitContinuation = false;
continue;
}
// No more messages, exit
break;
}
await emit({ type: "agent_end", messages: newMessages });
}

这里的 toolCalls 读自刚完成的模型消息。正常路径由 executeToolCalls 根据名称和参数执行工具,返回一个批次;hasMoreToolCalls 记录批次是否要求循环继续。每条结果写入 currentContext.messages,供后续模型读取,也写入 newMessages,供本次低层运行汇总新增消息。

本例的文件内容由此成为一条 ToolResultMessage,toolCallId 为 call_read_1|fc_read_1。下一次 Responses 请求转换遇到这条消息时,拆出 call_read_1,写成 function_call_output.call_id;文件文本成为 output。组合 ID 在统一消息内部保留两种身份,adapter 在 provider 输入边界恢复协议要求的关联。模型看到 README 内容后,生成“这个项目演示如何提供 HTTP API。”,主例的往返完成。

代码里的 length 分支说明执行决策还使用模型的结束原因:输出达到 token 上限时,调用参数可能已被截断,即使临时解析能得到对象,也不足以确认原参数完整。循环把这一批调用转成错误结果,让模型重新发起调用;error、aborted 则在前面的分支结束低层循环。消息里的内容和结束状态共同决定下一步。

工具批次还可以有多个独立调用。暂设一条 assistant 消息调用 A、B,B 更早完成:parallel 路径可能先发出 B 的 tool_execution_end,再发 A 的;最终 Promise.all 仍按调用数组顺序收集结果,构造 A、B 两条结果消息。完成顺序用来报告运行进度,调用顺序用来组装后续上下文。这个例子只说明并发的顺序关系;前面的 README 往返始终只有一个工具调用。

更外层的 AgentSession 在 message_end 后把消息交给 SessionManager 保存。后者用 entry 的 id、parentId 与当前 leaf 维护追加树,JSONL 可以保留多个分支;buildSessionProjection 从当前分支及压缩记录中形成模型上下文。例如,用户请求、读取调用、读取结果和概括沿父子关系相连,当前 leaf 指向概括;切换分支会改变下一次请求的投影,已有记录仍可保存在会话中。会话历史因而比一次模型输入保存更多材料。

结束事件也沿这些对象分层:turn_end 包含一次模型生成及其工具结果,因此本例先有读取调用的 turn,再有生成概括的 turn;agent_end 表示低层循环退出。外层 AgentSession 随后仍可能处理重试、压缩或排队输入,完成这些工作后另发 agent_settled。客户端要判断哪一层已经结束,需要消费对应层的事件。13

4.3 二次封装解决了什么,仍需保留什么

从 read_file 的往返看,框架把模型 API 变成了一组可以继续组合的对象。adapter 将原生 item 和参数事件转换成统一消息;聚合形成完整调用;执行器返回关联到调用的结果;循环将两者放进后续上下文。业务工具只需实现文件读取,上层编排则围绕这些对象组织任务。

03 / FRAMEWORK FLOW

框架中的消息流转

从当前上下文出发,观察完整模型消息、工具执行与结果回填怎样组成下一次生成。

Eino:消息怎样穿过执行图

从当前状态读入上下文,经过流聚合与工具节点,再回到下一次模型输入。

结果追加到 state.Messages → 模型节点再次读取同一状态

1 / state.Messages

输入

用户消息 / 既有历史

交给下一层

[]*AgenticMessage

模型节点读取当前消息列表,并从 ToolInfos 派生本次 model.WithTools 选项。工具声明与可执行实例分别持有。

查看固定源码

箭头表示本章选取的消息与对象交接路径;中间转换、异常分支和并发活动未全部绘制。这里展示固定源码中的职责关系,不运行真实 Agent。

交给下一层的内容Eino 中的表达Pi 中的表达本例中的用途
模型能力入口AgenticModel 的 Generate/Streamprovider stream提交当前上下文
结构化调用AgenticMessage.ContentBlocksAssistantMessage.content保存名称、参数和调用身份
完整模型消息StreamReader 聚合结果AssistantMessageEventStream.result()决定是否执行 read_file
工具结果与后续上下文工具节点及 state.MessagesToolResultMessage 及 currentContext.messages将 README 交给下一次生成
运行进度TypedAgentEventmessage_*、tool_execution_*、turn_*呈现生成和执行过程

转换时还需要保留原协议的材料。Eino 的 response metadata 与块 extension 保存部分 provider 信息;Pi 把模型来源、签名和工具身份放进统一消息,再在请求转换时处理它们。例如,更换模型后,reasoning 签名可能无法继续使用,工具 ID 也可能要同时改写调用与结果。统一类型让这些转换集中在 adapter,却仍需遵守目标协议的规则。

Pi 的 transformMessages 还会为历史中缺少结果的调用补一条 isError=true、文本为 No result provided 的消息。它告诉后续模型这次调用没有可用结果;它的来源是历史转换,和本例中 read_file 实际返回的文件文本有不同含义。读取会话或追踪执行时,应保留这一区别。

流与消息也需要明确的所有权。Pi 的 partial 是继续变化的共享对象,保存某一时刻的观察快照需要明确复制;Eino 的 StreamReader 要关闭,多个消费者要先分流。完整聚合会保留收到的消息片段,事件观察也可能积累尚未消费的内容。因此,读取流的接口还需要和消费速度、缓冲及最终消息的生命周期一起考虑:Eino Responses 入口的 Pipe 容量为 1,后续复制与聚合各有存储;Pi 的 EventStream.push 在没有等待者时进入没有容量上限的队列。流式消费与端到端有界内存,是两个需要分别判断的问题。14

框架已经把 LLM API 的调用与工具往返组织起来。下一步,要让这次工作归属一段持续存在的会话,由宿主提供工具和权限,让客户端观察进度、发起后续输入或中断。Codex app-server、Claude Code 和 DeepSeek Harness 正是在这些对象之上继续设计 Runtime 的接口与结构。

5. 三个 Agent Runtime 的结构设计

前面的样例中,模型提出 read_file("README.md") 调用,工具返回文件内容,模型再生成概括。把这条循环接到客户端,客户端还要能继续同一个会话、展示读取过程、取消正在进行的工作。有些工具由运行引擎提供,有些工具则需要调用宿主程序中的代码。

Agent Runtime 把循环放进会话,并为用户工作、工具活动和宿主交互提供接口。下面沿同一次 README 读取,分别看 Codex 如何标识工作,Claude Agent SDK 如何连接宿主与 Claude Code 引擎,DeepSeek Harness 如何把内部会话交给不同客户端。

5.1 Codex app-server:把用户工作组织为 Thread、Turn 与 Item

用户输入“读取 README.md,并用一句话概括项目用途”后,第一次模型调用产生文件读取请求,第二次模型调用产生概括。这两次生成都属于同一次用户工作。用户随后继续提问时,新工作还应接着使用前面的会话历史。

Codex app-server 用三个对象表达这层关系:

对象在这次往返中代表什么后续如何使用
Thread容纳这次输入和后续对话的会话。客户端用 threadId 向同一会话提交新输入。
Turn从“读取并概括”这条输入开始的工作。两次模型调用和工具执行都在这次工作中;客户端用 turnId 观察或控制它。
Item工作中的一个具体单元,例如用户消息、工具活动和最终 Agent 消息。客户端按 item 身份更新内容与状态,逐项呈现工作过程。

threadId、turnId 和 item ID 帮助客户端定位不同层次的对象。模型 function call 的 call_id 仍用于关联调用与工具结果。运行引擎连接这两组身份:模型对象参与下一次模型输入,Runtime 对象参与工作编排和客户端交互。相关类型分别定义在 codex-rs/app-server-protocol/src/protocol/v2/thread_data.rs 与 item.rs 中。

客户端通过 stdio 接入时,每行传输一个 JSON 对象。请求用 id 等待响应,运行事件用无 id 的通知持续发送。服务器也可以反向发起请求,等待客户端提供审批决策或工具结果。这里的 JSONL 规定消息的分帧方式;RPC envelope 规定每帧的交互身份,wire 中省略 jsonrpc 字段。

假定初始化和 Thread 创建已经完成,下面用短别名表示身份,选出这次工作的主要帧。中间的工具项与增量通知省略,内容和时间采用教学数据:

数据结构JSONL 记录 · 5.1 Codex app-server:把用户工作组织为 Thread、Turn 与 Item
array · 49 个子节点
$[5 项]
[0]{3 个字段}
id10
method"turn/start"
[1]{2 个字段}
id10
[2]{2 个字段}
method"turn/started"
[3]{2 个字段}
method"item/completed"
[4]{2 个字段}
method"turn/completed"

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
array
完整值
[
  {
    "id": 10,
    "method": "turn/start",
    "params": {
      "threadId": "th_demo",
      "input": [
        {
          "type": "text",
          "text": "读取 README.md,并用一句话概括项目用途。"
        }
      ]
    }
  },
  {
    "id": 10,
    "result": {
      "turn": {
        "id": "turn_demo",
        "items": [],
        "itemsView": "notLoaded",
        "status": "inProgress"
      }
    }
  },
  {
    "method": "turn/started",
    "params": {
      "threadId": "th_demo",
      "turn": {
        "id": "turn_demo",
        "items": [],
        "itemsView": "notLoaded",
        "status": "inProgress"
      }
    }
  },
  {
    "method": "item/completed",
    "params": {
      "threadId": "th_demo",
      "turnId": "turn_demo",
      "completedAtMs": 1791345605000,
      "item": {
        "type": "agentMessage",
        "id": "item_answer",
        "text": "这个项目演示如何提供 HTTP API。"
      }
    }
  },
  {
    "method": "turn/completed",
    "params": {
      "threadId": "th_demo",
      "turn": {
        "id": "turn_demo",
        "items": [
          {
            "type": "agentMessage",
            "id": "item_answer",
            "text": "这个项目演示如何提供 HTTP API。"
          }
        ],
        "itemsView": "summary",
        "status": "completed"
      }
    }
  }
]

第一帧中的 id:10 只关联这次 turn/start 请求。第二帧返回 turn_demo,状态还是 inProgress,客户端已经有了后续可以观察和控制的工作身份。turn/started 报告运行已经开始;item/completed 给出定稿的 item_answer;turn/completed 最后报告 Turn 的终态。这几个时点由不同消息表达,客户端应分别处理。

最后一帧再次出现 item_answer,因为 Turn payload 携带的是最终 Agent 消息摘要。显示层用同一 item 身份更新已有内容,完整的工作过程则由运行期间的 item 事件提供。源码中,codex-rs/app-server/src/request_processors/turn_processor.rs 的 turn_start_inner 构造初始 Turn;bespoke_event_handling.rs 的 emit_turn_completed_with_status 构造最终状态和摘要。

这个客户端协议接在前面的工具循环上。模型流中的 OutputItemDone 进入 codex-rs/core/src/stream_events_utils.rs 的 handle_output_item_done,ToolRouter 从完整模型 item 识别工具调用。核心记录该 item,再创建工具执行 future。工具结果转为 ResponseItem 写入 conversation history,codex-rs/core/src/session/turn.rs 中的循环从更新后的历史继续 sampling。读取请求、文件内容与概括因此能连续进入同一次 Turn,同时在客户端呈现为具体的运行活动。

文件读取也可以由宿主提供。若将本章的 read_file 注册为客户端动态工具,服务器用 item/tool/call 把调用交给客户端;客户端读取 README.md,返回内容与 success,运行引擎再把结果送回模型。这条动态工具路径在所分析实现中属于实验接口。另一种宿主交互是审批:item/commandExecution/requestApproval 要求客户端给出 decision,之后由引擎执行或拒绝命令。前一种响应带回工具结果,后一种响应决定引擎能否执行。

turnId 还让取消请求指向确定的工作。普通 turn/interrupt 会保存待答复的 RPC,再提交核心 Interrupt;等到核心发出 TurnAborted,app-server 才答复该请求并发布 interrupted 终态。启动期间空 turnId 的 interrupt 有单独的立即回应分支。相关处理位于 turn_processor.rs 与 bespoke_event_handling.rs。客户端应按各方法的确认条件推进状态,不能把 start 的回应时机直接套给 interrupt。

取消 Turn 后,已经创建的背景 terminal 仍有自己的生命周期。工作对象让客户端能表达“这次工作已经中断”;引擎还要由进程资源的所有者继续处理 terminal。会话、用户工作和执行资源各有身份,控制请求才能落到相应对象上。15

5.2 Claude Code / Agent SDK:宿主如何复用运行引擎

先从本地会话看 Claude Code 保存了什么。下面选取本机 Claude Code 的一段文件读取记录,保留与结构有关的字段,将身份替换为短别名,删除文件路径和工具返回正文。它来自实际会话,并非主例 README.md 的运行记录;中间的 attachment 记录没有列出:

数据结构JSONL 记录 · 5.2 Claude Code / Agent SDK:宿主如何复用运行引擎
array · 43 个子节点
$[3 项]
[0]{5 个字段}
type"assistant"
uuid"record-35"
parentUuid"record-34"
sessionId"session-A"
[1]{6 个字段}
type"user"
uuid"record-36"
parentUuid"record-35"
sessionId"session-A"
sourceToolAssistantUUID"record-35"
[2]{5 个字段}
type"assistant"
uuid"record-38"
parentUuid"record-37"
sessionId"session-A"

方向键浏览 · Enter 展开 / 折叠或解析字符串

所选字段
JSON path
$
类型
array
完整值
[
  {
    "type": "assistant",
    "uuid": "record-35",
    "parentUuid": "record-34",
    "sessionId": "session-A",
    "message": {
      "role": "assistant",
      "id": "message-A",
      "content": [
        {
          "type": "tool_use",
          "id": "tool-call-1",
          "name": "Read",
          "input": {
            "file_path": "[REDACTED_VALUE]"
          }
        }
      ]
    }
  },
  {
    "type": "user",
    "uuid": "record-36",
    "parentUuid": "record-35",
    "sessionId": "session-A",
    "message": {
      "role": "user",
      "content": [
        {
          "type": "tool_result",
          "tool_use_id": "tool-call-1",
          "content": "[REDACTED_TOOL_RESULT]"
        }
      ]
    },
    "sourceToolAssistantUUID": "record-35"
  },
  {
    "type": "assistant",
    "uuid": "record-38",
    "parentUuid": "record-37",
    "sessionId": "session-A",
    "message": {
      "role": "assistant",
      "id": "message-A",
      "content": [
        {
          "type": "tool_use",
          "id": "tool-call-2",
          "name": "Read",
          "input": {
            "file_path": "[REDACTED_VALUE]"
          }
        }
      ]
    }
  }
]

第一条记录里,message.content 保存一个 Read 调用。第二条记录把返回内容保存为 tool_result,通过 tool_use_id 对应原调用的 id。它的外层 type 和内层 role 都是 user,但内容来自工具执行;解析会话时,要继续读取内容块类型,才能区分工具结果和人工输入。

调用之外还有两组身份。sessionId 表示所属会话,uuid 标识这条日志记录,parentUuid 连接记录的父子关系;内层 message.id 则标识模型消息。第二条的 sourceToolAssistantUUID 指向产生调用的日志记录,与 tool_use_id 分别定位记录和调用。下面的箭头表示字段引用:

flowchart LR
    R35["日志 record-35"] -->|message.id| M["模型消息 message-A"]
    R38["日志 record-38"] -->|message.id| M
    R35 -->|tool_use.id| T["调用 tool-call-1"]
    R36["日志 record-36"] -->|tool_result.tool_use_id| T
    R36 -->|parentUuid / sourceToolAssistantUUID| R35

第三条又是 assistant,却仍属于 message-A。在原日志中,thinking 块和两个工具块分别落在不同记录里,共享同一 message.id;随后出现的独立模型消息才使用 message-B。因此,三条 assistant 记录不一定代表三次模型生成,工具结果后的一条 assistant 记录也不一定已经开始下一次生成。恢复内容需要同时保留记录链、消息身份和调用关联,不能只按行数或相邻角色判断。

磁盘会话还保存工作目录、分支、时间、CLI 版本以及 attachment 等记录。它们帮助 Runtime 维护会话,但不会全部作为聊天消息发送给模型。文件末尾只说明已经读到当前已写入的记录,不能据此判定一次用户工作已完成。16

本地 JSONL 展示的是会话如何保存;宿主接入运行引擎则使用另一层协议。下面再沿公开 Agent SDK 的实现,看输入、内容消息和控制请求如何进入宿主程序。

在 Eino 和 Pi 中,应用使用库来构造模型与工具的执行循环。Claude Agent SDK 让宿主程序连接已有的 Claude Code 运行引擎。以 Python SDK 为例,它启动 CLI 子进程,将 stdin/stdout 配置为 --input-format stream-json 和 --output-format stream-json --verbose,再通过 JSONL 交换输入、观察消息和控制请求。宿主复用引擎中的工具、权限和会话能力;Anthropic 模型 Client SDK 所封装的则是模型 API。

这层连接可从 src/claude_agent_sdk/_internal/transport/subprocess_cli.py 进入。随后,src/claude_agent_sdk/_internal/query.py 的 Query._read_messages 按消息类型分派:常规消息进入调用方消费的 stream,control_request 进入控制处理器,control_response 按 request_id 唤醒正在等待的请求。宿主可以一边接收模型输出,一边回答工具或审批请求。

在 README 往返中,AssistantMessage 中的 tool_use 表达模型提出的读取,工具结果可以通过 UserMessage 承载,后续 AssistantMessage 带回概括。初始化和压缩等运行变化通过 SystemMessage 表达。开启 partial messages 后,StreamEvent 还能提供原始 API 流事件,供客户端逐步显示内容。这些类型在 src/claude_agent_sdk/types.py 中定义。

假定宿主把 read_file 注册为 SDK MCP 工具。Claude Code 引擎需要调用它时,将内层 MCP 消息放进 subtype:"mcp_message" 的控制请求。Query 控制处理器从外层读取 request_id,按 subtype 进入权限、hook 或 MCP 分支,再按成功、取消或错误的结果处理控制回包。

query.py:完整的 Query._handle_control_request
async def _handle_control_request(self, request: SDKControlRequest) -> None:
"""Handle incoming control request from CLI."""
request_id = request["request_id"]
request_data = request["request"]
subtype = request_data["subtype"]
try:
response_data: dict[str, Any] = {}
if subtype == "can_use_tool":
permission_request: SDKControlPermissionRequest = request_data # type: ignore[assignment]
original_input = permission_request["input"]
# Handle tool permission request
if not self.can_use_tool:
raise Exception("canUseTool callback is not provided")
context = ToolPermissionContext(
signal=None, # TODO: Add abort signal support
suggestions=[
PermissionUpdate.from_dict(s)
for s in (
permission_request.get("permission_suggestions") or []
)
],
tool_use_id=permission_request.get("tool_use_id"),
agent_id=permission_request.get("agent_id"),
blocked_path=permission_request.get("blocked_path"),
decision_reason=permission_request.get("decision_reason"),
title=permission_request.get("title"),
display_name=permission_request.get("display_name"),
description=permission_request.get("description"),
)
response = await self.can_use_tool(
permission_request["tool_name"],
permission_request["input"],
context,
)
# Convert PermissionResult to expected dict format
if isinstance(response, PermissionResultAllow):
response_data = {
"behavior": "allow",
"updatedInput": (
response.updated_input
if response.updated_input is not None
else original_input
),
}
if response.updated_permissions is not None:
response_data["updatedPermissions"] = [
permission.to_dict()
for permission in response.updated_permissions
]
elif isinstance(response, PermissionResultDeny):
response_data = {"behavior": "deny", "message": response.message}
if response.interrupt:
response_data["interrupt"] = response.interrupt
else:
raise TypeError(
f"Tool permission callback must return PermissionResult (PermissionResultAllow or PermissionResultDeny), got {type(response)}"
)
elif subtype == "hook_callback":
hook_callback_request: SDKHookCallbackRequest = request_data # type: ignore[assignment]
# Handle hook callback
callback_id = hook_callback_request["callback_id"]
callback = self.hook_callbacks.get(callback_id)
if not callback:
raise Exception(f"No hook callback found for ID: {callback_id}")
hook_output = await callback(
request_data.get("input"),
request_data.get("tool_use_id"),
{"signal": None}, # TODO: Add abort signal support
)
# Convert Python-safe field names (async_, continue_) to CLI-expected names (async, continue)
response_data = _convert_hook_output_for_cli(hook_output)
elif subtype == "mcp_message":
# Handle SDK MCP request
server_name = request_data.get("server_name")
mcp_message = request_data.get("message")
if not server_name or not mcp_message:
raise Exception("Missing server_name or message for MCP request")
# Type narrowing - we've verified these are not None above
assert isinstance(server_name, str)
assert isinstance(mcp_message, dict)
mcp_response = await self._handle_sdk_mcp_request(
server_name, mcp_message
)
if mcp_response is None:
# JSON-RPC notifications get no reply, but the control
# request that carried one still expects an ack.
mcp_response = {"jsonrpc": "2.0", "result": {}}
response_data = {"mcp_response": mcp_response}
else:
raise Exception(f"Unsupported control request subtype: {subtype}")
# Send success response
success_response: SDKControlResponse = {
"type": "control_response",
"response": {
"subtype": "success",
"request_id": request_id,
"response": response_data,
},
}
await self.transport.write(json.dumps(success_response) + "\n")
except anyio.get_cancelled_exc_class():
# Request was cancelled via control_cancel_request; the CLI has
# already abandoned this request, so don't write a response.
raise
except Exception as e:
# Send error response
error_response: SDKControlResponse = {
"type": "control_response",
"response": {
"subtype": "error",
"request_id": request_id,
"error": str(e),
},
}
await self.transport.write(json.dumps(error_response) + "\n")

server_name 选中宿主注册的 MCP bridge,mcp_message 带着内层调用身份和参数。_handle_sdk_mcp_request 将它路由到对应 bridge,宿主中的 read_file 取得 README.md 内容。处理器把 MCP response 放入 response_data,再用原 request_id 构造外层 control_response,写回 CLI。引擎取得工具结果后继续调用模型,生成项目用途的概括。

这里同时完成两层交互。内层 MCP 使用 JSON-RPC;外层控制 envelope 使用 SDK 自定义的 type、request_id 和 subtype。内层 notification 可以没有回复,代码仍给承载它的外层控制请求生成 ack。因此,消费这一通道时要在各层自己的请求身份下判断完成。

若控制请求的 subtype 是 can_use_tool,SDK 调用宿主的权限 callback,回答是否允许执行;若是 hook_callback,SDK 按 callback ID 调用已注册的 hook。它们让宿主参与决策和运行扩展。SDK MCP bridge 则把工具实现放在宿主进程。配置中,tools 等字段决定模型可用的工具,allowed_tools 表达自动批准规则。这些选择决定了读取请求由谁实现、何时需要宿主介入。

输出侧也要保留对象关系。partial StreamEvent 供显示层累积增量,完整 AssistantMessage 供应用接受定稿内容。同一模型 message 的多个块可能分别产生 AssistantMessage,并共享 message ID;收到完整块时,应据此整理之前的增量显示,避免把同一内容重复追加。

当概括已经生成,ResultMessage 报告当前响应工作的结果,ClaudeSDKClient.receive_response() 会在第一个 ResultMessage 后返回。Query 还可能继续观察后台任务和 session_state_changed:result 结束的是当前 turn,run 是否收束由后续运行状态决定。调用方读取 is_error 等字段判断结果;关闭子进程则由 transport 的清理路径负责。一次回答有了结果、会话进入 idle 和进程退出,各有自己的观察点。

宿主程序因此有两项明确职责:处理自己提供的能力,并维护与 CLI 的连接。循环和内置工具由运行引擎承担,SDK 负责把消息、回调与控制请求接到宿主代码。17

5.3 DeepSeek Harness:从内部模型流与会话日志投影到 SDK / ACP

Codex 定义客户端可以控制的工作对象,Claude Agent SDK 连接宿主与 Claude Code 引擎。DeepSeek Harness 展示了另一层关系:同一次 README 读取先在内部会话中形成记录,再由 SDK 或 ACP adapter 交给不同客户端。

先沿模型与工具走一次。LlmRuntime 调用模型 adapter,BlockAssembler 将规范化的 block 流组装成模型消息,ReactLoopAgent 收到完整工具调用后执行 read_file。文件内容进入会话,Agent 再开启下一个模型 step,最终生成概括。这次工作包含两个模型 step 和一次工具执行。相关实现从 packages/llm/llm/src 的模型流进入 packages/core/agent-loop/src/agent.ts。

Session 为这条路径保存共同记录。log 中先后出现 turn、step、assistant message、tool/call、tool/result 等事件。第二次模型调用需要读取请求与文件内容,却不需要把每一条运行通知都当成对话文本。Session 因而通过 surface 操作构造模型可见历史:log 保留已经发生的事实,surface 表达哪些消息、按什么关系进入当前模型输入。packages/core/session/src/types.ts 定义事件,packages/core/session/src/index.ts 维护记录与 surface。

历史中的模型消息还可能保存 replayState,其中包含 response identity、签名等 adapter 续接材料。通用 blocks 用于表达内容,私有 envelope 用于后续向原 adapter 发请求。packages/llm/llm/src/index.ts 的 LlmRuntime.forAdapter 在构造请求时检查消息来源:只有拥有该历史路由的 adapter instance 能取得对应 replayState。切换 adapter 时,通用内容仍可进入请求,其他 adapter 的私有材料被过滤。这让内部会话同时保留可读内容与原模型协议所需的续接状态。

客户端观察从同一份 Session 记录导出。ACP bridge 收到已提交的 assistant/message 后,调用下面的函数,把消息中的 blocks 转成有序 updates。

updates.ts:完整的 assistantUpdates
export async function assistantUpdates(
ctx: Context,
session: Session,
event: SessionEvent<'assistant/message'>,
): Promise<SessionUpdate[]> {
const updates: SessionUpdate[] = []
for (const block of event.data.message.content) {
if (block.type === 'reasoning') {
if (block.text.length > 0) {
updates.push({
sessionUpdate: 'agent_thought_chunk',
messageId: event.data.message.id,
content: { type: 'text', text: block.text },
})
}
continue
}
const content = await assistantBlockToAcp(ctx, block)
if (content !== undefined) {
updates.push({
sessionUpdate: 'agent_message_chunk',
messageId: event.data.message.id,
content,
})
}
}
const usage = usageUpdate(ctx, session, event)
if (usage !== undefined) updates.push(usage)
return updates
}

输入 event 已经包含一条组装完成的模型消息。函数按 content 顺序遍历:reasoning 文本产生 agent_thought_chunk,普通内容经 assistantBlockToAcp 转换后产生 agent_message_chunk,这些内容 update 都带着原模型消息的 ID。最后,可取得的当前上下文占用以 usage update 追加。

当最终 text block 是“这个项目演示如何提供 HTTP API。”时,ACP 客户端收到相应的 agent_message_chunk。这个 chunk 来自已经提交的 block,内部此前的逐 token 流没有在这里逐 token 转发。packages/acp/acp/src/session.ts 的 onSessionEvent 将这些 updates 放入 outputTail,按序发送到客户端。

文件读取有自己的投影。updates.ts 的 toolCallUpdate 将 tool/call 的 callId 作为 ACP toolCallId,开始一个 in_progress 的工具活动;toolResultUpdate 读取结果消息的 toolCallId,把它更新为 completed 或 failed。客户端据此将读取请求和读取结果放在同一个工具对象中,模型消息 ID 则用来定位后续概括。

SDK 与 ACP 还选择了不同的响应时机。packages/sdk/server/src/server.ts 的 prompt 把输入交给 agent.followup(message) 后,返回 messageId。SDK 客户端随后通过 session.event 与 session.status 观察执行。ACP 的 session/prompt 先等待输入接纳;若输入已经入队,packages/acp/acp/src/session.ts 的 settleAfterQuiescence 继续等待 Agent idle 和 outputTail,再响应这次 prompt。SDK 返回输入身份,ACP 还等待该 prompt 的执行与输出交付。

交付时还会发生状态转换。packages/acp/acp/src/codec.ts 的 turnEndToStopReason 将内部 completed、aborted、blocked 映射为 ACP end_turn,多个内部原因收敛为同一个外部状态。显式客户端取消由 settlement 返回 cancelled;运行错误则在 settlement 中拒绝为 RPC 错误。只看外部 end_turn 无法区分上述内部原因,需要回到会话记录。类似地,ACP usage_update 表示当前上下文占用,模型调用的 input/output usage 仍有自己的记录。

这条内部记录路径也解释了“提交”的含义。Session.append() 把事件快照追加到内存 log,并同步发布给观察者;持久化插件另行进行 I/O。ACP 可以据已接受的事件发送客户端更新,但 turn/end 或 whenIdle() 本身没有提供磁盘 flush 的通用保证。需要等待持久化时,应使用存储层定义的 flush 边界。

DeepSeek Harness 的连接顺序由此完整了:模型流形成消息,工具执行形成结果,Session 保存事实并提供模型输入,外部 adapter 再按客户端协议发送观察消息和结算响应。新增客户端时,适配器需要处理的是身份、内容粒度和完成状态的转换。18

5.4 会话、宿主与客户端如何连接

同一次 README 读取经过三个实现后,已经能辨认各条路径。输入进入执行循环,完整调用交给工具,文件内容写入历史并送回模型;客户端观察整个用户工作的进展,宿主通过反向请求参与工具或审批。先切换实现,沿消息路径选择一个对象,查看它从哪一层接收什么、向下一层交付什么:

04 / RUNTIME BOUNDARIES

Runtime 的工作与消息流转

沿客户端输入、执行循环、宿主工具与结果事件,观察同一次工作跨越哪些边界。

Codex:一次输入落在哪些工作对象上

Thread 容纳会话,Turn 容纳这次用户工作,Item 呈现具体活动。

同一 Thread 接受后续 turn/start;本次 turn/completed 只结束这次工作

1 / turn/start

输入

threadId + 用户输入

交给下一层

初始 Turn

RPC id 关联 start 的回应。响应中 Turn 仍为 inProgress,客户端据 turnId 继续观察或控制。

查看固定源码

箭头表示本章选取的消息与对象交接路径;中间转换、异常分支和并发活动未全部绘制。这里展示固定源码中的职责关系,不运行真实 Agent。

客户端的 Thread、Turn、Item 与模型的 Response、输出 item、调用身份分属各自协议。切换下面的身份图,分别追踪 Codex 的两条身份轨、Pi 的组合 ID 回流,以及 Claude 日志里不同记录对同一消息和调用的引用:

IDENTITY / REFERENCE

沿着身份看清对象关系

点击对象或连接,查看关联字段与身份的作用范围。箭头表示包含、关联或转换;日志引用图保留样本列举顺序。

身份值沿用正文教学别名。左右属于不同协议,中间表示运行引擎与历史交接。

客户端 app-server 协议
模型 Responses 协议
所选关系 / 作用范围
运行引擎与历史交接

运行引擎处理模型结果、派发工具,并把工具输出写入历史,再继续 sampling。两轨有各自的对象身份;这条交接关系不保证 Thread、Turn、客户端 Item 与模型 Response、item 一一对应。

强调色标出当前关系涉及的字段;同名字段仍需按所在协议对象解释。

再把这些职责之间的关系放在一张图中:

flowchart LR
    P[模型 API] -->|原生输出| A[Provider adapter]
    A -->|请求| P
    A -->|结构化消息| L[执行循环]
    L -->|输入与续接材料| A
    L -->|完整工具调用| T[工具执行器]
    T -->|关联结果| L
    L -->|追加事实| S[会话与上下文]
    S -->|当前模型输入| L
    C[客户端] -->|开始与取消请求| L
    L -->|进度与运行结果| C
    L -->|审批或宿主工具请求| H[宿主扩展]
    H -->|决策或工具结果| L

方框表示职责,同一职责可以由引擎或宿主实现。选择集成方式时,可以沿着图中的箭头询问具体对象由谁维护:

要维护的关系Codex app-serverClaude CodeDeepSeek Harness
一次用户输入与内部执行Thread、Turn、Item 把会话、工作和具体活动连起来。CLI 执行循环,SDK 消费内容消息和运行状态。ReactLoopAgent 编排 steps,Session log 记录 turn 与执行事实。
工具调用与宿主代码内置工具由核心执行,动态工具可通过 item/tool/call 交给客户端。控制分派把权限 callback、hook 与 SDK MCP bridge 接入引擎。工具形成会话中的 call/result,外部 adapter 投影其状态。
会话记录与下一次模型输入工具结果写入 conversation history,核心继续 sampling。引擎使用会话上下文,SDK 观察消息与压缩等状态。surface 构造模型可见历史,LlmRuntime 过滤 adapter 私有 replayState。
内部状态与客户端完成消息turn/start 给初始对象,turn/completed 给终态与摘要。ResultMessage 给当前响应结果,后续运行状态决定 run 的收束。SDK 先返回 messageId,ACP 等 idle 与输出交付,再映射结束原因。

这些完成消息的作用范围随对象而定。Codex TurnStatus 包含 inProgress、completed、interrupted、failed;Claude ResultMessage 结束当前 turn,后台工作可以延续 run;DeepSeek Harness ACP 的 end_turn 已经压缩了内部原因。上一节 Pi 的 agent_settled 也属于自己的外层运行契约。将它们接入统一服务之前,先确定服务中的“工作”覆盖哪些活动,再转换各实现的状态。

会话延续还会带来存储和执行环境的问题。resume 或 fork 决定使用哪段历史,文件系统与执行环境的隔离需要另外安排。compaction 改变当前模型输入,长期保存的 transcript 可以继续存在。DeepSeek Harness 的 append 与 flush 则表明,客户端观察到事件、Agent 已 idle、记录已经落盘,也各有相应的完成条件。1917

到这里,read_file("README.md") 从模型响应进入了框架循环,又通过 Runtime 的工作对象、宿主回调和客户端事件成为一次可以观察与控制的工作。

参考资料

Footnotes

  1. Chat Completions 的请求、候选消息和停止原因。本章模型 API 官方文档读取于 2026-10-07;读取日期不冻结服务端部署或模型行为。协议样例为教学构造。 来源:Chat 请求与响应参考。 ↩ ↩2

  2. 应用执行函数并回填结果的往返。 来源:函数调用往返。 ↩

  3. SSE 的事件分帧与字段解析。 来源:SSE 格式与解析标准。 ↩

  4. 模型 API 的流式事件。 来源:Responses 流式指南;Chat 流式参考;Responses 流式参考。 ↩ ↩2 ↩3 ↩4

  5. Responses 的 items、函数结果与上下文续接。 来源:Messages 与 Items 对照;处理 function calls;手动管理 Responses 上下文;Conversation state。 ↩ ↩2

  6. Anthropic Messages 的 content blocks、工具结果顺序与停止原因。 来源:Messages 请求;多轮历史;工具回填与排序;Messages 流式过程;停止原因;工具执行分类。 ↩ ↩2

  7. Gemini GenerateContent 的 Content、Part、函数结果与签名。调研时 Google 的 GenerateContent 概念指南标记为 Legacy;本章只讨论该接口,不混入 Interactions API 的状态与事件。 来源:GenerateContent 结构参考;Schema 与 Type;FunctionCall 与 FunctionResponse;Gemini 函数往返;多轮对话;流式与 candidate;Thought signatures;GenerateContent 指南定位。 ↩

  8. 服务端工具、计量与 provider 续接材料。 来源:OpenAI 工具分类;Gemini code execution;Responses reasoning 回放。 ↩

  9. Responses 参数片段、定稿事件与流式函数调用。 来源:Responses streaming events;function-call streaming。 ↩

  10. 完整上下文与工具交接前的终态检查。本节 runTask 是伪代码,不是 Pi 的 API。 来源:Pi terminal checks。 ↩

  11. Eino 的 AgenticMessage、Responses 事件转换与块聚合。 来源:AgenticModel 声明;AgenticMessage 与块类型;流入口;delta 转换;块聚合;函数调用聚合;函数 item 定稿转换;事件分派。 ↩

  12. Eino 的状态写回、工具图与两类续接。 来源:状态 wrapper 的读取、聚合与写回;观察流的分支;工具节点及结果转换;结果写入状态;模型到工具的分支与回边;请求级工具声明;从状态派生请求工具选项;provider 自动续接;Runner 续接;TurnLoop 的运行连接。 ↩

  13. Pi 的模型事件、工具回填、会话历史与结束边界。 来源:ai 包;agent-core 包;coding-agent 包;消息类型;模型事件联合;最终消息提取;slot 与函数块创建;参数分支;函数 item 定稿;终态与未定稿工具检查;请求变换与消息消费;工具选择与回填;结果转换回 Responses;截断与错误分支;并行执行与有序回填;会话消息保存;历史树与上下文投影;低层结束路径;会话外层停定事件。 ↩

  14. 统一消息的续接材料、读者生命周期与事件缓冲。 来源:Eino response extensions;Pi 上下文变换与合成错误;Pi partial 契约;Eino 读者生命周期;完整流聚合;Eino Pipe 入口;Pi 事件队列。 ↩

  15. Codex app-server 的工作对象、工具派发、取消与客户端通知。 来源:对象定义;Item 类型;多次 sampling 的 turn 编排;协议说明;启动响应实现;最终摘要实现;工具派发入口;工具输出进入历史;审批处理;动态工具契约;interrupt 实现;回应时点;背景 terminal 契约。 ↩

  16. 本地 Claude Code 会话样本。取自本机 CLI 的实际会话并脱敏,仅保留结构;身份使用短别名,文件路径和工具正文已删除,中间 attachment 记录未列出。这与下面的 Python SDK 分别取证,不代表 SDK wire 格式,也不是 README 主例的运行记录。 ↩

  17. Claude Agent SDK 的 CLI transport、控制请求、消息和结果范围。本章未运行 SDK 内置的 CLI。Claude Code 内部循环未作为公开源码读取,相关叙述依据公开契约和可见 SDK 实现。 来源:官方 SDK 概述;子进程 transport;消息类型;read dispatcher;SDK 源码;bridge 路由与响应;工具与权限配置类型;官方权限契约;官方流输出契约;StreamEvent 类型;Query 结果处理;结果类型;receive_response。 ↩ ↩2

  18. DeepSeek Harness 的模型流、会话事实与 SDK / ACP 投影。 来源:项目仓库;内部流类型;Agent step;Session 事件定义;实时流入口;消息来源与私有材料;adapter 所有权;投影实现;会话事件消费;工具投影;SDK prompt;ACP settlement;ACP Prompt Turn 契约;结束原因映射;错误与取消结算;上下文占用计算。 ↩

  19. Runtime 的响应作用域、会话延续与持久化。 来源:Codex 状态与对象;Pi settled 出口;Codex fork 契约;Claude 会话契约;压缩契约;DeepSeek Harness append。 ↩