跳到主要内容

数据截至 (上游 commit 5053c08115bd)

一次对话的生命周期:三层循环与流式协议

30 秒导读: 用户在 Onyx 聊天框按下回车之后,后端不是"调一次模型然后返回",而是跑起三层嵌套的循环:最外层做一次性装配和线程编排,中间层反复"让模型选工具 → 跑工具 → 再问模型",最内层把一次 HTTP 流式响应拆成 reasoning / answer / tool_call 三股。这三层之间不靠函数返回值传数据,而是靠一个共享队列 + 一种叫 Packet 的信封。本章讲清这条骨架。

本章只讲骨架。上下文里每条消息放在哪、为什么这么排 → 02;工具怎么定义、怎么并行、模型不听话怎么兜底 → 03


1. 先建立直觉:按下回车之后

1.1 用户看到的

前端发一个 POST 到 /chat/send-chat-message,拿回一条 SSE(Server-Sent Events,服务端持续推送的长连接)流。屏幕上会依次出现:

  1. 一个"思考中"的折叠块(模型的推理内容);
  2. 一个"正在搜索……"的工具卡片,里面滚出查询词和命中文档;
  3. 正式答案一个字一个字往外冒,带蓝色的引用角标;
  4. 流结束。

依据:router prefix /chatbackend/onyx/server/query_and_chat/chat_backend.py:155,路由声明在 :538-539,处理函数 handle_send_chat_message:562

1.2 后端实际做的

同一段时间里,后端做的事远不止"转发模型输出":

阶段干了什么
装配查会话、跑权限校验、建 LLM 实例、拉历史、算 token 预算、预留助手消息 ID
循环最多 6 轮"问模型 → 模型要调工具 → 跑工具 → 把结果塞回历史 → 再问"
解码每轮里把 provider 的 chunk 流拆成推理、答案、工具调用参数三股,边拆边推给前端
收尾落库(消息、工具调用、引用文档)、清 Redis 上的"处理中"栅栏、必要时压缩历史

1.3 一句话类比

把它想成餐厅出餐:装配层是接单和备料(一次性),循环层是厨师"尝一口 → 加料 → 再尝"(可能来回几轮),解码层是传菜员边做边端(不等整道菜做完)。而 Emitter 就是那条传菜通道——厨房不用管客人还在不在座位上,菜照做照往通道上放。


2. 顶层全景:三个文件、三层循环

Onyx 把这条链路切成三个文件,职责严格分层,上层不碰下层的细节

┌──────────────────────────────────────────────┐
HTTP 请求 ───▶ │ 第 1 层:编排 │
│ process_message.py │
│ · build_chat_turn 一次性装配 │
│ · _run_models 开线程池、开写线程 │
└───────────────┬──────────────────────────────┘
│ 每个模型一条 worker 线程

┌──────────────────────────────────────────────┐
│ 第 2 层:多轮工具循环(agent 主循环) │
│ llm_loop.py · run_llm_loop │
│ 最多 MAX_LLM_CYCLES 轮: │
│ 问模型 → 有工具调用?→ 跑工具 → 结果进历史 → 再问│
└───────────────┬──────────────────────────────┘
│ 每轮一次

┌──────────────────────────────────────────────┐
│ 第 3 层:单次流式解码 │
│ llm_step.py · run_llm_step │
│ 把 provider 的 chunk 流拆成三股: │
│ reasoning / answer / tool_call 参数 │
└───────────────┬──────────────────────────────┘
│ 每个片段一个 Packet

Emitter ──▶ merged_queue ──▶ 写线程 ──▶ SSE

怎么读这张图: 从上往下是调用方向,从下往上(最后那条 Emitter 线)是数据回流方向。关键在于:数据不沿着调用栈往回 yield,而是横穿到一个共享队列上。这条设计是理解全篇的钥匙,§5 展开。

2.1 部件一句话职责

