跳到主要内容

主管道:一次请求的一生(Stream)

30 秒导读: Assistant.Stream 是 Yao 智能体运行时的中央循环——一次用户请求进来,从校验、装缓冲、选模型、拼历史,到调 LLM、跑工具、收尾落库,全部走这一个函数。本章按执行顺序端到端走一遍这条控制流,重点讲清楚"错误 / 中断 / panic 时,状态和落库怎么保证不丢"。Hook、工具循环、委派、记忆、沙箱各自的内部机制留给专章,本章只把它们当管道上的调用点串起来。

阅读前置:建议先看 装载篇——知道一个 Assistant 是怎么从 DSL 变出来的,再看它被"跑起来"会更顺。


1. 这是什么(零基础也能懂)

一句话定义: Stream 是"跑一次智能体"的主函数——你给它输入消息,它负责把整条流水线走完,并通过流式(SSE)把结果一段段吐给前端。

它解决什么问题: 一次 AI 对话请求,背后要做的事远不止"调一次大模型"。要先确认你有没有权限、要把历史对话捞出来拼进去、要挑一个合适的模型连接器、模型可能要求调工具、工具可能报错要重试、整个过程随时可能被用户按"停止"、也可能中途 panic——而且无论怎么结束,这次对话的消息和(出错时的)断点都得可靠落库,好让前端刷新后还能看到、甚至能续跑。

Stream 就是把这一大堆"必须按顺序做、且必须善后"的事,收进一个函数里管起来。

一句话直觉:Stream 想成一条工厂流水线的总控。原料(用户消息)从一头进,依次经过若干工位(校验 → 缓冲 → 选模型 → 拼历史 → 沙箱 → Create Hook → LLM → 工具重试 → Next Hook → 收尾),从另一头出成品(响应)。而流水线最关键的不是某个工位,是总控挂在门口的一排"善后钩子"(defer):不管中途哪个工位炸了、还是有人拉了急停,出厂时这排钩子一定会把"这批货记账、关灯锁门"做掉。这排 defer 就是本章的暗线。

它在代码里长这样(真实入口):

// agent/assistant/agent.go:21 func (ast *Assistant) Stream
func (ast *Assistant) Stream(
ctx *context.Context,
inputMessages []context.Message,
options ...*context.Options,
) (*context.Response, error)

调用者给三样东西:请求上下文 ctx(贯穿全程的"这次请求的世界")、输入消息 inputMessages、可选的 options(连接器、历史长度、skip 开关等)。返回一个 *context.Responseerror


2. 顶层全景(管道大概怎么转)

2.1 一张图先看清主干

下面这张图从上到下是执行顺序;左边一列是"正常往下走",右边标注的是出错时统一的三步善后。看图要点:几乎每个工位出错都走同一套善后(盖失败状态 → 发错误 stream_end → return),真正兜底的是最外层那两个 defer

用户请求


① checkPermissions ─────────── 无 Authorized → 直接 return err(还没进 defer 区)


② Interrupt.SetHandler 注册"停止按钮"回调


③ 合并 Options / MergeMetadata


④ EnterStack ──── defer done() ┐ 建栈(root/child),挂"归位"善后
│ │
▼ │ ← 从这里起,任何 return 都会触发已挂的 defer
⑤ InitBuffer ─ defer Flush ────┤ 建缓冲;挂"落库 + 收 panic + 记账"善后(核心)
│ │
▼ │
⑥ initializeCapabilities 解析连接器 → 拿 Capabilities ──┐
│ sendAgentStreamStart │
▼ │ 出错统一:
⑦ WithHistory → fullMessages 捞历史 + 去重叠,拼全量消息 │ finalStatus=Failed
│ BufferUserInput ├─ sendStreamEndOnError
▼ │ return nil, err
⑧ Sandbox V2 / Workspace ─ defer cleanup │
│ │
▼ │
⑨ Create Hook(可选)──── 若 Delegate → 提前委派,直接 return ─┘


⑩ BuildRequest → autoSearch → executeLLMStream(或 Sandbox 执行)


⑪ 工具调用 + 重试循环(maxToolRetries=3)


