qwen-code 的 ACP-over-HTTP 可断点续传会话流:基于 SSELast-Event-ID的事件重放设计
【免费下载链接】qwen-codeAn open-source AI coding agent that lives in your terminal.项目地址: https://gitcode.com/GitHub_Trending/qw/qwen-code
导读
本篇技术指南围绕 sse-resumable-stream.md 展开,剖析 qwen-code 守护进程(qwen serve)如何为 ACP-over-HTTP(Streamable HTTP)传输层补上会话事件流的断点续传(resumable)能力。文章以该设计文档为主体骨架,结合仓库中的 EventBus 环形缓冲区、SSE 流写入器、连接注册表与 TypeScript SDK 实现,讲解Last-Event-ID游标如何在重连时驱动环形缓冲重放、会话流"分离-宽限期-回收"机制如何让重放真正生效,以及重放正确性的两个关键守卫。读完你可以完整理解/acp会话事件流从"只支持实时(live-only)"到"支持断点续传"的底层原理、实现变更与边界限制,并掌握--event-ring-size等运维参数的实际作用。
一、问题背景:实时会话流的"断档丢失"
ACP-over-HTTP 传输层的会话事件流由GET /acp(携带Acp-Session-Id请求头)承载,在设计补丁之前,它是**纯实时(live-only)**的:
- 既不在 SSE 帧中输出
id:序号; - 也不在重连时读取客户端的
Last-Event-ID请求头。
当控制面的代理(ingress proxy)在对话中途空闲关闭这条长连接时(守护进程本身会发送retry: 3000,而代理又频繁掐断长 SSE 连接),客户端会重连并重新声明会话所有权,但守护进程在断档期间产生的所有内容帧都会丢失——这些帧通常是携带agent_thought_chunk/agent_message_chunk的session/update通知。最终这一轮对话仍然会走到终止状态(turn_complete会被产出或被合成),于是 UI 显示"已完成",但正文却是空的或被截断的。重新发送同样的提示词可以工作,这正是关键线索:丢的是传输断档,而不是模型输出。
设计文档将这一症状与现场证据记录在集成笔记的 §1.8(sdk-known-issues.md)中。与此同时,文档明确区分了两个相邻问题:
- §1.7:会话流上丢失的"进行中的 prompt 响应"(JSON-RPC 响应不在事件环中,属于另一个跟踪项);
- §1.8:丢失的"内容帧",即全部经由 bus 的
session/update事件——这正是本次补丁要解决的对象。
二、已具备的重放引擎:EventBus 环形缓冲区
这个补丁之所以"小而巧",是因为重放引擎早已建成并被充分验证,缺口仅仅是/acp传输层没有接上它。引擎位于 packages/acp-bridge/src/eventBus.ts:
- 单调递增的事件 id:每个会话一个单调
id,从 1 开始,在publish()中分配(nextId)。事件必须先通过可序列化门禁(DAEMON-011)才会占用 id,被拒绝的事件不会烧掉序号,避免其他订阅者看到 3 → 5 这样的空洞; - 有界环形缓冲区:每个会话一个环,默认深度
DEFAULT_RING_SIZE = 8000(源码注释解释了从 1000 上调到 8000 的原因:单次长 prompt 可能产生数百帧,真实负载可达数倍于此),运维可用qwen serve --event-ring-size <n>覆盖; - 带游标订阅:
subscribeEvents(sessionId, { lastEventId, signal })会在实时事件流入之前,先重放环内id > lastEventId的帧,并发出若干合成控制帧:replay_complete——重放排空信号(无论是否有帧被重放都会发出,让消费者确定性结束"追赶中"状态);state_resync_required——环已被驱逐(ring evicted)/ 守护进程重启导致 epoch 重置(epoch reset)/ 重放字节预算超限时发出,提示客户端状态已不可信、需调用 loadSession 恢复;client_evicted、slow_client_warning——慢客户端背压治理帧。
注意合成控制帧(client_evicted、state_resync_required、replay_complete等)不携带id,因此不会在会话单调序列中占据槽位——否则其他健康订阅者会在实时流和续传时看到空洞,破坏连续性。
REST 表面早已消费这套机制:GET /session/:id/events读取last-event-id头(server.ts→parseLastEventId),传给subscribeEvents,并用formatSseFrame为每帧序列化出 SSEid:行。而/acp传输层(dispatch.ts的pumpSessionEvents、sse-stream.ts的SseStream.send)在补丁前全部缺席,sse-stream.ts源码注释里甚至明说:"no ring-bufferid:sequencing — resumability is RFD Phase 4, deferred"。下表对比了两条表面的差异:
| 层次 | REST/session/:id/events | /acpGET(补丁前) |
|---|---|---|
读取Last-Event-ID请求头 | 是 | 否 |
将lastEventId传给subscribeEvents | 是 | 否(dispatch.ts pumpSessionEvents) |
输出 SSEid:行 | 是(formatSseFrame) | 否(SseStream.send只写data:) |
三、线上决策:SSEid:行,而非 payload 内_meta
两条 SSE 表面承载的载荷形态不同,这决定了续传游标的放置位置:
- REST流传输
BridgeEvent信封({ id, v, type, data, _meta }),SDK 解析器(sdk-typescript/src/daemon/sse.ts)从 JSON 信封的id字段提取游标(只读data:行); /acp流传输的是裸 JSON-RPC 2.0 对象(session/update通知、session/request_permission请求、响应等),它们没有承载 bus 游标的信封id——而且 JSON-RPC 的id语义是"请求 id",不能挪作他用。
因此/acp的续传游标采用标准 SSEid:行,理由充分:
- EventSource 原生:符合规范的 SSE 客户端(包括随仓库 vendored 的
AcpHttpTransport)会自动记录最后一条id:并在重连时自动回填Last-Event-ID请求头,无需自定义逻辑; - 载荷纯净:JSON-RPC 协议帧内不注入非标准的
_meta.qwen.eventId; - 与 REST 对齐:
formatSseFrame在 REST 上输出的就是同样的id:行,因此两条表面共享同一套eventBus id 与同一套Last-Event-ID语义。
需要明确的是:只有 bus 来源的帧携带id:(session/update、session/request_permission、守护进程推送的通知)。在会话流上"搭车"的JSON-RPC 响应/回复不是 bus 事件,不携带id:——它们不在环内、刻意不被重放跟踪(丢失的进行中 prompt 响应是 §1.7 单独跟踪的问题)。合成的终止帧(client_evicted、stream_error等)没有 bus id,同样不输出id:行,避免在客户端续传的单调序列中烧掉槽位。
四、实现变更清单:从总线到线上的一整套打通
设计文档列出 6 项核心变更,全部集中在packages/cli/src/serve/acp-http/:
- transport-stream.ts:
send(message, id?: number)。可选的id即 bus 事件 id,用于 SSE 游标跟踪。该文件定义了TransportStream接口与DeliveryResult('delivered' | 'outcome_unknown' | 'closed' | 'failed'); - sse-stream.ts:
send(message, id?)在id !== undefined时于data:行前追加id: ${id}\n(镜像 RESTformatSseFrame)。实现细节:帧按id:→data:→ 载荷 → 空行 的顺序写入,所有写入经单一writeChain串行化(心跳注释不能插队),并遵守背压(res.write返回 false 时等待drain);流自身拥有写失败处理——首次失败记日志并关闭,杜绝"僵尸流"; - ws-stream.ts:
send(message, id?)接受并忽略id——WebSocket 是有状态连接,无 SSE 重放(与AcpWsTransport.supportsReplay = false一致),且不能把 SSE 的id:框架泄漏进 WS 裸 JSON 帧; - connection-registry.ts:
sendSession(sessionId, frame, id?)把id透传给传输层。会话级 pre-attach缓冲区改为存储"一个已序列化的 UTF-8 载荷 + 其可选游标 + 预算租约"(PreparedFrame),这样被缓冲的帧保留游标、又不必持有源对象或在 attach 时重复序列化。连接级回复复用同一表示; - dispatch.ts:
translateEvent为 bus 事件把event.id透传给每次sendSession/binding.stream.send调用;pumpSessionEvents(conn, sessionId, signal, lastEventId?)将lastEventId转发给subscribeEvents——直接复用既有环形重放; - index.ts:
GET /acp会话流分支读取Last-Event-ID请求头(经由严格版parseLastEventId,与 REST 同为"仅接受十进制数字"规则),传给pumpSessionEvents。
关键点是:eventBus/bridge 零改动——引擎被原样复用。parseLastEventId被抽取为共享模块 packages/cli/src/serve/sse-last-event-id.ts,REST 与/acp两条表面共用同一套严格接受/拒绝规则与运维日志,不会漂移。该模块还提供parseEventEpochHeader(X-Qwen-Event-Epoch请求头,用于 DAEMON-001 的过期 epoch 检测)。
五、让续传真正生效:会话流的"分离 + 宽限期 + 回收"
id:/Last-Event-ID的管道打通是必要但不充分的——仅靠它,在真实流程中永远不会触发。原因在于:此前,当会话 SSE 流在传输层关闭时,GET 处理器会执行完整的closeSessionStream拆卸流程:从ownedSessions移除会话、abort 进行中的 prompt、分离 bridge 客户端。而在真实的 EventSource/代理时序中(旧 socket 先关闭,客户端后重连),携带Last-Event-ID的重连会在游标被读取之前就被所有权检查以403拒绝——而且正在产出内容的 prompt 已经被 abort 了,重放引擎根本没有东西可以重连。
因此补丁把"传输层会话流关闭"从"拆卸"改为"分离(detach)"(AcpConnection.detachSessionStream):
- 只停止流本身 + 其事件订阅;
- 保留 binding、所有权、进行中的 prompt、bridge-client 注册,持续一个宽限期
SESSION_GRACE_MS(镜像连接级CONN_GRACE_MS,二者在 index.ts 中均为10_000即 10 秒); - 宽限期内重连即回收(reclaim):
attachSessionStream清除宽限定时器,环形重放回填断档; - 若无人重连,宽限定时器执行完整拆卸——约束失控 prompt 的成本;
- 显式
session/close与连接拆卸(destroy)仍然立即完整拆卸; - GET 处理器依据
stream.isClosed分支:传输层关闭 → 分离 + 宽限;pump 结束而流仍打开(子进程结束 / 迭代器错误)→ 完整关闭(僵尸流)。
实现细节值得注意:attachSessionStream遵循"先安装新流,再关闭旧流"的顺序,且旧流 pump 的收尾逻辑以binding.stream做身份守卫(identity-guard)——只有"自己仍然是该会话的活跃流"时才执行操作。这样旧流在关闭时落入"分离 + 宽限"而非"拆卸",进行中的 prompt 得以存活(这正是第 18 轮评审 G1 修复的回归问题:重连曾无条件 abort 进行中的 prompt)。
六、两个重放正确性守卫
宽限期/回收机制让重放路径首次可达,两个潜在的正确性隐患也随之暴露,因此随本补丁一同发布:
守卫一:既无重复投递、也无静默丢失(缓冲区 ↔ 环)
被缓冲的 bus 事件同时也在 EventBus 环内(它正是为了拿 id 才发布进去的)。因此在续传(存在Last-Event-ID)时,attachSessionStream拿到游标后完全不冲刷携带 id 的缓冲帧——从客户端游标开始的环形重放成为游标之后所有 bus 事件的唯一投递路径。
这是刻意设计:帧"已发送给现已死亡的 socket 但客户端从未收到"时,其 id 低于缓冲区的 id 却高于客户端的游标——如果"先冲刷缓冲区、再把重放游标推进过缓冲区",就会静默丢弃这帧。让环接管所有 bus 事件,每个事件恰好投递一次、无缺口。
而无 id 帧(经replySession路由的 JSON-RPC 回复)不是环事件,环不会重投——但 attach 时也不能冲刷:若在重放前冲刷缓冲的session/prompt结果,它会跑到其前置内容块之前(客户端先看到"done"再看到正文——正是 §1.8 要修的截断正文故障)。因此在续传时,无 id 帧被延迟:留在缓冲区,由事件 pump 在重放排空后(仅在replay_complete时)通过flushBufferedSessionFrames释放,保持原始流顺序。
关键约束:绝不能挂在state_resync_required上冲刷——EventBus 在重放帧之前就发出该帧(随后仍会在末尾发出replay_complete),若在此冲刷会把回复放到被重放内容之前。纯实时场景(无Last-Event-ID⇒ 无重放 ⇒ 无replay_complete)由 pump 的循环后安全冲刷兜底;全新连接(无Last-Event-ID)没有环锚点,立即按序冲刷整个缓冲区(与补丁前行为一致)。
相关实现可见 connection-registry.ts 的sendSessionReply(带anchorId水印的延迟投递)、releaseDeferredSessionReplies(pump 每投递一个内容事件后按水印释放)、endReplayDeferral(replay_complete边界)与flushBufferedSessionFrames(无条件最终冲刷)。
守卫二:重放下幂等的permission_request
permission_request是携带 id 的环事件,因此游标位于"尚未答复的 permission"之前的重连会重放它。补丁后的translateEvent复用该bridgeRequestId在conn.pending中的既有条目(对追赶中的同一条出站 JSON-RPC id 重发),而不是新铸第二个 id + 条目——不会产生孤儿 pending,也不会让按_meta.requestId去重的客户端收到双重 prompt。
七、向后兼容性
补丁对三类消费者都保持兼容:
- 旧客户端不发送
Last-Event-ID→lastEventId为undefined→subscribeEvents从实时开始,行为与今天完全一致; - 新增
id:行是向后兼容的 SSE——忽略该字段的客户端不受影响;基于 EventSource 的客户端自动开始跟踪它并在重连时回填; - vendored SDK
AcpHttpTransport在本补丁中显式启用重放:源码中readonly supportsReplay = true(packages/sdk-typescript/src/daemon/AcpHttpTransport.ts),重连时回填Last-Event-ID请求头,断档帧从环中重放,§1.8 内容丢失无需守护进程进一步改动即被闭合。任何仍上报supportsReplay = false且省略该头的消费者,守护进程侧改动保持惰性。外部agent-web传输的开关翻转不在本仓库范围(见下文)。
此外,预附加队列被计数与字节双重约束:单个流最多持有 256 帧,单个逻辑连接最多 1024 帧 / 64 MiB,全部 ACP HTTP 挂载共享进程级 4096 帧 / 256 MiB 预算。溢出时关闭"精确所有者"(session 或逻辑连接),而不是驱逐更旧的帧;续传丢弃携带 id 的缓冲事件时会释放其保留的字节租约。REST 表面完全不受影响。
八、测试计划:从单测到端到端
设计文档的测试计划覆盖四个文件(其中三个可在此仓库直接找到):
sse-stream.test.ts:send(msg, 7)在data:前输出id: 7\n;send(msg)(无 id)省略id:行;顺序为id:→data:→ 空行;- transport.test.ts(端到端,经
/acp传输层):- 实时
session/update帧现在携带id:行; - 携带
Last-Event-ID: N的GET /acp把游标流入subscribeEvents;无头的新流行为与今天一致; - 溢出的
Last-Event-ID(>MAX_SAFE_INTEGER)→ 退化为纯实时; - 真实"先关后连"顺序:先关闭旧 SSE,再用
Last-Event-ID重连——断言200 而非 403(所有权被保留)且 prompt未被 abort(宽限/回收); - 被重放的
permission_request复用 pending 条目(相同出站 id);
- 实时
- connection-registry.test.ts:非续传 attach 冲刷整个缓冲区并逐个线程化
id;续传attach(有游标)跳过携带 id 的帧(环重放接管)但仍冲刷无 id 的 JSON-RPC 回复;detachSessionStream在宽限期内保留所有权/prompt、到期后拆卸;宽限期内重连即回收(取消待执行的拆卸); - ws-stream.test.ts:
send(msg, id)忽略 id——WS 线上帧是裸 JSON,无 SSEid:框架泄漏。
九、明确不在本次范围(仍延后)
设计文档如实记录了以下边界,避免读者误以为本补丁覆盖了全部续传场景:
- WebSocket / HTTP/2 传输:不携带 SSE
id:/Last-Event-ID语义,另行处理; - §1.7 跨连接 permission resolve:投票 POST 到与流式 prompt 不同的
Acp-Connection-Id上的问题,是独立且涉及安全敏感性的后续项。本补丁只让permission_request翻译在重放下幂等,不新增会话级 requestId resolve;"已决议 permission 的响应重放幂等"(将已决议结果记录在会话级有界 LRU 以重发记录票)也归入同一后续; - 会话流上丢失的进行中 prompt 响应:恢复的内容帧都走 eventBus 环,JSON-RPC 响应不是环事件;
- 外部
agent-webAcpHttpTransport的supportsReplay翻转:位于不同仓库,本 PR 已解阻; - 经导出的 SDK 传输进行 permission 投票:导出传输把
session/request_permission暴露为permission_request事件,但 SDK 的投票 API(respondToPermission/respondToSessionPermission)映射到一个 ACP 守护进程没有处理器的session/permission请求——守护进程只接受以 JSON-RPC响应(回显出站_qwen_perm_Nid)形式的投票。同时,无订阅者的会话回复泵(ensureSessionReplyPump)会打开真实GET /acp会话流,守护进程将其视为活跃流,导致仅回复泵挂载时触发的permission_request被路由到该流并被泵丢弃(泵只转发 JSON-RPC 响应),使 mediator 挂起。需要 permission 处理的消费者应按文档契约先打开subscribeEvents再发起会话 RPC; - 在导出
AcpHttpTransport的subscribeEvents循环内发起会话 RPC:会话/acp流是单读器,消费端异步生成器在两次yield之间停驻时读取器不排水,在事件处理循环内await会话路由 RPC(session/set_model、session/prompt等)会挂起直到消费者拉取下一个事件。修复方向是把会话读取器改为始终排水 JSON-RPC 回复的后台泵,仅将DaemonEvent入队给迭代器; SESSION_STREAM_REPLY_METHODS⇄replySession漂移的自动化守卫:SDK 的该方法集合必须镜像dispatch.ts中的replySession(...)调用点(不同包),任何一侧遗漏都会导致无回复泵的sendRequest挂起至 abort。正确的提取器需要轻量数据流分析(session/prompt的回复不在其case块内发出,而是由 prompt 完成处理器在异步完成后从不同调用点触发),长期修复方向是让守护进程主动宣告会话路由方法名作为共享真源。
十、小结:从"只实时"到"可续传"的完整链路
至此,/acp会话事件流的续传链路可以一句话概括:GET /acp读取Last-Event-ID→pumpSessionEvents传给subscribeEvents→ EventBus 环按游标重放断档帧 →SseStream.send为每帧附加id:行 → 传输层断连时分离而非拆卸(10 秒宽限期)→ 重连回收并续传。
补丁的本质是"接线"而非"发明":重放引擎(环形缓冲、游标订阅、合成控制帧)早已在 packages/acp-bridge/src/eventBus.ts 中完备存在并被 REST 表面验证,本补丁把id:序号、Last-Event-ID解析、缓冲帧游标透传、分离-宽限-回收生命周期逐一接入/acp,并用两个正确性守卫(缓冲↔环单次投递、permission 重放幂等)封堵了重放路径上的潜在缺陷。对于熟悉 SSE 断点续传或想在自己的 ACP 客户端中利用Last-Event-ID恢复断档内容的开发者,这份设计与配套源码(sse-stream.ts、connection-registry.ts、dispatch.ts、sse-last-event-id.ts)提供了完整的参考实现与测试基线。
【免费下载链接】qwen-codeAn open-source AI coding agent that lives in your terminal.项目地址: https://gitcode.com/GitHub_Trending/qw/qwen-code
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考