部件干什么在哪
build_chat_turn校验 + 装配,产出不可变的 ChatTurnSetupbackend/onyx/chat/process_message.py:593
_run_models开 N 条 worker 线程 + 1 条写线程,返回读出生成器process_message.py:1143
_stream_chat_turn串起装配和执行,统一处理各类异常process_message.py:1636
handle_stream_message_objects单模型公开入口(薄壳)process_message.py:1881
gather_stream非流式场景:把整条流吃干,聚合成一个响应对象process_message.py:2136
run_llm_loopagent 主循环,控制 tool_choice 状态机backend/onyx/chat/llm_loop.py:743
run_llm_step消费一次流式解码的所有 Packetemit 出去backend/onyx/chat/llm_step.py:1575
run_llm_step_pkt_generator真正的解码器,三股流分离llm_step.py:1074
EmitterPacket 打上 model_index 并投进共享队列backend/onyx/chat/emitter.py:8
Placement四元组坐标,告诉前端这个包渲染到哪个块backend/onyx/server/query_and_chat/placement.py:4

3. 第 1 层:装配与编排

3.1 build_chat_turn —— 一个"会 yield 的构造函数"

build_chat_turnprocess_message.py:593)的签名有点特别:

def build_chat_turn(...) -> Generator[AnswerStreamPart, None, ChatTurnSetup]:

既是生成器又有返回值。调用方这样用(process_message.py:1735-1741):一边 next() 取出它中途 yield 的包并转发给前端,一边等 StopIteration.value 拿到装配结果。

为什么要这样?因为装配阶段就有前端必须立刻知道的两条信息:

  • 新建会话时的 CreateChatSessionIDprocess_message.py:643);
  • 预留的助手消息 ID MessageResponseIDInfoprocess_message.py:970)——前端要靠它把后续 token 挂到正确的气泡上,而这条消息在模型还没开口之前就已经在数据库里占好坑了(reserve_message_id)。

装配阶段做的事按顺序是:会话解析 → 遥测 → 建 LLM 实例并查成本额度 → 校验用户文件归属 → 重建线性历史 → 定位父消息(支持重新生成、分支)→ 跑 Query Processing 钩子并落库用户消息 → 收集可读文件 ID → 摘要截断 → token 预算 → 决定检索参数 → 预留助手消息 ID → 转成不带 ORM 对象的 simple_chat_history → 重置停止信号、点亮"处理中"栅栏 → commit

3.2 ChatTurnSetup 的"脱离态安全"契约

装配结果是一个 frozen dataclass(backend/onyx/chat/chat_state.py:195)。它的 docstring 里写死了一条契约:里面所有 ORM 对象在 build_chat_turn 返回后都已 detached,下游只能读装配期 eager-load 过的列,碰任何懒加载关系都会炸 DetachedInstanceError

这不是洁癖,是被逼的:LLM 流可能跑几分钟,期间绝不能占着一条数据库连接。所以装配用完 expunge_all() 就把 session 关掉(process_message.py:1743),后面谁要写库谁自己开短连接。

3.3 ChatStateContainer:只累积状态,不承载逻辑

ChatTurnSetup(不可变的输入)对称的,是 ChatStateContainerchat_state.py:34)——可变的输出累加器。

它的设计约定只有一句:只累积、不决策。所有写方法都是 add_* / set_*,所有读方法都是 get_* 且返回拷贝,全部由一把 threading.Lock 保护(chat_state.py:46)。它不知道什么时候该落库,也不知道什么算"完成"。

好处在中断场景立刻兑现:用户点停止时,另一条线程可以随时对它做一次线程安全快照,把已经生成到一半的答案存下来(§7.2)。如果状态散落在循环的局部变量里,这件事做不到。

3.4 出口有三个

入口场景位置
handle_stream_message_objects单模型 SSE(UI 主路径、Slack bot、evals)process_message.py:1881
handle_multi_model_stream2–3 个模型并排对比,仅支持流式process_message.py:1921
gather_stream / gather_stream_full非流式 API:把整条流吃完再聚合process_message.py:2136 / :2001

gather_stream 值得单独说一句:它证明了流式协议是自足的——不需要额外的"非流式代码路径",只要把同一条包流按类型累加(AgentResponseDelta 拼答案、CitationInfo 收引用、AgentResponseStart 取最终文档),就能还原出完整响应。真正的非流式增强版 gather_stream_full 还会顺手读 ChatStateContainer 补上工具调用明细。


4. 第 2 层:run_llm_loop 主循环骨架

4.1 它要解决的小问题