⑫ Next Hook(可选)/ 无 hook 走 tool loop / 标准响应 ← 三选一收尾


⑬ 仅 root stack:sendAgentStreamEnd + CloseOutput


返回 finalResponse

2.2 主要工位一句话职责

阶段干什么关键符号 / 位置
① 权限没有授权信息就拒checkPermissions(permission.go:9)
② 中断注册"停止按钮"处理器Interrupt.SetHandlerhandleInterrupt(agent.go:762)
③ 选项取/建 Options,合并 metadataMergeMetadata(agent.go:60)
④ 栈建 root/child 栈,挂"归位"EnterStack(stack.go:200)
⑤ 缓冲建 buffer,挂总善后 deferInitBuffer / FlushBuffer(chat.go:98/195)
⑥ 连接器解析模型 + 拿能力位initializeCapabilitiesGetConnectorllm.ResolveConnector
⑦ 历史捞历史 + 去重叠 = 全量消息WithHistory(history.go:39)
⑧ 沙箱建沙箱 / 独立 workspaceinitSandboxV2 / initStandaloneWorkspace
⑨ Create Hook请求前 JS 钩子,可提前委派HookScript.Create(agent.go:225)
⑩ LLM组请求 + 自动检索 + 流式调用BuildRequest / executeLLMStream(llm.go:12)
⑪ 工具执行工具 + 失败重试executeToolCalls / buildToolRetryMessages
⑫ 收尾Next hook / tool loop / 标准processNextResponse / executeToolLoop / buildStandardResponse
⑬ 关流root 才发结束 + 关输出sendAgentStreamEnd / ctx.CloseOutput

2.3 主线走一遍(高层,不进代码)

一次成功请求:确认有权限 → 建好这次请求的栈和缓冲区 → 挑好模型 → 把历史对话和这次输入拼成完整消息 → (可选)Create Hook 看一眼要不要改选项或直接转给别的 agent → 组装请求调大模型,边调边把 token 流给前端 → 如果模型要求调工具就去调、调坏了最多重试 3 次 → (可选)Next Hook 决定还要不要继续、或就地收尾 → 如果是根请求,发送结束事件、关闭输出 → 返回。

贯穿全程的两个"角色":

  • finalStatus / finalError(agent.go:84-85): 两个局部变量,是这次请求"活着还是死了"的单一真相。默认 completed;任何一步出错就翻成 failed
  • defer FlushBuffer(agent.go:88): 挂在门口的总善后。函数从任何路径退出——正常 return、错误 return、panic——它都会拿着 finalStatus 去落库、记日志、恢复 logger 身份。

3. 核心机制(逐个拆,由浅入深)

3.1 门口那排 defer:为什么它是整章的骨架

要解决的小问题: 一次请求中途会在十几个地方 return err,还可能 panic。如果"落库、关流、状态记账"要在每个出口都手写一遍,必然漏。

思路: 把善后集中挂成 defer,由 Go 保证"无论从哪出去都执行"。Stream 里挂了三层,后进先出(LIFO)——所以退出时的实际执行顺序是从下往上。

进入顺序(挂 defer) 退出顺序(执行 defer,LIFO 反过来)
──────────────────── ──────────────────────────
EnterStack → done() 3. done() 栈归位/标完成
InitBuffer → FlushBuffer ↑ 先执行 2、再 1、最后 3?不——
Sandbox → cleanup 实际顺序:cleanup → FlushBuffer → done

准确说,挂载顺序是 done(④)→ FlushBuffer(⑤)→ sandboxCleanup(⑧),LIFO 退出时执行顺序为:沙箱清理 → 缓冲落库 → 栈归位

最关键的一层是 FlushBuffer 那个 defer,它自己还内含 panic 恢复:

// agent/assistant/agent.go:88 defer func() { ... }()
defer func() {
if r := recover(); r != nil { // ① 捕获 panic
finalStatus = context.ResumeStatusFailed
// ...把 r 转成 finalError...
defer panic(r) // ② 落库后再把 panic 抛回去
}
ast.FlushBuffer(ctx, finalStatus, finalError) // ③ 无论如何都落库
ctx.Logger.End(finalStatus == context.StepStatusCompleted, finalError)
ctx.Logger.RestoreAssistantID()
}()

