第 11 周 · Tether 源码课

LLM Streaming:从网络片段到类型化更新

预计 90 分钟 先修:第 9 周:Session Events、会 await foreach 与 CancellationToken、知道 HTTP response 可以分段到达

学完你能做到

  • 区分 SSE bytes、SseEvent、ChatResponseUpdate 与 SessionEvent
  • 解释 LlmStreamRequest 为什么先做 detached snapshot
  • 说明 IAsyncEnumerable 的逐项消费、backpressure 与 cancellation 边界
  • 用 ReplayLlm 和 SseParserTests 观察无网络与有传输语法的两条路径

课程进度

  1. 第 1 周
  2. 第 2 周
  3. 第 3 周
  4. 第 4 周
  5. 第 5 周
  6. 第 6 周
  7. 第 7 周
  8. 第 8 周
  9. 第 9 周
  10. 第 10 周
  11. 第 11 周
  12. 第 12 周
  13. 第 13 周
  14. 第 14 周
  15. 第 15 周
  16. 第 16 周

一句话先懂

Streaming 不是把最终字符串切几刀,而是让调用方在响应尚未结束时,逐个消费已经解析、与 Provider 无关的 typed updates。

网络 Provider 把 SSE bytes 变成 ChatResponseUpdate;ReplayLlm 绕过网络却实现同一 ILlmClient seam;Agent 再决定如何显示和提交这些更新。

模型可能数秒后才完成整段回答。若必须等完整 response 才交给调用方,用户看不到进度,取消也要跨过一大块不可观察工作。IAsyncEnumerable<T> 把“响应正在到达”变成 C# 可以逐项等待的控制流。

先看大图

LlmStreamRequest messages + options detached snapshot ILlmClient IChatClient seam Network Provider HTTP bytes → SSE wire grammar → typed ReplayLlm JSON script → typed IAsyncEnumerable ChatResponseUpdate Agent await foreach text · reasoning · function calls 逐项处理 + 聚合 runtime text delta live UI 可立即显示 assistant/message assembled + source seqs assistant/chunk durable stream fact
Agent 把 typed updates 规范化为 durable chunk,再由 assembled message 的 sourceEventSeqs 记录精确 provenance。

一个类比:行李传送带

机场传送带不会等所有行李都装齐才一次出现。你站在出口,每来一件就检查标签、取走并处理,然后再等下一件。

  • immutable request 是托运清单,交接后不能被旁边的人偷偷改写。
  • SSE frame 是运输层一节车厢,可能是 data、keepalive 或不完整尾帧。
  • ChatResponseUpdate 是拆掉 Provider 包装后的标准行李:text、reasoning、function call 等内容。
  • await foreach 是“取一件、处理一件、再请求下一件”。

类比边界:网络栈和 StreamReader 仍可能有内部 buffer,IAsyncEnumerable 不承诺每个 token 对应一个 TCP packet,也不保证零缓冲。它表达的是消费者可观察的异步序列与取消边界。

四层数据不要叫成同一个 “chunk”

层类型或表示谁负责解释
网络 bytesUTF-8 response streamHTTP transport
SSE frameSseEvent(Event, Data)SseParser 按空行分帧、合并 data lines
Provider-neutral updateChatResponseUpdatewire adapter 把 JSON grammar 转成 MEAI content
Session factassistant/chunk + assistant/message.sourceEventSeqsAgent 逐块 commit,并在成功或中断边界组装 provenance

一个 SSE frame 可能不产生可见文字,也可能更新 tool-call arguments;一个 typed update 也可同时含多项 content。不要从 packet、frame 或 token 数量推断 session event 数量。

源码放大镜:请求为何先 snapshot

src/Tether.Llm/LlmStreamRequest.cs 在构造时 snapshot:

  • messages 的角色、内容与顺序;
  • provider/model/sampling/stop sequences;
  • tool declaration 的 name、description 与 JSON schema;
  • caller cancellation。

MaterializeMessages() 与 MaterializeOptions() 每次都生成 fresh mutable objects。llm/stream waterfall listener 或 Provider 即使修改自己拿到的对象,也不能反写 Agent 已组装的 request authority。