模型一次开口只能说三种话:说答案、说"我要调工具"、说推理。如果它说要调工具,那这次对话还没结束——得真去跑工具,把结果喂回去再问一遍。循环的存在就是为了这个。

问题是:怎么保证它一定会停?

4.2 用 tool_choice 强行收敛

循环体是一个定长 forllm_loop.py:882):

for llm_cycle_count in range(MAX_LLM_CYCLES):
out_of_cycles = llm_cycle_count == MAX_LLM_CYCLES - 1

MAX_LLM_CYCLES 默认 6(backend/onyx/configs/chat_configs.py:12,可用同名环境变量覆盖)。源码注释里解释了为什么是 6:够跑完 web_search → open_url → web_search → open_url → open_url → 收尾 这条典型链(llm_loop.py:304-312)。

真正的收敛靠每轮开头那个三分支状态机(llm_loop.py:885-898):

条件tool_choice给模型的工具清单效果
forced_tool_id 非空REQUIRED只留那一个工具强制模型必须调它
out_of_cyclesran_image_genNONE[]断了模型的手,只能出答案
其它AUTO全部工具模型自己决定

两个细节容易漏:

  • 强制只强制一次。 命中 forced_tool_id 分支后立刻 forced_tool_id = Nonellm_loop.py:891),所以"强制搜索"只作用于第一轮,之后放回 AUTO
  • 最后一轮不是"劝"而是"断"。 out_of_cycles 时工具清单直接置空(llm_loop.py:895),不是靠 prompt 提醒模型该收尾了。同一个变量还会传给 select_reminder_textllm_loop.py:1003)去调整提醒语,但那只是锦上添花——真正的硬保证是空工具表。

4.3 一轮里发生什么

每轮:
① 定 tool_choice / final_tools llm_loop.py:741-754
② 拼 system prompt + 提醒消息 llm_loop.py:759-839 ─→ 详见 02 章
③ construct_message_history 截断 llm_loop.py:841
④ run_llm_step 流式解码 ★ llm_loop.py:860 ─→ 第 3 层
⑤ 模型不听话?文本里捞工具调用 llm_loop.py:881 ─→ 详见 03 章
⑥ run_tool_calls 并行跑工具 llm_loop.py:930 ─→ 详见 03 章
⑦ 结果写回 simple_chat_history llm_loop.py:1106-1158
⑧ 没有工具调用 → break llm_loop.py:1161

出口条件只有两个:这轮模型没要求调工具break),或者轮次跑满。跑满时因为第 ② 步已经把工具断掉了,模型只能吐答案,所以 break 不会漏。

4.4 循环结束后的三种结局

结局处理位置
既无答案也无工具调用EmptyLLMResponseError(OpenAI 空流会特判为额度耗尽)llm_loop.py:1389-1394:98
有工具调用但最终没答案RuntimeError,提示多半是工具输出畸形或 provider 配错llm_loop.py:1396-1401
正常emit 一个 OverallStop 收尾llm_loop.py:1403-1411

注意 run_llm_loop 的返回类型是 None——它不返回答案。答案早就通过 Emitter 流出去了,同时也一路写进了 ChatStateContainer。这是"数据横穿队列"设计的直接后果。


5. 第 3 层:run_llm_step 把一股流拆成三股

5.1 它要解决的小问题

provider 回来的是一串 chunk,每个 chunk 的 delta 里可能有 content(答案文本)、reasoning_content(推理)、tool_calls(工具调用参数的片段)。三种东西混在同一条流里,而前端要把它们渲染成三种完全不同的 UI 块。

5.2 原理演示

# 示意,非源码
for chunk in llm.stream(prompt, tools, tool_choice):
d = chunk.choice.delta
if d.reasoning_content: # 第一股:推理 → 折叠块
emit(ReasoningStart()) if first else None
emit(ReasoningDelta(reasoning=d.reasoning_content))
if d.content: # 第二股:答案 → 正文
close_reasoning_block() # 推理块必须先闭合
emit(AgentResponseDelta(content=d.content))
if d.tool_calls: # 第三股:参数碎片 → 工具卡片
close_reasoning_block()
accumulate(d.tool_calls) # 攒够了流末尾统一产出