三个巧妙点(务必理解):

  1. panic 不吞: recover 只是为了先把状态改成 failed、把 buffer 落库(否则崩溃就丢数据),然后用内层 defer panic(r) 原样重抛,不改变"崩了就崩了"的对外行为。
  2. 状态即真相: 落库读的是 finalStatus。前面任何错误路径都先写 finalStatus = ResumeStatusFailed 再 return,所以 FlushBuffer 拿到的一定是真实结局。
  3. 只有 root 真正落库: FlushBuffer 内部先判 ctx.Stack.IsRoot(),子栈(委派、A2A)直接返回——缓冲是共享的,只有根请求负责最终写库(chat.go:200-201)。

注意 ① 权限校验在 defer 挂载之前(agent.go:29):此时若失败,直接 return nil, err,不触发任何善后——因为还没建栈、没建缓冲,本来也没什么要善后的。这是唯一"裸退出"的阶段。

3.2 落库缓冲:为什么先攒着、最后一次性写

要解决的小问题: 一次请求会产生一堆消息(用户输入、助手回复、工具调用摘要)和多个执行步骤。若每条都实时写库,IO 频繁且难以在"出错时打断点续跑"。

思路: 用一个 ChatBuffer 在内存里攒——攒消息、攒每个阶段的 step 快照——退出时 FlushBuffer 一次性落。

每个阶段用 BeginStep / CompleteStep 括起来,形成可回放的步骤序列:

step 类型常量对应阶段
StepTypeHookCreateCreate Hook"hook_create"
StepTypeLLMLLM 调用(含重试)"llm"
StepTypeTool工具执行"tool"
StepTypeHookNextNext Hook"hook_next"

(依据:buffer.go:88-91)

resume(续跑)是这套设计的红利: FlushBuffer 里,成功不存 resume 步骤,只有出错/中断才存——正常跑完不需要断点:

// agent/assistant/chat.go:253 FlushBuffer 节选
// 3. Only save resume steps on error/interrupt (not on success)
if finalStatus != agentcontext.StepStatusCompleted {
steps := ast.convertBufferedSteps(ctx.Buffer.GetStepsForResume(finalStatus))
// ...chatStore.SaveResume(steps)...
}

BeginStep 每次开步前还会先 UpdateSpaceSnapshot 把请求级内存(ctx.Memory.Context)快照进 buffer(chat.go:179-182),这样断点续跑时能还原当时的临时状态。快照/续跑的细节属于 记忆与沙箱专章

3.3 连接器解析:一个 use::role 的四级降级

要解决的小问题: 助手 DSL 里写的模型可能是具体连接器(openai.gpt-4o),也可能是角色引用(use::light 表示"随便给我一个轻量模型"),还可能留空。运行时得把这三种都解析成一个真实连接器。

在管道里的位置:initializeCapabilities(agent.go:780)在发送 stream_start 之前就调用,因为输出适配器需要拿到模型的 Capabilities(能力位:支不支持视觉、工具等)才能正确转换首个事件。它内部走 GetConnectorllm.ResolveConnector

解析优先级(读这张图:从上往下,命中即停):

opts.Connector(Create hook 可能改过)
│ 空则用 ↓
ast.Connector(DSL 里写的,可能是 "use::light")

▼ ResolveConnector 内部:
├─ 有 "use::" 前缀 → 取出 role,connectorID 清空
├─ connectorID 非空 → 直接 select(显式连接器,最高优先)
└─ 走 role 解析(role 空则当 "default"):
① GetRoleBy(role, identity) 用户/团队级设置
② GetRole(role) 系统级该角色默认
③ GetRoleBy("default", id) 降级到 default 角色(用户/团队)
④ GetRole("default") 降级到 default 角色(系统)
⑤ 都没有 → 返回 error

▼ 回到 GetConnector 的 legacy fallback:
defaultConnector → 自动探测 findCapableConnector → 最终 "connector not specified"

