第 11 周 · Tether 源码课
LLM Streaming:从网络片段到类型化更新
学完你能做到
- 区分 SSE bytes、SseEvent、ChatResponseUpdate 与 SessionEvent
- 解释 LlmStreamRequest 为什么先做 detached snapshot
- 说明 IAsyncEnumerable 的逐项消费、backpressure 与 cancellation 边界
- 用 ReplayLlm 和 SseParserTests 观察无网络与有传输语法的两条路径
课程进度
- 第 1 周
- 第 2 周
- 第 3 周
- 第 4 周
- 第 5 周
- 第 6 周
- 第 7 周
- 第 8 周
- 第 9 周
- 第 10 周
- 第 11 周
- 第 12 周
- 第 13 周
- 第 14 周
- 第 15 周
- 第 16 周
一句话先懂
Streaming 不是把最终字符串切几刀,而是让调用方在响应尚未结束时,逐个消费已经解析、与 Provider 无关的 typed updates。
网络 Provider 把 SSE bytes 变成 ChatResponseUpdate;ReplayLlm 绕过网络却实现同一 ILlmClient seam;Agent 再决定如何显示和提交这些更新。
模型可能数秒后才完成整段回答。若必须等完整 response 才交给调用方,用户看不到进度,取消也要跨过一大块不可观察工作。IAsyncEnumerable<T> 把“响应正在到达”变成 C# 可以逐项等待的控制流。
先看大图
sourceEventSeqs 记录精确 provenance。一个类比:行李传送带
机场传送带不会等所有行李都装齐才一次出现。你站在出口,每来一件就检查标签、取走并处理,然后再等下一件。
- immutable request 是托运清单,交接后不能被旁边的人偷偷改写。
- SSE frame 是运输层一节车厢,可能是 data、keepalive 或不完整尾帧。
ChatResponseUpdate是拆掉 Provider 包装后的标准行李:text、reasoning、function call 等内容。await foreach是“取一件、处理一件、再请求下一件”。
类比边界:网络栈和 StreamReader 仍可能有内部 buffer,IAsyncEnumerable 不承诺每个 token 对应一个 TCP packet,也不保证零缓冲。它表达的是消费者可观察的异步序列与取消边界。
四层数据不要叫成同一个 “chunk”
| 层 | 类型或表示 | 谁负责解释 |
|---|---|---|
| 网络 bytes | UTF-8 response stream | HTTP transport |
| SSE frame | SseEvent(Event, Data) | SseParser 按空行分帧、合并 data lines |
| Provider-neutral update | ChatResponseUpdate | wire adapter 把 JSON grammar 转成 MEAI content |
| Session fact | assistant/chunk + assistant/message.sourceEventSeqs | Agent 逐块 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 流式响应。