重点看两处:推理块在答案或工具调用出现时必须先闭合;工具调用参数是攒到流结束才统一产出的,因为 JSON 参数在流中是逐字符切碎的。

5.3 真实实现

真身是 run_llm_step_pkt_generatorllm_step.py:1074),一个 Generator[Packet, None, tuple[LlmStepResult, bool]]。三股分流分别在 llm_step.py:1373(reasoning)、:1347(content)、:1355(tool_calls)。

工具调用在流结束后由 _extract_tool_call_kickoffsllm_step.py:1435、定义在 :374)统一产出,顺便按出现顺序给每个调用分配递增的 tab_index——这就是前端能把并行工具渲染成多个 tab 的原因。

run_llm_stepllm_step.py:1575)只是个消费壳,十行不到:

while True:
try:
packet = next(step_generator)
emitter.emit(packet)
except StopIteration as e:
llm_step_result, has_reasoned = e.value
return llm_step_result, has_reasoned

生成器负责"产出什么",壳负责"往哪儿送"。 这个切分让同一个解码器可以被 deep research 那条路径复用(换一个 emitter 和几个开关就行)。

5.4 推理会"吃掉"一个 turn_index

has_reasoned 这个返回值容易被忽略,但它是坐标系能对齐的关键。

推理块闭合时,_close_reasoning_if_active 会调 _increment_turnsturn_index 往前推一格(llm_step.py:1242:1024)——因为在前端眼里,推理块和后面的答案块是两个独立的渲染块

于是外层循环必须知道"这一步多花了一格坐标",所以它维护 reasoning_cycles,每次传给下一步的坐标是 Placement(turn_index=llm_cycle_count + reasoning_cycles)llm_loop.py:1051:876-877)。少了这个补偿,第二轮的包就会覆盖第一轮的推理块。


6. 流式协议:Emitter、Placement、Packet

6.1 Emitter:为什么不能一路 yield 回栈顶

最自然的写法是让 run_llm_step yield、run_llm_loop yield from_run_modelsyield from——数据沿调用栈冒泡。Onyx 没这么做。

原因很实际:工具在深处也要发包(搜索工具要推查询词和文档、Python 工具要推 stdout),而工具是被 run_tool_calls 在线程池里调起来的,根本不在生成器链上。要让它们也能 yield,整条链上每个函数都得改成生成器,工具接口也得变成生成器接口。

Emitteremitter.py:8)把这件事拍平:它只有一个方法。

def emit(self, packet: Packet) -> None:
if self._drain_done is not None and self._drain_done.is_set():
return
base = packet.placement or Placement(turn_index=0)
tagged = Packet(placement=base.model_copy(update={"model_index": self._model_idx}), obj=packet.obj)
self._merged_queue.put((self._model_idx, tagged))

三件事:查一下"还要不要发"、盖上 model_index 戳、投进共享队列(emitter.py:32-40)。谁拿到 emitter 谁就能发包,不管在栈的哪一层、哪条线程。

还有个 NullEmitteremitter.py:43):emit 是空实现。给那些在聊天流之外跑工具的调用方用——Search API、MCP server——它们要工具的返回值,不要流。

6.2 Placement:四元组坐标

包发出去了,前端怎么知道该渲染到哪儿?靠 Placementplacement.py:4)这四个字段:

字段含义前端拿它干什么
turn_index本条消息内的第几个迭代块,单调递增决定块的先后顺序
tab_index同一轮内并行工具调用的序号每个工具渲染成自己的 tab
sub_turn_index嵌套层级,None 表示顶层工具套工具(子 agent)时的缩进
model_index哪个模型产的,None 表示 LLM 未启动前的装配包多模型对比时的列

前端有一份镜像类型定义(web/src/app/app/services/streamingModels.ts:478-481),字段一一对应。

坐标的三个来源正好对上三层:turn_indexrun_llm_loop 按轮次给(外加推理补偿),tab_index_extract_tool_call_kickoffs 按工具顺序给,model_indexEmitter 统一盖戳。没有任何一层需要知道另外两层怎么编号。

03 讲子 agent 时只用到前三维(子 agent 靠 sub_turn_index 把自己的包嵌进父 tab),model_index 那一维只在本章的多模型对比里出现——同一个类,四维始终都在。

6.3 Packet:一个信封 + 一个判别联合

