源码地址:后端地址 前端地址
上一篇你跑通了环境,在终端看到curl -N刷出来的 SSE 事件流——event: TEXT_DELTA、event: TOOL_CALL_START、event: DONE一行行往外蹦。但那些事件名是谁定的?框架的Flux<AgentEvent>怎么变成前端看得懂的 JSON 信封?用户关了浏览器标签页,后端怎么知道该停下来而不是误报模型故障?
这一篇钻进流式协议的三个核心类:SseFrontendAdapter、AgentEventMapper、SseEnvelope。它们加起来不到 250 行代码,却是整条请求链路上框架与平台的边界线。
先搞清楚 SSE 是什么
SSE(Server-Sent Events)是个 HTTP 单向推送协议,响应头Content-Type: text/event-stream,body 是一段纯文本流,格式极简:
event: TEXT_DELTA data: {"text":"你好"} event: DONE data: {"replyId":"abc123"}每个事件两行:event:是事件名,data:是 JSON payload,事件之间空行分隔。浏览器原生EventSourceAPI 或fetch+ ReadableStream 都能消费。跟 WebSocket 比,SSE 走的是标准 HTTP,不需要协议升级,Nginx 反向代理天然支持(只要别开 buffering)。
loser-agent 选 SSE 不选 WebSocket,原因很实际:聊天场景是单向推送——用户发一条消息进来,服务端推一堆事件出去,不需要双向实时通信。SSE 比 WebSocket 简单一个数量级,Spring MVC 的SseEmitter开箱即用。
三次变换:从框架事件到浏览器
回到第 01 集那张调用顺序图的第七段,AguiChatController第 228-285 行。整条流式链路就三次变换:
HarnessAgent.streamEvents() → Flux<AgentEvent> (框架产出) SseFrontendAdapter.adapt() → Flux<SseEnvelope> (平台适配) emitter.send(SseEmitter.event()) → SSE 文本行 (Spring 推送)逐个看。
第一次变换:框架产出Flux<AgentEvent>
第 229 行:
Flux<AgentEvent>rawFlux=((HarnessAgent)agent).streamEvents(userMsg,rc);AgentScope 把 Agent 的整个执行过程统一建模为事件流。ReAct 循环的每一步——思考开始、思考增量、思考结束、工具调用开始、参数增量、调用结束、工具结果、最终文本——都是一个AgentEvent子类。框架返回的是 Reactor 的Flux,不是List,意味着事件是异步流式产出的:模型生成一个 token 就推一个TextBlockDeltaEvent,不等整句生成完。
这行代码是框架与平台的唯一交接点——框架在这里交出Flux<AgentEvent>,之后的事全归平台管。
第二次变换:SseFrontendAdapter适配为信封
第 235 行:
Flux<SseEnvelope>adapted=SseFrontendAdapter.adapt(rawFlux);SseFrontendAdapter是个纯静态适配器,全部代码 33 行,核心就一个方法:
publicstaticFlux<SseEnvelope>adapt(Flux<AgentEvent>events){returnevents.handle((event,sink)->{SseEnvelopeenv=AgentEventMapper.toEnvelope(event);if(env!=null){sink.next(env);}else{log.debug("Dropped event: {}",event.getType());}});}用 Reactor 的handle操作符逐个处理框架事件:调AgentEventMapper.toEnvelope()把框架事件映射成平台信封,映射成功就往下传,映射返回null就丢弃(DEBUG 日志记录,不静默丢)。
这里有个设计决策值得注意——为什么用handle而不是map。map是 1:1 变换,不能丢弃元素;handle的sink可以调next()也可以不调,天然支持过滤。框架会产出一堆TEXT_BLOCK_START、MODEL_CALL_*、SUBAGENT_EXPOSED等调试/保留事件,这些前端不需要,AgentEventMapper对它们返回null,handle静默丢弃。
信封长什么样:SseEnvelope
SseEnvelope是个 Javarecord,两个字段:
publicrecordSseEnvelope(Stringevent,StringdataJson){}event是 SSE 事件名(如TEXT_DELTA),dataJson是已经序列化好的 JSON 字符串。用record不是随便选的——它不可变、自带equals/hashCode/toString,作为流中的数据载体干净利落。
JSON 序列化在AgentEventMapper.json()里做(第 178-187 行),用LinkedHashMap保证字段顺序稳定。别小看这个——前端按到达顺序拼接THINKING_START的content增量,字段顺序不稳定会导致 JSON 解析抖动。
第三次变换:Spring 推送 SSE 文本
第 254 行的sendSse方法(587-597 行):
privatevoidsendSse(SseEmitteremitter,SseEnvelopeenv){try{emitter.send(SseEmitter.event().name(env.event()).data(env.dataJson()));}catch(IOExceptione){emitterLost=true;log.warn("[AguiChatController] SSE send failed for event={} (client disconnect?)",env.event());}}SseEmitter.event().name(x).data(y)是 Spring MVC 的 SSE 构建器,输出就是标准 SSE 文本格式。sendSse的try-catch是流式协议最容易出事的地方——用户关了浏览器标签页,TCP 连接断开,下一次emitter.send()抛IOException。这时候不能当模型故障报,得标记emitterLost = true。
事件映射表:16 种框架事件 → 11 种 SSE 事件名
AgentEventMapper.toEnvelope()是整个适配层的核心,一个方法 100 多行,16 个instanceof分支。把框架的 16 种事件映射成前端 11 种 SSE 事件名(加 3 种 Plan 业务事件)。
前端loser-agent-web/src/types/chat.ts第 6-24 行定义了完整枚举,这是前后端的协议契约:
| SSE 事件名 | 框架事件类 | 前端怎么用 |
|---|---|---|
TEXT_DELTA | TextBlockDeltaEvent | 逐 token 拼接助手回复文本 |
THINKING_START | ThinkingBlockStartEvent+ThinkingBlockDeltaEvent | 思考块开始/增量,前端按到达顺序拼接 |
THINKING_END | ThinkingBlockEndEvent | 思考块结束 |
TOOL_CALL_START | ToolCallStartEvent | 工具调用开始,带 callId + toolName |
TOOL_CALL_DELTA | ToolCallDeltaEvent | 参数 JSON 增量,前端累加 |
TOOL_CALL_END | ToolCallEndEvent | 工具调用参数结束 |
TOOL_RESULT | ToolResultEndEvent | 工具执行结果,带 isError 标志 |
HITL_CONFIRM | HintBlockEvent+RequireUserConfirmEvent | 人工审批弹窗,用户点确认/拒绝 |
MESSAGE_COMPLETE | AgentResultEvent | 完整回复文本 + replyId |
ERROR | ExceedMaxItersEvent+AllToolsDeniedEvent | 错误事件,带 code + recoverable |
DONE | AgentEndEvent | 流结束标志 |
几个映射决策值得细看。
思考增量复用THINKING_START。第 81-85 行,ThinkingBlockDeltaEvent映射的 SSE 事件名不是THINKING_DELTA而是THINKING_START。注释写得很直白:
// Reuse THINKING_START with content — frontend appends by design// (matches chat.ts ThinkingStartData comment: "content 可能是空串或一段增量")returnjson("THINKING_START",Map.of("content",nullToEmpty(e.getDelta())));框架把思考拆成了 Start/Delta/End 三个事件,但前端只需要"思考块的内容"——Start 给空串表示块开始,Delta 给增量,前端按到达顺序拼接。少一个事件类型,前端逻辑更简单。这是一个「框架粒度细、平台粒度粗」的典型适配。
工具结果不带文本。第 108-116 行,ToolResultEndEvent在 AgentScope 2.0.1 里只带ToolResultState(OK/ERROR),不带结果文本——文本走另一个TOOL_RESULT_TEXT_DELTA事件流,而这个事件被平台丢弃了(AgentEventMapper注释第 173-174 行说TOOL_RESULT_TEXT_DELTA属于 drop 之列)。所以TOOL_RESULT的result字段恒为空串,前端只看isError标志。工具结果的完整文本走持久化层(ac_agent_block)回查,不走 SSE。
HITL 两个来源映射到同一个事件名。HintBlockEvent(第 117-121 行)和RequireUserConfirmEvent(第 122-147 行)都映射成HITL_CONFIRM。前者带hint提示文本,后者没有——平台从ToolUseBlock的名字拼了一句"Confirm tool execution: xxx"。后者还额外塞了toolCalls数组(第 134-142 行),因为前端审批弹窗需要原样回传confirmResults才能恢复 ASKING 状态的工具调用。
框架事件被丢弃的那些
AgentEventMapper第 173-175 行:
// TEXT_BLOCK_START / END, TOOL_RESULT_START / TEXT_DELTA,// MODEL_CALL_*, SUBAGENT_EXPOSED, CUSTOM and others -> drop (debug only / reserved)returnnull;框架产出的TEXT_BLOCK_START/END是文本块的边界标记,但TEXT_DELTA已经给了每个 token,前端拼完就是完整文本,边界标记多余。MODEL_CALL_*是模型调用前后的元数据事件(用于审计),但审计走AuditedModel装饰器另有一条线,不需要走 SSE。SUBAGENT_EXPOSED是子代理暴露事件,当前前端不渲染。
丢掉不等于没用——toolAuditor::dispatch(第 234 行)在适配之前就挂了doOnNext,审计器从原始Flux<AgentEvent>拦截事件做异步落库,不受 SSE 丢弃影响。适配层只管前端需要什么,不管平台需要什么——这是职责分离。
错误处理:客户端断开 vs 模型故障
subscribe的三个回调里,onError最容易出事。第 256-273 行:
err->{toolAuditor.flushRemaining();if(emitterLost||isClientDisconnect(err)){log.info("[AguiChatController] SSE client disconnected, skip stream-failure report: {}",err.getClass().getSimpleName());}else{log.error("[AguiChatController] SSE error",err);reportStreamFailure(agentCode);}emitter.completeWithError(err);}核心逻辑:只有真实模型故障才报给熔断器,客户端断开不算。
为什么这点重要?用户关浏览器标签页时,Tomcat 检测到 TCP 断开,抛ClientAbortException;Spring 6 包成AsyncRequestNotUsableException。如果把这个当模型故障报给LlmConfigRegistry的熔断器,连续几个用户关页面就会触发熔断,切到备用模型——一个完全正常的模型被用户关页面的行为搞熔断了,这是 V2.4.1 踩过的坑。
isClientDisconnect()方法(607-624 行)遍历 cause 链识别三类断开:
privatebooleanisClientDisconnect(Throwableerr){Throwablet=err;while(t!=null){if(tinstanceofAsyncRequestNotUsableException)returntrue;if(tinstanceofIllegalStateException&&t.getMessage()!=null&&t.getMessage().contains("ResponseBodyEmitter has already completed")){returntrue;}Stringname=t.getClass().getName();if(name.equals("org.apache.catalina.connector.ClientAbortException"))returntrue;if(tinstanceofIOException&&t.getMessage()!=null&&t.getMessage().contains("中止了一个已建立的连接"))returntrue;t=t.getCause();}returnfalse;}四个判断条件对应四种断开表现:Spring 6 的AsyncRequestNotUsableException、emitter已完成后再 send 的IllegalStateException(总是跟在客户端断开的IOException后面)、Tomcat 的ClientAbortException、以及 Windows 环境下中文操作系统的IOException消息「中止了一个已建立的连接」。
emitterLost标志(第 600 行)是volatile的——因为sendSse的IOException和后续subscribe的onError可能在不同 Reactor 线程上发生。第一次sendSse失败时翻成 true,后续onError看到emitterLost直接跳过熔断上报,不用再走 cause 链。
Plan 事件:业务事件不走框架
AgentEventMapper还有一组不走toEnvelope()的方法——第 50-70 行的planUpdate/taskUpdate/planStatusChange。它们是 V2.10 Plan Mode 加的,产出自PlanRuntimeListener(平台自己的监听器),不是框架的AgentEvent。
调用链在第 230-232 行:
if(planListener!=null){rawFlux=rawFlux.doOnNext(ev->planListener.onEvent(ev).forEach(extra->sendSse(emitter,extra)));}planListener监听框架事件,但产出的是平台自己的SseEnvelope——PLAN_UPDATE/TASK_UPDATE/PLAN_STATUS_CHANGE。这些信封不经过SseFrontendAdapter.adapt(),直接sendSse推给前端。这是「平台在框架事件流上挂旁路」的模式:监听框架事件、产出业务事件、共用同一条 SSE 通道但不复用框架的事件类型。
HITL 确认事件的二次加工
第 236-238 行有一步可选的enrichHitlConfirm:
if(!ruleSourceByTool.isEmpty()){adapted=adapted.map(env->enrichHitlConfirm(env,ruleSourceByTool));}V2.14 角色差异化权限加的:如果配置了按工具的 HITL 规则,HITL_CONFIRM事件会被反序列化、注入ruleSource字段(标注规则来源:「角色规则(xxx)」或「全局规则」),再序列化回去推给前端。前端弹窗显示规则来源,让用户知道这条审批是哪条规则触发的。
这是适配层之后的「二次加工」——SseFrontendAdapter做框架→平台的基础映射,Controller 再做平台内部的业务增强。两层分离,各管各的。
一个请求的生命周期时序
把上面所有片段串起来,一次请求的 SSE 生命周期是这样走的:
- HTTP 进来:
POST /agui/chat,produces = text/event-stream,Spring MVC 返回SseEmitter后释放 Servlet 线程 - SSE 通道建立:设
X-Accel-Buffering: no关 Nginx 缓冲,注册三个回调做连接统计 - 框架产出事件:
streamEvents()返回Flux<AgentEvent>,ReAct 循环开始推事件 - 旁路监听:
planListener.onEvent()监听框架事件,产出 Plan 业务信封直接推送 - 审计拦截:
toolAuditor::dispatch从原始事件流拦截做异步落库 - 适配映射:
SseFrontendAdapter.adapt()把框架事件映射成SseEnvelope,丢弃前端不需要的 - HITL 增强:如有工具规则配置,
enrichHitlConfirm注入ruleSource - 审计上下文注入:
contextWrite把 threadId/runId/traceId 塞进 Reactor Context - 订阅推送:
subscribe的 onNext 逐条sendSse推前端;HITL_CONFIRM 做拦截计数 - 结束:onComplete 审计落库 + 记忆抽取 +
emitter.complete();onError 区分断开/故障
关键在第 8 步——审计上下文通过 Reactor Context 传递,不是 ThreadLocal。因为 SSE 的事件推送在 Reactor 调度线程上,线程不固定,ThreadLocal 会丢。Reactor Context 是响应式编程的「线程安全上下文」,跟着Flux走,不跟线程走。这个设计第 19 集审计篇细讲。
小结
记住三件事:
- 三次变换:框架
Flux<AgentEvent>→SseFrontendAdapter.adapt()→Flux<SseEnvelope>→SseEmitter.send()→ SSE 文本行。适配层 33 行代码,AgentEventMapper100 多行映射逻辑,是框架与平台的边界线。 - 事件映射表:框架 16 种事件映射成前端 11 种 SSE 事件名(加 3 种 Plan 业务事件)。映射不是 1:1——思考增量复用
THINKING_START、工具结果不带文本、两种 HITL 事件合流,都是「框架粒度细、平台粒度粗」的适配决策。 - 错误区分:客户端断开(关标签页/网络断)不报熔断器,只有真实模型故障才报——
isClientDisconnect()遍历 cause 链识别四类断开标志,emitterLostvolatile 标志做跨线程协调。
下一篇我们回到旅程的起点——灰度路由。RequestEnvResolver怎么从请求里识别环境,GrayRoutingService怎么按用户算版本号,AgentVariant的三种 key 拼装规则又是怎么回事。第 01 集说「找不到工厂就 refresh 再找」,那个 refresh 背后的 SmartLifecycle 三阶段注册,第 04 集全展开。