(依据:ResolveConnector resolve.go:28-82;GetConnector legacy fallback 见 agent.go:662-676)

关键细节: use:: 前缀常量是 RolePrefix = "use::"(resolve.go:13)。role 解析同时带 identity(来自 ctx.Authorized),所以同一个 use::light,不同用户/团队可以指向不同的实际模型——这是多租户按角色路由的基础。连接器体系的更多内容不在本章。

3.4 历史组装:去掉"客户端重复带的"那段

要解决的小问题: 前端有时会把历史消息也塞进请求里。如果直接把"库里的历史 + 请求带的历史 + 新输入"拼起来,历史就重了。

思路: 从库里加载历史后,找出"历史的后缀与输入的前缀"最长重叠段,把输入里重叠的部分砍掉,再拼:

库里历史: [h1 h2 h3 h4]
本次输入: [h3 h4 u1 u2] ← h3 h4 是客户端重复带的
└──┬──┘
overlapIndex = 2,砍掉输入前 2 条
拼成 fullMessages: [h1 h2 h3 h4 u1 u2]
cleanInput 只留: [u1 u2] ← 只有它进 BufferUserInput

(依据:WithHistory history.go:39;findOverlapIndex history.go:368 从最长可能重叠往下试,命中即返回)

在管道里的位置与善后: ⑦ 出错时,WithHistory 内部不硬失败——加载历史报错只是 Logger.Warn 然后当"无历史"继续(history.go:62-71),尽量让对话能进行。真正 return 的只有更外层的错误。

拼好后立刻 BufferUserInput(historyResult.InputMessages)——注意传的是 cleanInput(去重后的),避免把重复输入也存进库。且 Skip.History 为真时跳过(agent.go:154)。历史加载/转换的更多规则(工具调用摘要、action 摘要)见 history.go,本章不展开。

3.5 Create Hook:管道上的第一个"可提前跳车"点

要解决的小问题: 有时不该盲目调大模型——可能想在请求前改改选项,甚至直接判断"这该交给另一个专门 agent 办",省掉一次 LLM 调用。

在管道里的位置: ⑨ 仅当 ast.HookScript != nil 时执行。用 BeginStep(StepTypeHookCreate) / CompleteStep 括起来。

两条出路:

  1. 正常返回 createResponse + 可能被改过的 opts,继续往 LLM 走。
  2. 提前委派:createResponse.Delegate != nil,直接调 handleDelegation 转给目标 agent,然后——注意收尾也要做全:
// agent/assistant/agent.go:262 Create hook 提前委派后的收尾
if ctx.Stack != nil && ctx.Stack.IsRoot() {
ast.sendAgentStreamEnd(ctx, streamHandler, streamStartTime, "completed", nil, nil)
ctx.CloseOutput() // root 才关流
}
return delegateResponse, nil // 跳过 LLM 和 Next hook

这里体现一个贯穿全章的规矩:任何提前 return 的分支,只要自己是 root,就得补做 sendAgentStreamEnd + CloseOutput,否则前端的流永远收不到结束标记。Hook 的 JS 桥、Delegate 的语义属于 Hooks 专章多智能体专章,本章只标出"它在管道第 ⑨ 步、能提前跳车"。

3.6 工具重试循环:可重试 vs 不可重试的分岔

要解决的小问题: 大模型返回的工具调用,参数可能写错(比如漏字段、类型不对)。这类错误其实可以"把错误结果喂回给模型、让它改一次再试";但如果是 MCP 服务内部炸了,重试没意义。

思路: 最多试 maxToolRetries = 3 次(agent.go:365)。每轮执行完看错误性质决定去留:

执行 executeToolCalls

├─ 全成功 ──────────────────► CompleteStep,break 退出循环

└─ 有错误 → 遍历看 result.IsRetryableError

├─ 没有一个可重试 ──► "non-retryable errors" 直接 return err(MCP 内部问题)

├─ 是最后一次尝试 ──► "failed after 3 attempts" return err

└─ 可重试且还有次数:
buildToolRetryMessages(把错误结果拼成对话)
→ executeLLMForToolRetry(让模型改)
→ 模型没返回工具调用 → return err(视为放弃)
→ 否则更新 currentResponse,进入下一轮