class Packet(BaseModel):
placement: Placement
obj: Annotated[PacketObj, Field(discriminator="type")]

backend/onyx/server/query_and_chat/streaming_models.py:490

PacketObj 是一个 40 多个成员的判别联合(streaming_models.py:435-488),靠 type 字段区分。所有 type 字符串集中定义在 StreamingType 枚举里(streaming_models.py:14)——一处真源,Pydantic 靠它做判别式反序列化,前端靠它做 switch。

按用途分四类:

类别代表成员说明
控制SectionEndOverallStopTopLevelBranchingChatHeartbeat不带内容,只标记流的结构
答案AgentResponseStart / AgentResponseDelta / CitationInfo正文与引用
推理ReasoningStart / ReasoningDelta / ReasoningDone折叠块
工具SearchToolStartPythonToolDeltaOpenUrlDocuments ……每种工具一组

两个控制包值得单独讲:

  • SectionEndstreaming_models.py:76):一个工具跑完就发一个,成功失败都发backend/onyx/tools/tool_runner.py:219-225)。前端靠它把工具卡片从"转圈"切到"完成"。源码注释诚实地说了这个包"严格来说不必要,将来可以删"。
  • OverallStopstreaming_models.py:80):整条流的终点。正常结束由 run_llm_loop 发(llm_loop.py:1403);用户点停止时由写线程发,带上 stop_reason="user_cancelled"process_message.py:1514-1521)。两条路径产出同一种终止包,前端只有一个收尾分支。

另外 TopLevelBranchingllm_loop.py:1099-1107)是个提前量:一旦发现这轮有多个工具调用,先发一个"接下来会分 N 支"的通知,免得前端先渲染第一个工具再重排。


7. 并发与生命周期

7.1 _run_models 的线程拓扑

_run_modelsprocess_message.py:1143)是整章工程含量最高的地方。它开出的东西:

merged_queue (无界)

worker[0] ──emit──────▶│
worker[1] ──emit──────▶│──▶ 写线程 _drain_to_completion
worker[2] ──emit──────▶│ │
├──▶ StreamBufferWriter (落 Redis,可重放)
└──▶ tee ──▶ 读出生成器 ──▶ SSE ──▶ 浏览器

怎么读这张图: 左边 N 条 worker 各跑一个完整的 run_llm_loop;中间写线程是唯一消费者;右边分两路——一路进可重放缓冲(跨 pod 恢复用),一路进 tee 队列给读出协程。

关键在于:写线程和读出协程是解耦的。读出协程可以随时死掉(客户端断线),写线程照样跑到底。docstring 里那句话说得很直白:"写线程是这次运行的生命线,永远跑到完成;读出端可以随意死掉"(process_message.py:1218-1219)。

N=1 和 N>1 走的是同一条路径process_message.py:1159-1161)。单模型不是特例,只是 n_models == 1。这消掉了一整套"单模型专用代码",代价是单模型也要过一次线程池——对一个动辄几十秒的 LLM 调用来说完全可以忽略。

每条 worker 线程都拿 contextvars.copy_context() 的独立副本(process_message.py:1594),因为一个 Context 对象不能被多线程同时 enter。tenant id 这类上下文变量就是靠这个跨线程传的。

7.2 三个协调原语各管什么

原语类型谁设置效果
drain_donethreading.Event停止按钮路径(:1392)或写线程崩溃(:1446Emitter.emit 直接 return,worker 剩余输出全丢;客户端断线不设它
reader_gonethreading.Event读出生成器的 finally:1500_publish 不再往 tee 塞东西,但写线程继续跑、继续写缓冲
teequeue.Queue写线程 → 读出协程的交接队列;_STREAM_DONE 哨兵表示流真的完了

这三者的区别是本章最需要记住的一点:

  • 客户端断线 → 只设 reader_gone。生成继续,答案照样落库,用户刷新页面还能看到完整回答。
  • 用户点停止 → 设 drain_done。worker 是不可中断的(阻塞在 provider 的 HTTP 流上),所以做法不是"叫停",而是把它们的出口堵死,然后立刻拿已有的部分状态落库。

7.3 worker 自完成:不等队友

_run_modelfinally 块里做两件事(process_message.py:1420-1422):先 _persist_model_outcome(model_idx, _PersistContext.WORKER),再往队列投 _MODEL_DONE

也就是说,每个 worker 跑完立刻把自己的结果落库,不等写线程、不等其他模型。多模型对比时模型 A 早完 30 秒,用户就早 30 秒能在库里看到 A 的答案。

落库入口 _persist_model_outcomeprocess_message.py:1225)用 persist_lock + persisted[] 数组保证每个模型恰好落库一次,可以从任何线程调用。四个调用点用一个枚举 _PersistContextprocess_message.py:1112)标注,纯粹为了日志归因:

