【AgentScope 2.0】03-一条 SSE 事件流是怎么从框架走进浏览器的:loser-agent 流式协议全解
2026/9/12 20:18:56 网站建设 项目流程

源码地址:后端地址 前端地址

上一篇你跑通了环境,在终端看到curl -N刷出来的 SSE 事件流——event: TEXT_DELTAevent: TOOL_CALL_STARTevent: DONE一行行往外蹦。但那些事件名是谁定的?框架的Flux<AgentEvent>怎么变成前端看得懂的 JSON 信封?用户关了浏览器标签页,后端怎么知道该停下来而不是误报模型故障?

这一篇钻进流式协议的三个核心类:SseFrontendAdapterAgentEventMapperSseEnvelope。它们加起来不到 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而不是mapmap是 1:1 变换,不能丢弃元素;handlesink可以调next()也可以不调,天然支持过滤。框架会产出一堆TEXT_BLOCK_STARTMODEL_CALL_*SUBAGENT_EXPOSED等调试/保留事件,这些前端不需要,AgentEventMapper对它们返回nullhandle静默丢弃。

信封长什么样:SseEnvelope

SseEnvelope是个 Javarecord,两个字段:

publicrecordSseEnvelope(Stringevent,StringdataJson){}

event是 SSE 事件名(如TEXT_DELTA),dataJson是已经序列化好的 JSON 字符串。用record不是随便选的——它不可变、自带equals/hashCode/toString,作为流中的数据载体干净利落。

JSON 序列化在AgentEventMapper.json()里做(第 178-187 行),用LinkedHashMap保证字段顺序稳定。别小看这个——前端按到达顺序拼接THINKING_STARTcontent增量,字段顺序不稳定会导致 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 文本格式。sendSsetry-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_DELTATextBlockDeltaEvent逐 token 拼接助手回复文本
THINKING_STARTThinkingBlockStartEvent+ThinkingBlockDeltaEvent思考块开始/增量,前端按到达顺序拼接
THINKING_ENDThinkingBlockEndEvent思考块结束
TOOL_CALL_STARTToolCallStartEvent工具调用开始,带 callId + toolName
TOOL_CALL_DELTAToolCallDeltaEvent参数 JSON 增量,前端累加
TOOL_CALL_ENDToolCallEndEvent工具调用参数结束
TOOL_RESULTToolResultEndEvent工具执行结果,带 isError 标志
HITL_CONFIRMHintBlockEvent+RequireUserConfirmEvent人工审批弹窗,用户点确认/拒绝
MESSAGE_COMPLETEAgentResultEvent完整回复文本 + replyId
ERRORExceedMaxItersEvent+AllToolsDeniedEvent错误事件,带 code + recoverable
DONEAgentEndEvent流结束标志

几个映射决策值得细看。

思考增量复用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_RESULTresult字段恒为空串,前端只看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 的AsyncRequestNotUsableExceptionemitter已完成后再 send 的IllegalStateException(总是跟在客户端断开的IOException后面)、Tomcat 的ClientAbortException、以及 Windows 环境下中文操作系统的IOException消息「中止了一个已建立的连接」。

emitterLost标志(第 600 行)是volatile的——因为sendSseIOException和后续subscribeonError可能在不同 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 生命周期是这样走的:

  1. HTTP 进来POST /agui/chatproduces = text/event-stream,Spring MVC 返回SseEmitter后释放 Servlet 线程
  2. SSE 通道建立:设X-Accel-Buffering: no关 Nginx 缓冲,注册三个回调做连接统计
  3. 框架产出事件streamEvents()返回Flux<AgentEvent>,ReAct 循环开始推事件
  4. 旁路监听planListener.onEvent()监听框架事件,产出 Plan 业务信封直接推送
  5. 审计拦截toolAuditor::dispatch从原始事件流拦截做异步落库
  6. 适配映射SseFrontendAdapter.adapt()把框架事件映射成SseEnvelope,丢弃前端不需要的
  7. HITL 增强:如有工具规则配置,enrichHitlConfirm注入ruleSource
  8. 审计上下文注入contextWrite把 threadId/runId/traceId 塞进 Reactor Context
  9. 订阅推送subscribe的 onNext 逐条sendSse推前端;HITL_CONFIRM 做拦截计数
  10. 结束:onComplete 审计落库 + 记忆抽取 +emitter.complete();onError 区分断开/故障

关键在第 8 步——审计上下文通过 Reactor Context 传递,不是 ThreadLocal。因为 SSE 的事件推送在 Reactor 调度线程上,线程不固定,ThreadLocal 会丢。Reactor Context 是响应式编程的「线程安全上下文」,跟着Flux走,不跟线程走。这个设计第 19 集审计篇细讲。

小结

记住三件事:

  1. 三次变换:框架Flux<AgentEvent>SseFrontendAdapter.adapt()Flux<SseEnvelope>SseEmitter.send()→ SSE 文本行。适配层 33 行代码,AgentEventMapper100 多行映射逻辑,是框架与平台的边界线。
  2. 事件映射表:框架 16 种事件映射成前端 11 种 SSE 事件名(加 3 种 Plan 业务事件)。映射不是 1:1——思考增量复用THINKING_START、工具结果不带文本、两种 HITL 事件合流,都是「框架粒度细、平台粒度粗」的适配决策。
  3. 错误区分:客户端断开(关标签页/网络断)不报熔断器,只有真实模型故障才报——isClientDisconnect()遍历 cause 链识别四类断开标志,emitterLostvolatile 标志做跨线程协调。

下一篇我们回到旅程的起点——灰度路由。RequestEnvResolver怎么从请求里识别环境,GrayRoutingService怎么按用户算版本号,AgentVariant的三种 key 拼装规则又是怎么回事。第 01 集说「找不到工厂就 refresh 再找」,那个 refresh 背后的 SmartLifecycle 三阶段注册,第 04 集全展开。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询