ILlmClient 直接扩展 Microsoft.Extensions.AI.IChatClient。Agent 使用 GetStreamingResponseAsync,而 AgentPipelineStream.cs 在调用前后验证 cancellation、Provider identity、non-null stream 与 non-null updates。

网络 Provider 与 ReplayLlm

  • SseParser.ParseAsync 支持 CRLF、multi-line data、comment keepalive 与 read-split frame;EOF 前没有空行的 partial frame 被丢弃。
  • MultiProviderClient 冻结本次 route/profile,再选择 wire grammar;它本身不做隐藏 retry,retry 属于 Agent 的 request-error policy。
  • ReplayLlm 从 JSON queue 直接 yield 一个 typed update;它不制造 SSE,也不联网。

Backpressure 与 Cancellation

await foreach 每次请求下一项前,会先执行当前循环体。若 UI listener 或 Agent 对当前 update 的处理变慢,producer 不会通过这条异步枚举无限地向调用方推送下一项——这就是序列层面的 backpressure。

取消有两个来源:request cancellation 与 enumeration cancellation。GuardStream 把它们链接,并在开始、每个 update 前后、枚举结束时检查。Provider 的 SseParser 也把 token 传给每次 read。取消不是“丢下后台 reader 就返回”,而是贯穿消费链。

动手实验:一项 update 不等于一个 SSE frame

实验目标:分别观察 transport parser 与 immutable request,再用 ReplayLlm 证明无网络实现仍遵守同一 typed stream seam。

1. 预测

对以下原始 SSE,写出会产生几个 SseEvent,每个 Data 是什么:

: keepalive
event: delta
data: {"text":
data: "hello"}

data: second

data: incomplete

注意最后一帧没有 blank-line delimiter。

2. 运行基础测试

dotnet test tests/Tether.Llm.Tests/Tether.Llm.Tests.csproj \
  --filter "FullyQualifiedName~SseParserTests|FullyQualifiedName~LlmStreamRequestTests"

3. 运行 ReplayLlm

apps/Tether.Cli/replies/text.json 只有:

[{"kind":"text","text":"hello-world"}]
course_home=$(mktemp -d)
dotnet run --project apps/Tether.Cli -- \
  --home "$course_home" \
  --replies apps/Tether.Cli/replies/text.json \
  --send "stream once"

ReplayLlm 为这一 turn yield 一个 ChatResponseUpdate,其中一个 TextContent 已含完整 hello-world。CLI 会立刻输出 text op;session snapshot 则在流结束后包含完整 assistant/message,不会因为叫 streaming 就自动变成 11 个字符事件。

4. 解释

SSE 预测答案是两个 complete frames:第一项 Data 为两条 data line 用换行连接的 {"text":\n"hello"},第二项为 second;keepalive 忽略,trailing partial 丢弃。Replay path 没有这一步,却在 ILlmClient 之后与网络 Provider 汇合。

检查理解

1. UI 已收到三个 agent/text-delta,能否推出 session 已 durable 三条 chunk event?

查看答案

不能。runtime delta 与 session event 是不同通道;当前 Agent 聚合 typed updates 后提交完整 assistant/message,assistant/chunk 还只是已声明 descriptor。

2. 为什么 Provider 不应悄悄重试一个已经向 caller yield 过部分内容的 request?

查看答案

caller 已观察到前缀,隐藏 retry 可能重复文字、reasoning 或 tool-call 片段。Tether 把可见 retry 放在 Agent request-error policy,由拥有 turn/step 语义的一层决定。

3. IAsyncEnumerable 是否保证一次 yield 就对应模型的一个 token?

查看答案

不保证。Provider wire grammar 决定如何把 frame 合成 ChatResponseUpdate;一个 update 可含多项 content,一个 token 也不是稳定的网络分帧单位。

本周带走

  • SSE frame、typed update、runtime delta 与 session event 是四层不同事实。
  • LlmStreamRequest 先 snapshot;每次 materialize 都给 listener/provider fresh objects。
  • IAsyncEnumerable 给出逐项等待、序列层 backpressure 与贯穿链路的 cancellation。

下一周模型不只返回 text,还会声明 tool calls;我们将沿安全闸门追踪调用直到 durable result。更多 wire 细节可查LLM 流式响应。

在 GitHub 上编辑此页