context触发场景
WORKERworker 自己跑完(最常见)
STOP_BUTTON停止按钮,强制认领并存部分状态
NORMAL写线程看到所有 _MODEL_DONE 后兜底
POST_STEPS收尾阶段最后兜底

stop_button=True 是唯一能"抢跑"的路径(process_message.py:1244):正常情况下没成功也没出错的模型不会被落库(还在跑),但停止按钮会强行认领,把 in-flight 状态存下来,之后 worker 自己那次调用就变成 no-op。

7.4 一个模型崩了不影响别的

worker 抛异常时不会终止整条流(process_message.py:1413-1418):分类错误信息、把 API key 从消息里抹掉、把异常投进队列。写线程收到异常时发一个带 model_indexStreamingError,但故意不减 models_remaining——因为 finally 一定会再投一个 _MODEL_DONE,那才是唯一的完成信号(process_message.py:1529-1533)。

API key 脱敏做了两遍(emit 前 :1296、publish 前 :1416),因为 litellm_exception_to_error_msg 的兜底分支会原样返回 str(e),而那里面可能嵌着 key。


8. 中断与落库

8.1 停止信号:Redis 栅栏 + 空闲轮询

用户点停止 → 前端打 POST /chat/stop-chat-session/{chat_session_id}backend/onyx/server/query_and_chat/chat_backend.py:1360,处理函数 stop_chat_session)→ set_fence(chat_session_id, cache, True) 在缓存里种一个键(backend/onyx/chat/stop_signal_checker.py:23,调用点 chat_backend.py:1379)。TTL 10 分钟。

检查侧是 is_connectedstop_signal_checker.py:38)——没有这个键就代表可以继续。装配期把它包成闭包塞进 ChatTurnSetup.check_is_connectedprocess_message.py:1050-1051)。

写线程的排空循环里这样用(process_message.py:1501-1526):

try:
model_idx, item = merged_queue.get(timeout=_CANCEL_POLL_INTERVAL_S) # 0.05s
except queue.Empty:
if stream_buffer is not None:
stream_buffer.flush()
if not setup.check_is_connected():
... # 落库 → 发 user_cancelled → drain_done.set() → return
continue

有个诚实的细节值得说明:这是空闲时才轮询,不是严格的 50ms 定时器。只有当队列在 50ms 内没吐出新包时才会去查一次停止信号(_CANCEL_POLL_INTERVAL_Sprocess_message.py:1120)。token 密集流出的时候检查会被推后——但那正是用户不太可能在此刻按停止、且推后一点也无害的时候。同一个超时值还顺带承担了缓冲刷盘的节拍。

8.2 两道 Redis 栅栏,别搞混

栅栏语义TTL位置
chatsessionstop_fence停止信号:键存在 = 用户要求停10 分钟stop_signal_checker.py:6
chatprocessing_fence处理中标记:键存在 = 这个会话有流在跑,值是 run_id30 分钟backend/onyx/chat/chat_processing_checker.py:10

第二道栅栏有个容易忽略的问题:长任务可能跑超 30 分钟。所以写线程每 60 秒重新盖一次章(_FENCE_REFRESH_INTERVAL_Sprocess_message.py:1123:1356-1364),失效的栅栏会被恢复端读成"写线程已死"。

栅栏的清除权归写线程,不归请求生成器(_run_post_stepsprocess_message.py:1307-1311)。因为请求生成器可能因为客户端断线提前退出,而那时候生成还在跑。_stream_chat_turnfinally 里只有 run_started == False 时才自己清(process_message.py:1871)。

8.3 可重放缓冲:换台机器也能接着看