(依据:重试循环 agent.go:369-489;可重试判定 result.IsRetryableError agent.go:415)

buildToolRetryMessages 拼了什么(agent.go:803): 严格照 OpenAI 的工具应答格式——① 之前的全部消息 → ② 带 tool_calls 的 assistant 消息 → ③ 每个工具一条 role:tool 的结果消息(含错误)→ ④ 一条 system 消息解释"请修正重试"(取自 i18n assistant.agent.tool_retry_prompt)。模型据此看到"我上次这么调、结果报了这个错",从而修正。

沙箱模式跳过这段: 整个工具执行块的入口条件是 !ast.HasSandboxV2()(agent.go:363)——沙箱里 Claude CLI 自己在内部处理工具调用,主管道不插手。工具执行 executeToolCalls、MCP 落地属于 工具循环专章

3.7 三条收尾分支:请求怎么"结束"

要解决的小问题: 工具跑完(或压根没工具)之后,这次请求可能还没完——Next Hook 可能说"再来一轮",或者需要把工具结果喂回模型继续对话。得有一个清晰的三选一。

三选一(依据:agent.go:501-596):

if ast.HookScript != nil: ← 分支 A:有 Next Hook
HookScript.Next(...) → processNextResponse(...)
(hook 决定继续/委派/自定义响应,细节见专章)

else if 有工具结果 && 非沙箱 && tool loop 没被禁: ← 分支 B:无 hook,走工具循环
executeToolLoop(...)
└─ 失败则降级:buildLoopFallbackDelegate → handleDelegation
再失败 → buildStandardResponse(兜底)

else: ← 分支 C:标准响应
buildStandardResponse(...)
  • 分支 A 交给 Next Hook 自主决策;processNextResponse 处理其返回(next.go:11)。
  • 分支 B 是"没写 hook 但模型调了工具"的默认智能——把工具结果喂回模型继续,直到收敛;executeToolLoop 失败还有一层 __yao.loop_fallback 委派兜底,再不行才 buildStandardResponse
  • 分支 C 最简单:没工具、或沙箱模式、或 tool loop 被禁,直接包一个标准响应。

三分支的内部逻辑分属 Hooks 专章工具循环专章。本章只强调:它们都汇流到同一个出口——⑬ 的 root 收尾。

3.8 root stack 才关流:嵌套调用为什么"不许关灯"

要解决的小问题: 一次请求里,主 agent 可能委派给子 agent、子 agent 又调 MCP……这些嵌套调用都走同一个 Stream 函数。如果每一层都发"流结束"、都关闭输出,前端会收到好几个结束标记、流被提前关死。

思路:ctx.Stack.IsRoot()(root 即 ParentID == "",stack.go:115)守门——只有最外层那次调用负责发 sendAgentStreamEndCloseOutput:

// agent/assistant/agent.go:604 正常收尾
if ctx.Stack != nil && ctx.Stack.IsRoot() {
ast.sendAgentStreamEnd(ctx, streamHandler, startTime, "completed", nil, completionResponse)
ctx.CloseOutput() // 发 [DONE],关输出写入器
} else {
// 嵌套调用:只记 debug 日志,不关流
}

同样地,sendAgentStreamStart(agent.go:705)也判 IsRoot()——保证一次 agent 执行只有一个 stream_start / stream_end,哪怕内部调了好几次 LLM。这条"root 才关流"的规矩在前面每个提前 return 分支里都被反复遵守(§3.5 就是一例)。


