网络、流式与取消
0. 先纠正一个容易埋雷的模型
Section titled “0. 先纠正一个容易埋雷的模型”用户点下“生成”以后,页面上只有一条进度线,后端却至少有三种互相独立的状态:
- 命令状态:创建任务的HTTP请求有没有被服务端接受;
- 任务状态:排队、运行、完成、失败或取消;
- 连接状态:事件流正在连接、已打开、断线退避或永久关闭。
这三种状态不能互相代替。
POST /runs返回成功,不代表任务已经完成;- SSE连接断开,不代表后台任务已经取消;
- 用户关闭页面,也不保证工作进程已经停下;
- 客户端没收到完成事件,不代表服务端没提交结果。
可靠系统首先要承认这些状态可能短暂不一致,再提供查询、续传、去重和终态收敛。
1. 一个请求经过了哪些层
Section titled “1. 一个请求经过了哪些层”对常见的HTTP/1.1或HTTP/2 HTTPS请求,可以先用下面的顺序排查:
DNS → TCP → TLS → HTTP → Router → Handler → Queue/Worker → StorageHTTP/3不走TCP,而是在QUIC上集成传输与TLS握手。但排障思路相同:先确定失败发生在哪一层,再讨论应用重试。
| 阶段 | 成功意味着什么 | 常见失败 | 客户端此时知道什么 |
|---|---|---|---|
| DNS | 域名解析成可连接地址 | NXDOMAIN、解析超时、缓存旧地址 | 应用请求通常还没发出 |
| TCP或QUIC | 到目标建立传输连接 | 拒绝、丢包、握手超时 | 服务端应用可能尚未看到请求 |
| TLS | 双方完成加密参数与证书校验 | 证书、SNI、协议版本错误 | HTTP语义还没开始 |
| HTTP | 请求行、Header和Body被传输并得到响应 | 4xx、5xx、响应中断 | 要结合方法与状态码判断 |
| 应用处理 | Router、鉴权、业务代码执行 | 校验失败、依赖超时、进程崩溃 | 请求可能已产生部分副作用 |
| 异步任务 | 命令转成持久任务并由Worker执行 | 队列积压、Worker失联、重复消费 | HTTP连接可能早已结束 |
这里最麻烦的是“结果未知”窗口:请求Body可能已经送达服务端,但客户端在响应回来前断线。客户端只看到网络错误,无法凭这个错误判断服务端是否已经写入数据库。
所以,网络错误不是业务失败证明。对于有副作用的命令,恢复必须依靠任务ID、幂等(Idempotency)幂等Idempotency同一个业务意图重复提交时,不会额外产生新的业务副作用。打开术语条目 →记录或可查询的业务状态。
2. HTTP层需要一个明确契约
Section titled “2. HTTP层需要一个明确契约”HTTP不是“传JSON的管道”。方法、状态码和响应表示共同构成契约。
一个长任务可以使用两步模型:
POST /runs HTTP/1.1Idempotency-Key: run-20260721-001Content-Type: application/json
{"prompt":"...","model":"..."}HTTP/1.1 202 AcceptedContent-Type: application/jsonLocation: /runs/run_01J...
{"runId":"run_01J...","state":"queued","events":"/runs/run_01J.../events"}202 Accepted只说明服务端接受了处理,不是完成证明。客户端随后通过任务资源读取当前状态,通过事件资源接收增量更新:
GET /runs/{runId} 当前权威状态GET /runs/{runId}/events 历史回放 + 实时事件DELETE /runs/{runId} 请求取消,而非承诺已经撤销这种设计把“创建一次任务”和“观察一个任务”分开。SSE重连只重建观察连接,不会再次执行POST /runs。
3. SSE不是“不断返回一点字符串”
Section titled “3. SSE不是“不断返回一点字符串””服务端事件流(Server-Sent Events, SSE)服务端事件流Server-Sent Events, SSE服务器通过一个持续的 HTTP 连接向浏览器单向发送文本事件。打开术语条目 →是一个单向、UTF-8文本事件协议。服务端响应类型是text/event-stream,一条消息由空行结束。
下面是一组完整事件,而不是随意切开的HTTP Chunk:
id: 41event: run.accepteddata: {"runId":"run_01J...","state":"queued"}
id: 42event: run.deltadata: {"text":"网络错误"}
id: 43event: run.completeddata: {"state":"completed","resultVersion":7}四个字段各有职责:
| 字段 | 职责 | 不应承担的职责 |
|---|---|---|
event | 区分状态、增量、工具事件、心跳和终态 | 不要把所有消息都塞成message再猜类型 |
data | 承载当前事件的数据 | 不要依赖跨事件拼接才能解析一个JSON对象 |
id | 更新客户端最后确认的事件标识(Event ID)事件标识Event ID标识事件在一个有序事件流中的位置,用于断线续传、去重和回放。打开术语条目 → | 不要用随机UI组件ID代替顺序游标 |
retry | 给原生EventSource提供重连等待建议 | 不能代替服务端限流或指数退避策略 |
以冒号开头的注释行可以作为Keep-alive。它只能维持中间代理或连接的活跃判断,不能证明任务仍在推进:
: keep-alive 2026-07-21T21:00:00Z原生EventSource会在连接意外关闭后尝试重连,并维护最后事件ID。浏览器可以在后续请求中发送Last-Event-ID。如果使用fetch()自行解析流,则客户端需要自己保存游标、退避和取消读取器。
事件ID应该指向什么
Section titled “事件ID应该指向什么”稳定的事件ID应满足:
- 在同一个任务内顺序可比较,或至少可以作为服务端游标;
- 只有事件已经进入可回放存储后,才对外发送该ID;
- 重连时可以从“严格晚于最后确认ID”的位置继续;
- 客户端重复收到同一ID时,应用结果不再改变。
不要在写入持久事件之前先推送到Socket。否则客户端看到id: 43后断线,服务端却只能从42恢复,协议出现无法解释的缺口。
4. 重连恢复的是观察,不是重新执行
Section titled “4. 重连恢复的是观察,不是重新执行”连接恢复可以按下面的顺序做:
1. 读取本地 lastAppliedEventId2. 查询 GET /runs/{id} 的权威任务状态3. 如果任务已到终态,获取最终结果并停止重连4. 如果仍在运行,从 lastAppliedEventId 之后回放5. 回放追上最新位置后,切换到实时尾部6. 对每个事件先去重,再更新UI如果协议提供“订阅已有任务”,第一条事件最好是当前状态快照,后面再跟增量。这可以缩小“先查询状态、后建立流”之间的竞态窗口。A2A协议的任务订阅就要求先返回当前Task,再发送后续状态或Artifact更新,直到终态关闭流。
一个客户端可以用下面的纯函数守住事件顺序:
type RunEvent = { id: number; type: 'accepted' | 'delta' | 'completed' | 'failed' | 'cancelled'; data: unknown;};
type CursorState = { lastAppliedId: number; terminal: boolean;};
function applyEvent(state: CursorState, event: RunEvent): CursorState { if (event.id <= state.lastAppliedId) return state; // 重放或重复投递 if (state.terminal) throw new Error('event-after-terminal-state');
return { lastAppliedId: event.id, terminal: ['completed', 'failed', 'cancelled'].includes(event.type), };}这段代码没有解决缺口。如果当前游标是41,下一条却是44,客户端应暂停应用并请求回放,而不是假设42和43无关紧要。
5. Timeout、Deadline与Cancellation不是同一个东西
Section titled “5. Timeout、Deadline与Cancellation不是同一个东西”超时(Timeout)超时Timeout操作超过明确期限后主动停止等待,并进入可分类的失败路径。打开术语条目 →是“最多等多久”的配置或观测结果;截止时刻(Deadline)截止时刻Deadline一个绝对时间点;到达后,整条调用链都不应再为本次操作启动新工作。打开术语条目 →是一个绝对截止时刻;取消(Cancellation)取消Cancellation调用方明确通知正在执行的工作尽快停止,并释放相关资源。打开术语条目 →是沿调用链传播的停止信号。
它们的关系可以写成:
remaining = parentDeadline - now - localReserveDeadline预算要逐层递减。 网关拿到30秒,不应该再给下游一个全新的30秒;否则三层依赖串起来可能运行90秒,最外层早已放弃。
function childTimeoutMs(parentDeadlineMs: number, reserveMs = 250): number { const remaining = parentDeadlineMs - Date.now() - reserveMs; if (remaining <= 0) throw new Error('deadline-exceeded-before-start'); return remaining;}Go的context.Context把Deadline、取消信号和请求级值跨API边界传下去。父Context取消时,派生Context也取消;子任务只能接收父级停止信号,不应反过来取消父任务。
JavaScript中的AbortSignal承担相似的通知职责:
async function loadToolResult(url: string, parentSignal: AbortSignal) { const deadline = AbortSignal.timeout(8_000); const signal = AbortSignal.any([parentSignal, deadline]);
const response = await fetch(url, { signal }); if (!response.ok) throw new Error(`tool-http-${response.status}`); return response.json();}取消必须一直传到真正占资源的地方:HTTP读取器、数据库查询、队列任务、子进程和内部Pipeline。只把前端按钮改成“已取消”没有释放任何资源。
同时要承认一个边界:取消不是时间倒流。已经提交的事务、已经发送的邮件或已经执行的外部写入,不能靠关闭连接撤销。它们需要补偿动作、幂等记录或人工处理。
为什么Buffer不能替代取消
Section titled “为什么Buffer不能替代取消”下游提前停止读取时,上游可能永久阻塞在发送操作上。Go Pipeline示例把这种情况称为资源泄漏:阻塞的Goroutine仍占内存,栈上的引用也妨碍垃圾回收。
随手增大Buffer只是把阻塞推迟。只要生产总量或下游读取量发生变化,Buffer仍会填满。正确设计需要一个广播式停止信号,让所有上游阶段都能从阻塞点退出。
6. 重试之前先回答“第一次发生了什么”
Section titled “6. 重试之前先回答“第一次发生了什么””重试(Retry)重试Retry在可分类的临时失败后再次尝试;是否安全取决于副作用、幂等性和剩余时间。打开术语条目 →不是统一包一层循环。下面是最小的重试决策表:
| 操作与失败位置 | 第一次可能有副作用吗 | 自动重试 | 条件 |
|---|---|---|---|
| DNS或连接建立前失败 | 通常没有到达应用 | 可以 | 有总Deadline与退避 |
GET读取状态时断线 | 不应有业务写入 | 可以 | GET本身保持安全语义 |
POST发送前本地校验失败 | 没有 | 不需要 | 先修请求 |
POST发送后、响应前断线 | 结果未知 | 默认不可以 | 先查任务或使用幂等键 |
| 带幂等键的创建命令超时 | 可能已经创建 | 可以重放同一命令 | Key与请求摘要必须一致 |
| SSE连接断开 | 后台任务通常仍在运行 | 只重连事件流 | 不能重发创建命令 |
429或临时503 | 由接口语义决定 | 有条件 | 尊重Retry-After、加Jitter |
永久校验错误4xx | 请求本身无效 | 不可以 | 修改输入后才是新请求 |
一个幂等记录至少要绑定:
(scope, idempotencyKey, requestHash) -> status, taskId, response相同Key加相同请求,应返回同一个任务或结果;相同Key加不同请求,应明确冲突,不能静默复用。
重试还必须受总预算限制:
for (let attempt = 0; attempt < maxAttempts; attempt += 1) { const remaining = deadlineMs - Date.now(); if (remaining < minimumAttemptBudgetMs) throw new Error('retry-budget-exhausted');
try { return await callOnce({ timeoutMs: remaining, idempotencyKey }); } catch (error) { if (!isTransient(error) || !isRetrySafe(error)) throw error; await sleep(withJitter(backoffMs(attempt), remaining)); }}maxAttempts、Deadline、退避和Jitter要一起存在。缺任何一个都可能形成Retry Storm。
7. 背压从一条简单关系开始
Section titled “7. 背压从一条简单关系开始”设事件到达速率为λ,客户端稳定处理速率为μ。当λ > μ时:
积压增长率 = 到达速率 - 处理速率 = λ - μBuffer只能容纳一段时间的差值。比如UI每秒只能渲染20次,而服务端每秒推送100个Token事件,逐事件更新DOM会让队列、内存和交互延迟持续上升。
背压(Backpressure)背压Backpressure下游处理不过来时,限制上游继续生产数据的流量控制机制。打开术语条目 →可以在不同层处理:
- 生产端:降低非必要事件频率,把多个Token合成一批;
- 队列端:设置有界容量,拒绝或降级,而不是无限增长;
- 客户端:按动画帧或固定窗口批量渲染;
- 协议端:把可丢弃的进度事件与不可丢的状态事件分开;
- 运营端:根据Lag、队列长度和最老消息年龄扩容或限流。
TCP流量控制只能避免把接收窗口无限塞满。它不知道run.completed比十条打字动画更重要,因此应用层仍要定义丢弃、合并和保留策略。
一个常见策略是:
delta事件:允许批量合并progress事件:只保留最新值state-change事件:必须持久化并按序交付terminal事件:必须持久化、可查询、可回放8. SSE、持久事件流与交付语义
Section titled “8. SSE、持久事件流与交付语义”SSE解决浏览器如何持续接收事件,不自动提供持久化、消费确认或Exactly-once。
如果后端使用Redis Streams一类持久流:
XADD追加事件并生成有序ID;- Consumer Group把新消息分给组内消费者;
- 已投递但未确认的消息进入Pending Entries List(PEL);
XACK表示消费者已经处理;- 消费者失联后,其他消费者可以Claim超时的Pending消息。
这是一种可重投模型。消息可能被处理一次以上,因此副作用仍要幂等。所谓Exactly-once通常不是传输层送来的魔法,而是“可重投交付 + 幂等业务写入 + 去重记录”共同得到的业务效果。
浏览器SSE游标和Worker消费确认是两套游标:
Worker cursor 决定后台处理到了哪里Event cursor 决定任务事件写到了哪里Client cursor 决定这个浏览器应用到了哪里把它们混成一个lastId,故障恢复时就无法分辨“任务没处理”“事件没发布”还是“客户端没看到”。
9. 一条可诊断的Trace应该记录什么
Section titled “9. 一条可诊断的Trace应该记录什么”网络日志如果只有request failed,几乎没有排障价值。一次长任务的追踪记录(Trace)追踪记录Trace跨步骤记录一次任务的时间、公开状态、调用、结果和关联身份。打开术语条目 →至少要能关联:
traceId, requestId, runId, idempotencyKeyHTTP method, route, statusconnectionOpenedAt, firstEventAt, lastEventAtlastProducedEventId, lastDeliveredEventId, lastAppliedEventIdqueueWaitMs, processingMs, streamDurationMsattempt, retryReason, backoffMscancelRequestedAt, cancelObservedAt, terminalState常用派生指标:
- Time to First Byte(TTFB):连接建立到首字节;
- Time to First Event(TTFE):创建任务到第一个有意义事件;
- Inter-event Gap:相邻事件空档;
- Reconnect Count:一次任务重连次数;
- Replay Count:恢复时补发事件数;
- Cancellation Latency:发出取消到Worker停止的时间;
- Queue Lag:最新事件与消费者游标的距离。
| 现象 | 先查什么 | 可能根因 | 不该立即做什么 |
|---|---|---|---|
| 一直没有首事件 | DNS/握手、HTTP状态、代理Buffer、Queue Wait | 代理缓存、Worker未取任务、鉴权失败 | 直接无限重试创建命令 |
| 流每隔固定时间断开 | 空闲超时、心跳间隔、中间代理配置 | 没有Keep-alive、平台连接上限 | 把心跳当任务进度 |
| 重连后文字重复 | Last-Event-ID、客户端去重、回放起点 | 从>= cursor而非> cursor回放 | 仅在UI字符串层去重 |
| 任务完成但UI仍转圈 | 终态是否持久化、客户端是否查询状态 | 完成事件丢失、连接先断 | 再执行一次任务 |
| 用户取消后资源仍增长 | 取消信号是否到Worker/子进程 | 只关闭前端连接、下游忽略Signal | 只缩短前端Timeout |
| 延迟持续上升但CPU不高 | 队列长度、最老事件年龄、渲染频率 | 下游慢、无界Buffer、单消费者卡住 | 继续增大Buffer |
| 同一外部动作执行两次 | Idempotency记录、Attempt、Worker PEL | 超时后盲重试、消息重投 | 把日志去重当业务去重 |
10. 一个完整的恢复状态机
Section titled “10. 一个完整的恢复状态机”客户端状态可以保持很少,但每个转换要有证据:
IDLE └─ create accepted ─→ CONNECTINGCONNECTING ├─ stream open ─────→ STREAMING ├─ terminal query ──→ TERMINAL └─ transient error ─→ BACKOFFSTREAMING ├─ ordered event ───→ STREAMING ├─ event gap ───────→ REPLAYING ├─ terminal event ──→ TERMINAL └─ connection lost ─→ BACKOFFBACKOFF ├─ deadline remains ─→ CONNECTING └─ budget exhausted ─→ RECOVERY_REQUIREDREPLAYING ├─ caught up ───────→ STREAMING └─ history missing ─→ SNAPSHOT_REQUIRED终态必须可查询,不能只存在于一条可能丢失的SSE消息里。事件历史允许按保留策略清理,但清理后要提供状态快照;否则老客户端永远无法追上。
11. 实验:不要只画正常路径
Section titled “11. 实验:不要只画正常路径”实验A:断线续传
Section titled “实验A:断线续传”准备固定事件1..20,在客户端收到7后强制断线:
- 重连请求带上游标
7; - 服务端从
8开始回放; - 故意让
9重复一次,确认客户端只应用一次; - 故意删除
11,确认客户端检测到缺口; - 收到终态后关闭连接并停止重试。
验收不是“最后文字看起来一样”,而是lastAppliedEventId严格单调,副作用计数不增加。
实验B:取消传播
Section titled “实验B:取消传播”启动三个阶段:读取输入、调用慢工具、写入结果。分别在三个阶段取消,记录:
- 哪个阶段观察到Signal;
- 取消延迟;
- 是否仍有未退出任务;
- 是否发生已提交副作用;
- 是否需要补偿。
验收是取消后资源回到基线,而不只是HTTP请求返回AbortError。
实验C:背压
Section titled “实验C:背压”让生产端从每秒10个事件逐步升到每秒200个,客户端固定每秒处理50个。记录队列长度、最老事件年龄、内存和渲染延迟。再比较:
- 无界队列;
- 固定Buffer;
- Delta批量合并;
- 有界队列加生产端限速。
你应该看到固定Buffer只延后故障,限速或合并才改变长期积压趋势。
12. 完成本课的检查线
Section titled “12. 完成本课的检查线”你应当能独立回答:
- HTTP请求失败发生在DNS、连接、TLS、HTTP还是应用层?
- 网络断线后,第一次写入是失败、成功还是未知?
- SSE的
id、event、data、retry分别做什么? - 为什么重连事件流不等于重试创建命令?
- Deadline如何逐层缩短,取消如何传播到Worker?
- 哪些操作可以自动重试,哪些必须先查幂等记录?
- 当
λ > μ时,系统在哪一层施加背压? - 如何用Trace区分代理Buffer、Worker积压、事件丢失与客户端渲染慢?
如果答案里只有“加Timeout”“失败就重试”,这部分还没有学完。