StreamBufferWriterbackend/onyx/chat/stream_buffer.py:64)把发出去的每一行 NDJSON 原样存进缓存,zlib 压缩 + 分块编号。目的很具体:任何一个 api-server pod 都能重放/续读一次正在跑的流stream_buffer.py:1-11)。

设计上的克制体现在三点:

  • 写失败绝不影响主流程。 异常一律吞掉记日志,把这次运行降级成"不可恢复"(stream_buffer.py:135-150)。
  • 超出上限就截断。 超过 CHAT_STREAM_BUFFER_MAX_BYTES 直接标 truncated,不再写(stream_buffer.py:116-126)。
  • 读端必须容忍缺块。 read_stream_chunks 遇到缺失或解压失败就把 gap=True 返回,调用方要退回读数据库里的最终消息,而不是拼一条断裂的流(stream_buffer.py:229-247)。

顺序上有个精细的地方:mark_done() 必须在 _run_post_steps() 清栅栏之前调(process_message.py:1579-1585),否则恢复端可能读到"栅栏已清但缓冲没标完成"的中间态,把尾巴丢掉。

8.4 落库时机

真正写数据库的是 save_chat_turnbackend/onyx/chat/save_chat.py:169),由 llm_loop_completion_handleprocess_message.py:1971)调起。

llm_loop_completion_handle 干的第一件事是在任何写库之前,用 getter 把 ChatStateContainer 的状态整个快照下来process_message.py:1983-1991)。因为在停止按钮路径上 worker 线程可能还在跑,直接读属性不是线程安全的。

它同时决定被停止时存什么(process_message.py:1996-2009):已经生成了部分答案就存 答案 + " ... \n\n生成已被用户停止",一个字都没有就存"生成已被用户停止"。

save_chat_turn 内部按八步走:更新消息正文和 token 数 → 建 SearchDoc 记录 → 建 tool_call ↔ doc 映射 → 建引用号 ↔ doc 映射(只保留流中真正 emit 过的引用号save_chat.py:284-286)→ 关联文档到消息 → 建 ToolCall 行 → 写引用 → 附上答案里真正引用到的代码解释器产物文件。

最后,落库完还会顺手判断历史是不是该压缩了(process_message.py:2080-2090)——历史压缩挂在上一轮的收尾,而不是下一轮的开头。


9. 那条岔路:deep research

worker 里有个分叉(process_message.py:1372-1388):

if n_models == 1 and setup.new_msg_req.deep_research:
if setup.chat_session.project_id:
raise RuntimeError("Deep research is not supported for projects")
run_deep_research_llm_loop(...)
else:
run_llm_loop(...)

三个约束一目了然:deep research 只在单模型下可用(多模型入口会直接拒掉,process_message.py:1954-1960)、不支持 Project、走完全不同的第 2 层实现(backend/onyx/deep_research/dr_loop.py)。

但它只换第 2 层。第 1 层(装配、线程、队列、落库)和第 3 层(run_llm_step_pkt_generator)原封不动复用——is_deep_research 只是解码器上的一个开关,作用是让 tool_choice == REQUIRED 时的前置文本被当成推理而非答案(llm_step.py:1258-1272)。它多出的那些包(DeepResearchPlanDeltaResearchAgentStart 等,streaming_models.py:355-391)也只是 PacketObj 联合里多几个成员。

深入见 03


10. 巧妙之处(可借鉴)

① 用共享队列换掉生成器链。 让"发包"变成一个可以传递的能力(Emitter),而不是调用栈的结构属性。代价是失去了背压——队列是无界的,靠 drain_done 兜底防止无消费者时无限增长(process_message.py:1568-1570)。emitter.py:32

② 写线程和读出协程分离。 生成的可靠性不再依赖 HTTP 连接的存活。客户端断线只是"少了个看客",不是"任务取消"。process_message.py:1604-1631

③ N=1 走 N>1 的路。 不为单模型留特例分支,把并发路径变成唯一路径——反过来也保证了并发路径天天被主流量走,不会腐烂。process_message.py:1159-1161

④ 助手消息 ID 在模型开口前就预留。 前端从第一个 token 起就知道该往哪个气泡里塞,也让"部分落库"有确定的目标行。process_message.py:962-973