4. 巧妙之处(可借鉴)

  • 状态单一真相 + defer 落库:finalStatus/finalError 两个变量当"结局",所有错误路径先写它再 return,善后统一读它——把"十几个出口的记账"收敛成一处(agent.go:84-108)。
  • panic 先落库再重抛: recover 不是为了吞掉崩溃,而是抢在崩溃传播前把 buffer 存下来,再 defer panic(r) 原样抛回(agent.go:90-100)。数据不丢,行为不变。
  • 成功不存 resume,出错才存: 断点续跑数据只在失败/中断时落,正常路径零额外写入(chat.go:253)。
  • root 守门: 一个 IsRoot() 判断,让"同一个 Stream 函数"既能当主入口、又能当被嵌套的子调用,而不重复发流事件、不重复落库(agent.go:604、chat.go:200、sendAgentStreamStart)。
  • 历史去重叠: 用"历史后缀 ∩ 输入前缀"的最长匹配砍掉客户端重复带的历史,避免上下文里出现重复对话(history.go:368)。
  • 连接器四级降级 + identity: use::role 让 DSL 只声明"要什么档次的模型",实际选哪个交给运行时按用户/团队解析,天然支持多租户(resolve.go:28)。

5. 边界与局限(诚实)

  • 权限校验极简: checkPermissions 当前只校验 ctx.Authorized != nil(permission.go:9-14),没有更细的能力/配额检查——细粒度授权不在这一层。
  • 工具重试上限硬编码为 3: maxToolRetries := 3 写死在函数里(agent.go:365),不可配置。
  • 非 root 不落库: 子栈(委派/A2A)共享 buffer 但不 flush;若你期望子 agent 独立持久化,这里不会发生——落库是 root 的专属职责(chat.go:200-201)。
  • 沙箱模式绕过主工具循环: HasSandboxV2() 为真时,§3.6 工具重试与部分收尾分支被跳过,由 Claude CLI 内部接管(agent.go:363、317)——主管道对那段是"黑盒"。
  • 历史加载失败静默降级: 加载历史出错只 Warn 不中断(history.go:62),对话会在"无历史"下继续——好处是鲁棒,代价是用户可能察觉不到上下文丢了。

6. 横向对比

同组其它章各讲管道上的一个"调用点",本章是把它们串起来的骨架:

你想深入的点去哪章
一个 Assistant 怎么从 DSL 装载出来01-loading.md
Create/Next Hook 的 JS 桥与边界注入03-hooks-jsapi.md
工具怎么落到真实 MCP、执行细节04-toolloop-mcp.md
Stack、委派、A2A 的多智能体编排05-multiagent-stack.md
四层记忆、快照续跑、沙箱执行器06-memory-sandbox.md

总览(笼子而非动物的视角)见 index.md


7. 代码地图(导航索引)

用符号名 grep 定位,比行号抗漂移。

主题文件路径符号名
中央循环主函数agent/assistant/agent.go(*Assistant).Stream
权限校验(裸退出)agent/assistant/permission.gocheckPermissions
中断处理器agent/assistant/agent.gohandleInterrupt
建栈 + 归位 deferagent/context/stack.goEnterStack / (*Stack).IsRoot
缓冲初始化agent/assistant/chat.goInitBuffer
步骤括号agent/assistant/chat.goBeginStep / CompleteStep
落库总善后agent/assistant/chat.goFlushBuffer
步骤类型常量agent/context/buffer.goStepTypeLLM / StepTypeTool / StepTypeHookCreate / StepTypeHookNext
连接器解析入口agent/assistant/agent.go(*Assistant).GetConnector
角色/连接器四级降级agent/llm/resolve.goResolveConnector / RolePrefix
能力位初始化agent/assistant/agent.goinitializeCapabilities
历史组装 + 去重叠agent/assistant/history.goWithHistory / findOverlapIndex
缓存用户输入agent/assistant/chat.goBufferUserInput
沙箱 / 独立工作区agent/assistant/sandbox_v2.goHasSandboxV2 / initStandaloneWorkspace
Create Hook 调用点agent/assistant/agent.goHookScript.Create(Stream 内)
组请求 / 自动检索 / LLM 流agent/assistant/build.go, search.go, llm.goBuildRequest / shouldAutoSearch / executeLLMStream
工具执行agent/assistant/mcp.goexecuteToolCalls
工具重试消息拼装agent/assistant/agent.gobuildToolRetryMessages
收尾三分支agent/assistant/next.go, loop.goprocessNextResponse / executeToolLoop / buildStandardResponse
root 关流agent/assistant/agent.gosendAgentStreamStart / sendAgentStreamEnd