⑤ 断手而不是劝退。 最后一轮把工具清单置空,用类型约束代替 prompt 约束。llm_loop.py:892-895

⑥ 状态容器只累积不决策。 把"要不要存、存什么"的判断全推到调用方,换来了"任意时刻可安全快照"。chat_state.py:34

⑦ 停止和正常结束产出同一种终止包。 前端不需要为取消写第二套收尾逻辑。process_message.py:1514-1521


11. 边界与局限

  • worker 不可中断。 停止按钮不能真的掐断在飞的 provider 请求,只能堵住出口丢弃后续输出。token 还在烧,只是没人看了。process_message.py:1522-1524
  • 停止检测只在队列空闲时发生。 token 密集流出时检测会被推后(§8.1)。
  • merged_queue 无界。 写线程死了才会靠 drain_done 止血;写线程活着但慢,队列就会涨。
  • 多模型上限 2–3 个。 硬编码校验,超出直接返回 VALIDATION_ERRORprocess_message.py:1946-1953
  • run_llm_step_pkt_generator 是雷区。 docstring 第二行写着 "DO NOT TOUCH THIS FUNCTION BEFORE ASKING YUHONG"(llm_step.py:1096)——三股流的边界条件(推理块闭合、空 answer 恢复、XML 兜底)互相纠缠,改一处容易崩另一处。
  • 缓冲可能有洞。 缓存驱逐或超限截断时,恢复端只能退回数据库里的最终消息,看不到中间过程。stream_buffer.py:8-11
  • MAX_LLM_CYCLES 是全局的。 不能按 persona 或按工具集配,装了工具多的 MCP 就得改环境变量。configs/chat_configs.py:12

12. 代码地图

主题文件路径符号名
一次性装配backend/onyx/chat/process_message.pybuild_chat_turn
线程池 + 排空循环backend/onyx/chat/process_message.py_run_models_run_model_drain_to_completion
读出协程 / 心跳backend/onyx/chat/process_message.py_read_stream_publish
一次一存的落库闸backend/onyx/chat/process_message.py_persist_model_outcome_PersistContext
统一异常与入口backend/onyx/chat/process_message.py_stream_chat_turnhandle_stream_message_objectshandle_multi_model_stream
非流式聚合backend/onyx/chat/process_message.pygather_streamgather_stream_full
落库编排backend/onyx/chat/process_message.pyllm_loop_completion_handle
agent 主循环backend/onyx/chat/llm_loop.pyrun_llm_loopMAX_LLM_CYCLES
空响应分类backend/onyx/chat/llm_loop.pyEmptyLLMResponseError_build_empty_llm_response_error
三股流解码backend/onyx/chat/llm_step.pyrun_llm_step_pkt_generator_close_reasoning_if_active_increment_turns
解码消费壳backend/onyx/chat/llm_step.pyrun_llm_step
工具调用坐标分配backend/onyx/chat/llm_step.py_extract_tool_call_kickoffs
发包通道backend/onyx/chat/emitter.pyEmitterNullEmitter
渲染坐标backend/onyx/server/query_and_chat/placement.pyPlacement
包类型backend/onyx/server/query_and_chat/streaming_models.pyPacketPacketObjStreamingTypeOverallStopSectionEnd
输入 / 输出容器backend/onyx/chat/chat_state.pyChatTurnSetupChatStateContainer
停止栅栏backend/onyx/chat/stop_signal_checker.pyset_fenceis_connected
处理中栅栏backend/onyx/chat/chat_processing_checker.pyset_processing_statusget_processing_run_id
可重放缓冲backend/onyx/chat/stream_buffer.pyStreamBufferWriterread_stream_chunks
落库backend/onyx/chat/save_chat.pysave_chat_turn
HTTP 入口(发消息 / 停止)backend/onyx/server/query_and_chat/chat_backend.pyhandle_send_chat_messagestop_chat_session
前端坐标镜像web/src/app/app/services/streamingModels.tsPlacement

接着读: 每轮循环里的第 ②③ 步(system prompt 怎么拼、历史怎么截断、文件放哪)→ 02 上下文工程;第 ⑤⑥ 步(工具定义、并行执行、模型不听话的兜底、deep research 子 agent)→ 03 工具与子 agent