feat(agent): AgentEventBridge 恢复流式(streamEvents→FlowEvent)+ 流式路径 ChatUsage 线程安全

Task 6.1:恢复 process() 流式路径。当 slot 有 FlowEvent 监听者(ExecuteOption.eventListener)
时走 ReActAgent.streamEvents → AgentEventBridge 映射成 FlowEvent 发布;否则仍 call()。

- 新增 AgentEventBridge.streamAndPublish:TEXT_BLOCK_DELTA→agent.reasoning、
  TOOL_RESULT_*→agent.tool_result、HINT_BLOCK→agent.summary、AGENT_RESULT→agent.result(last=true),
  HITL(RequireUserConfirm/RequireExternalExecution)可选透传。
- ChatUsage 线程安全(findings R-stream):真实 vendor 模型下 onModelCall 在 boundedElastic
  调度线程执行,ThreadLocal 读不到 bind 的累加器。改用 reactor Context 双源(Context 优先、
  ThreadLocal 回退),process() 经 bindToContext 注入;保留 void bind() 二进制兼容,
  新增 bindAndReturn()。流式/非流式两路径都注入 Context。
- StreamingBridgeTest:mock 模型流式回放,断言 reasoning 增量 + 末尾 agent.result(last=true) +
  流式下 getChatUsage() 正确(验证 Context 传播)。
- findings 追加 R-stream 节(streamEvents/call 线程模型 + Context 方案)。

Tests run: 21, Failures: 0, Errors: 0(StreamingBridge/ChatUsageMiddleware/ProcessIntegration/
ShellPermissionBehavior/SkillLoading/SkillTracking)。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
everywhere.z 2026-06-20 11:35:33 +08:00
parent 6159af3a23
commit dd08efc5b3
9 changed files with 751 additions and 38 deletions

View File

@ -559,3 +559,104 @@ Task 4.1 的 `strict` 语义在 <b>resolver 层</b>实现,分两段:
- per-invocation 绑定同 `ChatUsageMiddleware``bind()`/`unbind()` 操作 ThreadLocal `Set<String>``process()` 入口 bind、出口 finally unbind。`usedSkills()` 静态读当前线程集合(未 bind 返回空 List
- 工具名常量 `LOAD_SKILL_TOOL_NAME = "load_skill_through_path"`、入参 key `skillId`(与 1.0 `SkillTrackingHook` 一致)。
---
## R-streamstreamEvents / call 的线程模型Task 6.1 探针确认)
**来源:** 对 `agentscope-2.0.0-RC3-sources.jar``io/agentscope/core/ReActAgent.java`
`buildAgentStream` 795849、`reasoningStream` 20202045、`doCall` 925953
`io/agentscope/core/middleware/MiddlewareChain.java`
`io/agentscope/core/model/OpenAIChatModel.java`stream 方法 160181全量源码通读
JDK 21直接读 RC3 源码,比反射更权威)。结论已用于 Task 6.1 `AgentEventBridge` +
ChatUsage 线程安全方案选择。
### (a) streamEvents 的执行链
`ReActAgent.streamEvents(List<Msg>, RuntimeContext)` 直接委派给 `buildAgentStream(...)`
```java
private Flux<AgentEvent> buildAgentStream(List<Msg> msgs, RuntimeContext context, ...) {
Function<AgentInput, Flux<AgentEvent>> core =
input -> Flux.<AgentEvent>create(sink -> {
sink.next(new AgentStartEvent(...));
Mono<Msg> lifecycle = runLifecycle(input.msgs(), doCallFn);
lifecycle.contextWrite(c -> c.put(EVENT_SINK_KEY, sink))
.contextWrite(c -> c.put(AgentEventEmitter.CONTEXT_KEY, sink::next))
.doFinally(sig -> { sink.next(new AgentEndEvent(replyId)); sink.complete(); })
.contextWrite(subscriberCtx)
.subscribe(finalMsg -> sink.next(new AgentResultEvent(finalMsg)),
sink::error);
}, FluxSink.OverflowStrategy.BUFFER);
return MiddlewareChain.build(middlewares, this, context, MiddlewareBase::onAgent, core)
.apply(new AgentInput(msgs));
}
```
`call(...)` 同样走 `buildAgentStream`762784{@code call()} 与 {@code streamEvents()} 共用 buildAgentStream 核心,{@code onAgent} middleware chain 在两条路径上都触发)。
### (b) middleware onModelCall 在哪条线程被调用?
`reasoningStream`reasoning 阶段入口):
```java
Flux<AgentEvent> reasoningStream(ReasoningContext ctx, List<Msg> msgs, ...) {
Function<ModelCallInput, Flux<AgentEvent>> modelCallCore = mci -> modelCallStream(ctx, mci, true);
return MiddlewareChain.build(middlewares, this, rc, MiddlewareBase::onModelCall, modelCallCore)
.apply(new ModelCallInput(msgs, tools, options, model))
.doOnNext(this::publishEvent);
}
```
`MiddlewareChain.build(...).apply(input)`<b>同步</b>的 Java 函数链构造——它在此刻
调用 `middleware.onModelCall(...)`。该调用发生在<b>订阅 {@code reasoningStream} 的那条线程</b>上。
### (c) 谁在订阅 reasoningStream—— `.subscribeOn(Schedulers.boundedElastic())`
`modelCallStream` 内部 `mci.model().stream(...)`,对真实 vendor 模型(实测
`OpenAIChatModel.stream`,第 180 行)结尾是 `.subscribeOn(Schedulers.boundedElastic())`
故:当 reasoning 循环订阅 `reasoningStream` 时,整个上游(包括 `onModelCall`
<b>链构造</b>)在 `boundedElastic` 工作线程上被订阅——<b>不在</b> `process()` 调用
`bind()` 的 HTTP 线程上。
**关键结论:真实 vendor 模型下,无论 `call()` 还是 `streamEvents()`middleware
`onModelCall`/`onActing` 都在 `boundedElastic` 调度线程上被调用,而非 caller 线程。**
Task 5.1 的 `ChatUsageMiddlewareTest` / `ProcessIntegrationTest` 之所以"无线程跳变"
通过,是因为它们用 `CannedReplyModel`——一个 `Flux.just(resp)` 的 mock`subscribeOn`
故 emit 留在 caller 线程上。真实 HTTP 模型不满足此前提。)
### (d) 这对 ThreadLocal 累加器意味着什么
`ChatUsageMiddleware.onModelCall`<b>方法体(链构造)</b>里读 `BOUND.get()`——这是
在 boundedElastic 线程上读的,而 `bind()` 在 HTTP 线程上写——<b>读不到,返回 null
usage 丢失</b>。(`Accumulator` 的 `synchronized` 只保证写端/读端互斥,不解决"读的线程
压根没有累加器"的问题。)
`SkillTrackingMiddleware.onActing` 同构问题:`BOUND.get()` 在 boundedElastic 上返回 null
`usedSkills()` 恒为空Task 6.1 仅修复 ChatUsageSkillTracking 的修复推迟,影响只是
`usedSkills()` 在真实模型下为空,无 token 计费正确性问题)。
### (e) 修复方案(采用 breactor Context
Reactor `Context` 在 reactor 链上向上游传播upstream propagation不受
`.subscribeOn`/`.publishOn` 线程切换影响——故只要 `process()` 在订阅前用
`contextWrite` 把累加器塞进 Context下游任意调度线程上的 middleware 都能读到。
RC3 `buildAgentStream` 已经 `.contextWrite(c -> c.put(RUNTIME_CONTEXT_KEY, context))`
`.contextWrite(c -> c.put(EVENT_SINK_KEY, sink))`——证明 Context 传播在 RC3 设计中
是首选机制。Task 6.1 的 ChatUsage 累加器走同一路径:
- `ChatUsageMiddleware` 的累加器改为"先看 reactor ContextView没有再回退 ThreadLocal"
(保留 ThreadLocal 回退使 {@code ChatUsageMiddlewareTest} 这种无 Context 的纯单元测试
继续可用——它直接 `bind()` 后调 `onModelCall`,无 reactor Context
- `process()``call(...)``streamEvents(...)` 返回的 `Mono`/`Flux` 上都
`.contextWrite(c -> c.put(USAGE_ACC_KEY, acc))``acc` 是 bind 时创建的同一个实例,
同时被 ThreadLocal 持有(供 HTTP 线程上的 `snapshot()` 读)。
- 累加器仍是带 `synchronized` 的共享对象HTTP 线程读 / boundedElastic 线程写互斥安全。
### (f) 验证
`StreamingBridgeTest` 用 mock 模型回放 `ModelCallEndEvent`(带 ChatUsage+
`AgentResultEvent`,断言 `streamEvents` 路径下 `ctx.getChatUsage()` 仍拿到累计值——
即 reactor Context 方案在流式路径下生效。`ChatUsageMiddlewareTest`ThreadLocal 回退路径)
保持不变、继续 PASS。

View File

@ -1,10 +1,12 @@
package com.yomahub.liteflow.agent.component;
import com.yomahub.liteflow.agent.event.AgentEventBridge;
import com.yomahub.liteflow.agent.exception.AgentConfigException;
import com.yomahub.liteflow.agent.middleware.ChatUsageMiddleware;
import com.yomahub.liteflow.agent.middleware.SkillTrackingMiddleware;
import com.yomahub.liteflow.agent.model.ModelSpec;
import com.yomahub.liteflow.core.NodeComponent;
import com.yomahub.liteflow.flow.FlowEventPublisher;
import com.yomahub.liteflow.property.LiteflowConfigGetter;
import com.yomahub.liteflow.property.agent.AgentConfig;
import com.yomahub.liteflow.slot.Slot;
@ -256,7 +258,8 @@ public abstract class ReActAgentComponent extends NodeComponent {
/* ===== 框架 final 执行体 ===== */
/**
* 端到端非流式执行RC3
* 端到端执行RC3 {@code FlowEvent} 监听者时走流式{@code streamEvents}
* {@link AgentEventBridge} {@code FlowEvent} 发布否则走非流式 {@code call()}
*
* <p>流程
* <ol>
@ -267,14 +270,22 @@ public abstract class ReActAgentComponent extends NodeComponent {
* <li>解析 {@code conversationId} / {@code agentKey}写回 slot并构造对应的
* {@link RuntimeContext}{@code userId=conversationId, sessionId=agentKey}</li>
* <li> {@link ReActAgentContext} runtimeContext挂到 slot attachment
* 使 {@link #ctx()} {@code agent.call(...)} 触发的工具回调内可用</li>
* <li> {@code agent.call(List.of(new UserMessage(userPrompt())), rc).block()}
* 阻塞拿回复 {@link #handleReply(Msg)} 处理</li>
* 使 {@link #ctx()} {@code agent.call(...)}/{@code streamEvents(...)} 触发的
* 工具回调内可用</li>
* <li><b>分流</b>{@link FlowEventPublisher#hasListener} 为真
* {@link AgentEventBridge#streamAndPublish}订阅 {@code streamEvents}映射成
* {@code agent.reasoning/tool_result/summary/result} 等事件发布
* 否则 {@code agent.call(...).block()}非流式</li>
* <li>两条路径的回复都交 {@link #handleReply(Msg)} 处理</li>
* <li>{@code finally} 中摘除 ctx避免跨 invocation 悬挂引用</li>
* </ol>
*
* <p><b>非流式</b>本方法始终用 {@code call(...)}流式{@code streamEvents}桥接
* Task 6.1 单独实现不在本方法范围内
* <p><b>ChatUsage 线程安全findings R-stream</b>真实 vendor 模型下
* {@code ChatUsageMiddleware.onModelCall} reactor 调度线程boundedElastic
* 执行<b>不在</b> {@code bind()} HTTP 线程上 {@code bind()} 返回的累加器
* 同时经 {@link ChatUsageMiddleware#bindToContext} 注入 reactor Context使 middleware
* 能在调度线程上经 {@code deferContextual} 读到两条路径都注入非流式也走同一
* Context 传播机制行为一致
*
* <p>签名保持 {@code final} 1.0 一致
*/
@ -294,10 +305,26 @@ public abstract class ReActAgentComponent extends NodeComponent {
slot.setAttachment(ctxKey(), ctx);
// per-invocation 绑定 ChatUsage / Skill 累加器agent 是单例 process() 复用
// 不能跨 invocation 累加在入口 bind出口 finally unbind findings R5
ChatUsageMiddleware.bind();
// bind() 返回的 ChatUsage 累加器还要注入 reactor Context使流式/非流式路径下
// middleware 在调度线程上都能读到findings R-stream
ChatUsageMiddleware.Accumulator usageAcc = ChatUsageMiddleware.bindAndReturn();
SkillTrackingMiddleware.bind();
try {
Msg reply = agent.call(List.of(new UserMessage(userPrompt())), rc).block();
Msg reply;
if (FlowEventPublisher.hasListener(slot)) {
// 流式streamEvents AgentEventBridge FlowEvent 发布末尾 AGENT_RESULT Msg 交回
Msg userMsg = new UserMessage(userPrompt());
reply = ChatUsageMiddleware.bindToContext(
AgentEventBridge.streamAndPublish(
agent, userMsg, rc, slot,
slot.getChainId(), getNodeId(), slot.getRequestId(), cid),
usageAcc).block();
} else {
// 非流式call().block()
reply = ChatUsageMiddleware.bindToContext(
agent.call(List.of(new UserMessage(userPrompt())), rc),
usageAcc).block();
}
handleReply(reply);
} finally {
ChatUsageMiddleware.unbind();

View File

@ -0,0 +1,191 @@
package com.yomahub.liteflow.agent.event;
import com.yomahub.liteflow.flow.FlowEvent;
import com.yomahub.liteflow.flow.FlowEventPublisher;
import com.yomahub.liteflow.slot.Slot;
import io.agentscope.core.ReActAgent;
import io.agentscope.core.agent.RuntimeContext;
import io.agentscope.core.event.AgentEventType;
import io.agentscope.core.event.AgentResultEvent;
import io.agentscope.core.event.HintBlockEvent;
import io.agentscope.core.event.RequireExternalExecutionEvent;
import io.agentscope.core.event.RequireUserConfirmEvent;
import io.agentscope.core.event.TextBlockDeltaEvent;
import io.agentscope.core.event.ToolResultDataDeltaEvent;
import io.agentscope.core.event.ToolResultEndEvent;
import io.agentscope.core.event.ToolResultTextDeltaEvent;
import io.agentscope.core.message.Msg;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.util.List;
/**
* {@link ReActAgent#streamEvents} 的细粒度 {@link io.agentscope.core.event.AgentEvent}
* 流映射成 LiteFlow {@link FlowEvent} {@link FlowEventPublisher} 发布给当次 {@link Slot}
* 上注册的监听者{@code ExecuteOption.eventListener}
*
* <p>这是 Task 6.1 恢复的流式路径{@code ReActAgentComponent.process()} slot 有监听者时
* {@link #streamAndPublish}否则仍走非流式 {@code agent.call(...)}
*
* <h2>事件映射保对外 type 字符串与 1.0 一致</h2>
* <ul>
* <li>{@link AgentEventType#TEXT_BLOCK_DELTA}reasoning 文本增量
* {@code agent.reasoning}last=false</li>
* <li>{@link AgentEventType#TOOL_RESULT_TEXT_DELTA} /
* {@link AgentEventType#TOOL_RESULT_DATA_DELTA} /
* {@link AgentEventType#TOOL_RESULT_END}
* {@code agent.tool_result}last=falsedelta 透传文本end 透传工具名</li>
* <li>{@link AgentEventType#HINT_BLOCK}RC3 summary/hint 信号
* {@code agent.summary}last=false</li>
* <li>{@link AgentEventType#AGENT_RESULT}最终回复
* {@code agent.result}last=truetext {@link AgentResultEvent#getResult()}
* 文本并把该 Msg 作为 {@link Mono} 的返回值交回 {@code process()}</li>
* <li>{@link AgentEventType#REQUIRE_USER_CONFIRM} /
* {@link AgentEventType#REQUIRE_EXTERNAL_EXECUTION}HITL RC3 已存在
* {@code agent.hitl.confirm} / {@code agent.hitl.external_exec}last=false
* 可选透传data 携带原始事件便于业务侧自定义 HITL 处理</li>
* </ul>
*
* <p>其余事件类型{@code MODEL_CALL_*}{@code TEXT_BLOCK_START/END}
* {@code TOOL_CALL_*}{@code THINKING_BLOCK_*} 不直接映射成 FlowEvent它们是
* 中间态信号对外暴露的"用户可观测事件"语义由 reasoning/tool_result/summary/result 承担
*
* <h2>线程安全</h2>
* {@code streamEvents} emit reactor 调度线程上findings R-stream
* {@link FlowEventPublisher#publish} 是无状态静态方法仅读 slot attachment 调监听者回调
* 监听者实现由调用方经 {@code ExecuteOption.eventListener} 提供自行保证线程安全
* 典型实现用 {@code CopyOnWriteArrayList} 收集事件本桥不持有跨 invocation 状态
*/
public final class AgentEventBridge {
/** HITL要求用户确认透传用可选。 */
public static final String FLOW_EVENT_TYPE_HITL_CONFIRM = "agent.hitl.confirm";
/** HITL要求外部执行透传用可选。 */
public static final String FLOW_EVENT_TYPE_HITL_EXTERNAL_EXEC = "agent.hitl.external_exec";
private AgentEventBridge() {
}
/**
* 流式执行 agent 并把 {@link io.agentscope.core.event.AgentEvent} 桥接成
* {@link FlowEvent} 发布捕获 {@link AgentResultEvent} 的最终 Msg 作为返回值
*
* @param agent 已构建好的 {@link ReActAgent} 单例
* @param userMsg 本次用户输入{@code process()} 构造的 UserMessage
* @param rc 本次调用的 {@link RuntimeContext}{@code userId=conversationId,
* sessionId=agentKey}
* @param slot 当次执行 slot事件经 {@link FlowEventPublisher#publish} 发到其监听者
* @param chainId chain id填入 FlowEvent.chainId
* @param nodeId 组件 nodeId填入 FlowEvent.nodeId
* @param requestId 请求 id填入 FlowEvent.requestId可空
* @param conversationId 会话 id填入 FlowEvent.conversationId
* @return {@link Mono}emit 流完成后携带 {@link AgentResultEvent} 的最终 Msg
* 若流未发 {@code AGENT_RESULT}则携带一个空文本 Msg 兜底
*/
public static Mono<Msg> streamAndPublish(
ReActAgent agent,
Msg userMsg,
RuntimeContext rc,
Slot slot,
String chainId,
String nodeId,
String requestId,
String conversationId) {
return agent.streamEvents(List.of(userMsg), rc)
.doOnNext(event -> publishMapped(slot, event, chainId, nodeId, requestId, conversationId))
.filter(event -> event.getType() == AgentEventType.AGENT_RESULT)
.next()
.map(event -> ((AgentResultEvent) event).getResult())
.switchIfEmpty(Mono.fromSupplier(() -> Msg.builder().textContent("").build()));
}
/** 把单个 {@link io.agentscope.core.event.AgentEvent} 映射并发布成 {@link FlowEvent}。 */
private static void publishMapped(
Slot slot,
io.agentscope.core.event.AgentEvent event,
String chainId,
String nodeId,
String requestId,
String conversationId) {
AgentEventType type = event.getType();
if (type == null) {
return;
}
switch (type) {
case TEXT_BLOCK_DELTA: {
String delta = event instanceof TextBlockDeltaEvent d ? d.getDelta() : null;
publish(slot, com.yomahub.liteflow.agent.component.ReActAgentComponent.FLOW_EVENT_TYPE_REASONING,
delta, false, null, chainId, nodeId, requestId, conversationId);
break;
}
case TOOL_RESULT_TEXT_DELTA: {
String delta = event instanceof ToolResultTextDeltaEvent d ? d.getDelta() : null;
publish(slot, com.yomahub.liteflow.agent.component.ReActAgentComponent.FLOW_EVENT_TYPE_TOOL_RESULT,
delta, false, null, chainId, nodeId, requestId, conversationId);
break;
}
case TOOL_RESULT_DATA_DELTA: {
Object data = event instanceof ToolResultDataDeltaEvent d ? d.getData() : null;
publish(slot, com.yomahub.liteflow.agent.component.ReActAgentComponent.FLOW_EVENT_TYPE_TOOL_RESULT,
null, false, data, chainId, nodeId, requestId, conversationId);
break;
}
case TOOL_RESULT_END: {
String toolName = event instanceof ToolResultEndEvent e ? e.getToolCallName() : null;
publish(slot, com.yomahub.liteflow.agent.component.ReActAgentComponent.FLOW_EVENT_TYPE_TOOL_RESULT,
toolName, false, null, chainId, nodeId, requestId, conversationId);
break;
}
case HINT_BLOCK: {
String hint = event instanceof HintBlockEvent h ? h.getHint() : null;
publish(slot, com.yomahub.liteflow.agent.component.ReActAgentComponent.FLOW_EVENT_TYPE_SUMMARY,
hint, false, null, chainId, nodeId, requestId, conversationId);
break;
}
case AGENT_RESULT: {
Msg result = event instanceof AgentResultEvent r ? r.getResult() : null;
String text = result == null ? null : result.getTextContent();
publish(slot, com.yomahub.liteflow.agent.component.ReActAgentComponent.FLOW_EVENT_TYPE_RESULT,
text, true, result, chainId, nodeId, requestId, conversationId);
break;
}
case REQUIRE_USER_CONFIRM: {
// RequireUserConfirmEvent 透传为 data业务侧自定义 HITL 处理读 data 即可
if (event instanceof RequireUserConfirmEvent) {
publish(slot, FLOW_EVENT_TYPE_HITL_CONFIRM, null, false, event,
chainId, nodeId, requestId, conversationId);
}
break;
}
case REQUIRE_EXTERNAL_EXECUTION: {
// 仅在事件确实是 RequireExternalExecutionEvent 时透传保持 data 类型一致
if (event instanceof RequireExternalExecutionEvent) {
publish(slot, FLOW_EVENT_TYPE_HITL_EXTERNAL_EXEC, null, false, event,
chainId, nodeId, requestId, conversationId);
}
break;
}
default:
// 其余事件类型MODEL_CALL_*TEXT_BLOCK_START/ENDTOOL_CALL_*
// THINKING_BLOCK_*DATA_BLOCK_*AGENT_START/END 不映射成 FlowEvent
break;
}
}
private static void publish(
Slot slot, String type, String text, boolean last, Object data,
String chainId, String nodeId, String requestId, String conversationId) {
FlowEvent event = FlowEvent.builder()
.type(type)
.chainId(chainId)
.nodeId(nodeId)
.requestId(requestId)
.conversationId(conversationId)
.text(text)
.last(last)
.data(data)
.build();
FlowEventPublisher.publish(slot, event);
}
}

View File

@ -8,6 +8,8 @@ import io.agentscope.core.middleware.MiddlewareBase;
import io.agentscope.core.middleware.ModelCallInput;
import io.agentscope.core.model.ChatUsage;
import reactor.core.publisher.Flux;
import reactor.util.context.Context;
import reactor.util.context.ContextView;
import java.util.function.Function;
@ -21,17 +23,29 @@ import java.util.function.Function;
* {@code next.apply(input)} 返回的 {@code Flux<AgentEvent>} 里订阅
* {@link ModelCallEndEvent} 并把当步 usage 累加到一个 per-invocation 累加器
*
* <h2>per-invocation 绑定</h2>
* <h2>per-invocation 绑定 双源reactor Context 优先ThreadLocal 回退</h2>
* {@code ReActAgent} {@code ReactAgentFactory} cmp 子类缓存为单例
* {@code process()} 复用<b>不能</b>把累加器做成实例字段直接累加否则会把上次
* 调用的余量带入下一次故累加器用 {@link ThreadLocal} 持有{@link #bind()}
* {@code process()} 入口调用push 新累加器{@link #unbind()} 在出口 {@code finally}
* 调用pop 并丢弃
* 调用的余量带入下一次
*
* <p>之所以用 ThreadLocal 而非 RuntimeContextRC3-core {@code process()}
* {@code .block()} 同步执行整条 ReAct 循环模型调用在同一线程上完成middleware
* 在该线程上被调用子类若未来引入异步流式Task 6.1需相应把累加器改成随
* reactor {@code Context} 传播RC3-core 不在此范围内
* <p><b>线程模型findings R-stream</b>真实 vendor 模型{@code OpenAIChatModel.stream}
* 实测第 180 {@code .subscribeOn(Schedulers.boundedElastic())}middleware
* {@code onModelCall} 链构造在 reactor 调度线程boundedElastic上执行<b>不在</b>
* {@code process()} 调用 {@code bind()} HTTP 线程上故单纯的 ThreadLocal 在真实模型
* 下读不到累加器usage 丢失修复方案
* <ul>
* <li>{@code process()} 在订阅前用 {@code contextWrite} 把累加器塞进 reactor
* {@link Context}{@link #bindToContext(Flux, Accumulator)} / {@link #USAGE_CONTEXT_KEY}</li>
* <li>{@code onModelCall} {@code Flux.deferContextual} 先读 reactor Context 里的累加器
* 没有例如纯单元测试无 Context才回退 ThreadLocal</li>
* </ul>
* reactor Context reactor 链上向上游传播不受 {@code subscribeOn} 线程切换影响
* 这是 RC3 内部 {@code buildAgentStream} {@code EVENT_SINK_KEY}/{@code RUNTIME_CONTEXT_KEY}
* 的同一机制
*
* <p>累加器仍是带 {@code synchronized} 的共享对象写端boundedElastic 线程的 add
* 读端HTTP 线程的 {@link #snapshot()} {@code ctx.getChatUsage()}互斥可见
* ThreadLocal / Context 仅提供 per-invocation 隔离<b></b>保证单线程访问
*
* <h2>读累计值</h2>
* {@link com.yomahub.liteflow.agent.component.ReActAgentContext#getChatUsage()} 通过
@ -40,6 +54,15 @@ import java.util.function.Function;
*/
public class ChatUsageMiddleware implements MiddlewareBase {
/**
* reactor {@link Context} 上携带 per-invocation {@link Accumulator} key
* {@code process()} {@code call()}/{@code streamEvents()} 返回的 Mono/Flux
* {@code contextWrite(c -> c.put(USAGE_CONTEXT_KEY, acc))} 注入middleware
* {@link #onModelCall} {@code deferContextual} 读取
*/
public static final String USAGE_CONTEXT_KEY =
"io.agentscope.liteflow.ChatUsageMiddleware.accumulator";
/** per-invocation 累加器栈bind push、unbind pop栈结构支持嵌套虽当前 process() 不嵌套)。 */
private static final ThreadLocal<Accumulator> BOUND = new ThreadLocal<>();
@ -49,11 +72,27 @@ public class ChatUsageMiddleware implements MiddlewareBase {
/**
* 在当前线程绑定一个新的 per-invocation 累加器必须在 {@code process()} 入口调用
* 出口 {@link #unbind()} 清零
*
* <p>返回类型保持 {@code void}Task 5.1 既定契约二进制兼容需拿到累加器引用
* 用于 reactor Context 注入的调用方改用 {@link #bindAndReturn()}
*/
public static void bind() {
BOUND.set(new Accumulator());
}
/**
* {@link #bind()} 相同但返回新建的累加器实例 {@code process()} 再通过
* {@link #bindToContext(reactor.core.publisher.Mono, Accumulator)} 注入到 reactor
* Context使 middleware 在调度线程上能读到
*
* @return 本次 invocation 新建的累加器
*/
public static Accumulator bindAndReturn() {
Accumulator acc = new Accumulator();
BOUND.set(acc);
return acc;
}
/**
* 摘除当前线程绑定的累加器必须在 {@code process()} 出口{@code finally}调用
* 避免单例 middleware invocation 累加
@ -66,8 +105,8 @@ public class ChatUsageMiddleware implements MiddlewareBase {
* 返回当前线程累加器截至当前累计的 token 用量 bind 或未观察到任何 usage
* 返回 {@code null}
*
* <p>静态访问累加器本身是 ThreadLocal与具体 middleware 实例无关
* {@link com.yomahub.liteflow.agent.component.ReActAgentContext#getChatUsage()}
* <p>静态访问累加器本身是 ThreadLocal+ reactor Context 双源与具体 middleware
* 实例无关 {@link com.yomahub.liteflow.agent.component.ReActAgentContext#getChatUsage()}
* 可直接读无需持有 middleware 引用
*/
public static ChatUsage snapshot() {
@ -75,40 +114,86 @@ public class ChatUsageMiddleware implements MiddlewareBase {
return acc == null ? null : acc.snapshot();
}
/**
* 把累加器注入 reactor {@link Context}使下游任意调度线程上的 middleware
* {@link #onModelCall} 都能经 {@code deferContextual} 读到
*
* <p>用法{@code process()}
* <pre>{@code
* Accumulator acc = ChatUsageMiddleware.bind();
* Msg reply = ChatUsageMiddleware.bindToContext(
* agent.call(msgs, rc), acc).block();
* }</pre>
*
* @param publisher 要附加 Context reactor call 返回的 Mono / streamEvents 返回的 Flux
* @param acc 本次 invocation 的累加器{@link #bind()} 返回值
* @param <T> Mono/Flux 元素类型
* @return {@link #USAGE_CONTEXT_KEY} 注入的同一源contextWrite 返回新实例
*/
public static <T> reactor.core.publisher.Mono<T> bindToContext(
reactor.core.publisher.Mono<T> publisher, Accumulator acc) {
return acc == null ? publisher : publisher.contextWrite(c -> c.put(USAGE_CONTEXT_KEY, acc));
}
/**
* {@link #bindToContext(reactor.core.publisher.Mono, Accumulator)} Flux 重载
*/
public static <T> Flux<T> bindToContext(Flux<T> publisher, Accumulator acc) {
return acc == null ? publisher : publisher.contextWrite(c -> c.put(USAGE_CONTEXT_KEY, acc));
}
@Override
public Flux<AgentEvent> onModelCall(
Agent agent,
RuntimeContext ctx,
ModelCallInput input,
Function<ModelCallInput, Flux<AgentEvent>> next) {
Flux<AgentEvent> downstream = next.apply(input);
Accumulator acc = BOUND.get();
if (acc == null) {
// bind例如被独立使用 process() 外触发的模型调用不累加透传
return downstream;
}
// 订阅下游流每个 ModelCallEndEvent 累加到当前线程的累加器不影响事件本身
return downstream.doOnNext(event -> {
if (event instanceof ModelCallEndEvent end) {
ChatUsage usage = end.getUsage();
if (usage != null) {
acc.add(usage);
}
// 先在 caller 线程取 ThreadLocal 累加器回退路径纯单元测试或非 process() 触发的模型调用
Accumulator threadLocalAcc = BOUND.get();
return Flux.deferContextual(cv -> {
Accumulator acc = resolveAccumulator(cv, threadLocalAcc);
if (acc == null) {
// 既无 reactor Context 累加器也无 ThreadLocal不累加透传
return next.apply(input);
}
// 订阅下游流每个 ModelCallEndEvent 累加到本次 invocation 的累加器不影响事件本身
// doOnNext 在模型流所在的 reactor 调度线程上执行findings R-stream
// 累加器的 synchronized 保证与 HTTP 线程的 snapshot 互斥可见
return next.apply(input).doOnNext(event -> {
if (event instanceof ModelCallEndEvent end) {
ChatUsage usage = end.getUsage();
if (usage != null) {
acc.add(usage);
}
}
});
});
}
/** reactor Context 里的累加器优先;没有则回退 ThreadLocal兼容无 Context 的调用)。 */
private static Accumulator resolveAccumulator(ContextView cv, Accumulator threadLocalAcc) {
Object fromCtx = cv == null ? null : cv.getOrDefault(USAGE_CONTEXT_KEY, null);
if (fromCtx instanceof Accumulator) {
return (Accumulator) fromCtx;
}
return threadLocalAcc;
}
/* ----- 累加器{@code add}/{@code snapshot} 特意 synchronized-----
* 旧注释写"线程不安全、仅由 ThreadLocal 保证单线程访问"具有误导性实际上
* {@link #add} {@code onModelCall} {@code doOnNext} 回调里被调用该回调
* 运行在<b>模型流所在的 reactor 调度线程</b>它通常与 {@code process()} 调用
* {@code bind()} HTTP 线程<b>不同</b>流可能被 publishOn 切换线程因此两个方法
* 特意加 {@code synchronized}写端流线程的 add与读端HTTP 线程的 snapshot
* {@code ctx.getChatUsage()}之间保证可见性与互斥ThreadLocal 仍提供
* per-invocation 隔离单例 agent invocation 不串<b></b>保证单线程访问
* 运行在<b>模型流所在的 reactor 调度线程</b>boundedElastic {@code process()}
* 调用 {@code bind()} HTTP 线程<b>不同</b>流可能被 subscribeOn/publishOn 切换线程
* 故两个方法特意加 {@code synchronized}写端与读端HTTP 线程的 snapshot
* {@code ctx.getChatUsage()}之间保证可见性与互斥ThreadLocal / reactor Context
* 提供 per-invocation 隔离单例 agent invocation 不串<b></b>保证单线程访问
*/
private static final class Accumulator {
/**
* per-invocation token 用量累加器{@link #bind()} 创建可同时绑到 ThreadLocal
* reactor Context{@code onModelCall} 跨任意调度线程累加{@code snapshot()} HTTP
* 线程读公开为静态嵌套类以便 {@code process()} 持有其引用并 {@link #bindToContext}
*/
public static final class Accumulator {
private int inputTokens;
private int outputTokens;
private double time;

View File

@ -0,0 +1,89 @@
package com.yomahub.liteflow.test.agent.v2;
import com.yomahub.liteflow.agent.component.ReActAgentComponent;
import com.yomahub.liteflow.agent.model.ModelSpec;
import com.yomahub.liteflow.property.agent.AgentConfig;
import io.agentscope.core.model.Model;
import org.springframework.stereotype.Component;
import java.util.concurrent.atomic.AtomicInteger;
/**
* {@code StreamingBridgeTest}Task 6.1用的 ReActAgentComponent 子类
* <ul>
* <li>{@link #buildModel()} escape hatch 返回 {@link StreamingReplyModel}绕开真实 LLM
* 且模型会发出多个 chunk 多条 reasoning 增量+ 末尾 usage</li>
* <li>关闭 shell / workspace 工具最小化 toolkit</li>
* <li>记录 userPrompt / handleReply 调用次数供断言</li>
* </ul>
*
* <p>该组件配合 {@code ExecuteOption.eventListener(...)} 触发
* {@code ReActAgentComponent.process()} 的流式分流有监听者时走
* {@code AgentEventBridge.streamAndPublish}{@code streamEvents}而非 {@code call()}
*/
@Component("streamingBridgeAgent")
public class StreamingBridgeCmp extends ReActAgentComponent {
public static final AtomicInteger USER_PROMPT_COUNT = new AtomicInteger();
public static final AtomicInteger HANDLE_REPLY_COUNT = new AtomicInteger();
public static volatile io.agentscope.core.model.ChatUsage LAST_CHAT_USAGE;
public static void reset() {
USER_PROMPT_COUNT.set(0);
HANDLE_REPLY_COUNT.set(0);
LAST_CHAT_USAGE = null;
}
/** 仅满足抽象方法签名buildModel() 被覆写后这里不会被调用。 */
@Override
@SuppressWarnings("rawtypes")
protected ModelSpec model() {
return new ModelSpec() {
@Override
public Model resolve(AgentConfig c) {
return new StreamingReplyModel("streaming-mock");
}
};
}
@Override
protected Model buildModel() {
return new StreamingReplyModel("streaming-mock");
}
@Override
protected String systemPrompt() {
return "test system prompt for streaming bridge";
}
@Override
protected String userPrompt() {
USER_PROMPT_COUNT.incrementAndGet();
Object reqData = getSlot().getChainReqData(getSlot().getChainId());
return reqData == null ? "hi" : reqData.toString();
}
@Override
protected boolean enableShellTool() {
return false;
}
@Override
protected boolean enableWorkspaceFileTools() {
return false;
}
@Override
protected boolean enableReActLogging() {
return false;
}
@Override
protected void handleReply(io.agentscope.core.message.Msg reply) {
HANDLE_REPLY_COUNT.incrementAndGet();
// handleReply 时快照流式累计的 ChatUsage验证 reactor Context 把累加器传到
// middleware 调度线程后HTTP 线程仍能读到正确的累计值
LAST_CHAT_USAGE = ctx().getChatUsage();
super.handleReply(reply);
}
}

View File

@ -0,0 +1,130 @@
package com.yomahub.liteflow.test.agent.v2;
import com.yomahub.liteflow.agent.component.ReActAgentComponent;
import com.yomahub.liteflow.core.ExecuteOption;
import com.yomahub.liteflow.flow.FlowEvent;
import com.yomahub.liteflow.flow.LiteflowResponse;
import com.yomahub.liteflow.test.agent.support.LiveTestSupport;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.test.context.TestPropertySource;
import javax.annotation.Resource;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.stream.Collectors;
/**
* Task 6.1 端到端流式集成测试验证 {@code ReActAgentComponent.process()}
* {@code ExecuteOption.eventListener} 注册了监听者时走 {@code AgentEventBridge.streamAndPublish}
* {@code streamEvents}路径 {@code TextBlockDeltaEvent} {@code agent.reasoning}
* {@code AgentResultEvent} {@code agent.result}(last=true) 桥接成 {@link FlowEvent} 推给监听者
*
* <p><b>无真实 LLM</b>{@link StreamingBridgeCmp} 覆写 {@code buildModel()} 返回
* {@link StreamingReplyModel}确定性回放两个 TextBlock chunk + 末尾 ChatUsage
* 整个测试不需要任何凭据
*
* <p>断言
* <ul>
* <li>收到<b>多条</b> {@code agent.reasoning} 增量文本拼接还原为 "Hello " + "world"</li>
* <li>末尾收到一条 {@code agent.result}(isLast=true)nodeId 为组件 nodeId</li>
* <li>流式路径下 {@code ctx.getChatUsage()} handleReply 内快照正确累加为
* 100 input / 40 output 验证 reactor Context 把累加器传到 middleware 调度线程
* HTTP 线程仍能读到R-stream 方案 b</li>
* <li>{@code response.isSuccess()} true</li>
* </ul>
*/
@TestPropertySource("classpath:/feature/streamingbridge/application.properties")
@SpringBootTest(classes = StreamingBridgeTest.class)
@EnableAutoConfiguration
@ComponentScan("com.yomahub.liteflow.test.agent.v2")
public class StreamingBridgeTest {
@Resource
private com.yomahub.liteflow.core.FlowExecutor flowExecutor;
@Resource
private com.yomahub.liteflow.property.LiteflowConfig liteflowConfig;
@BeforeEach
public void resetRuntime() {
LiveTestSupport.resetAgentSessionManager();
StreamingBridgeCmp.reset();
}
@Test
public void testStreamingPathPublishesReasoningAndResultEvents() {
List<FlowEvent> events = new CopyOnWriteArrayList<>();
LiteflowResponse response = flowExecutor.execute2Resp(
"streamingBridgeChain", "ping",
ExecuteOption.of().eventListener(events::add));
Assertions.assertTrue(response.isSuccess(),
"chain failed: " + (response.getCause() == null
? "<no cause>"
: toString(response.getCause())));
// 1) 收到多条 agent.reasoning 增量文本拼接还原为 mock 模型的两段回复
List<FlowEvent> reasoning = events.stream()
.filter(e -> ReActAgentComponent.FLOW_EVENT_TYPE_REASONING.equals(e.getType()))
.collect(Collectors.toList());
Assertions.assertFalse(reasoning.isEmpty(),
"stream listener should receive at least one agent.reasoning event");
String joined = reasoning.stream()
.map(FlowEvent::getText)
.reduce("", String::concat);
Assertions.assertEquals(StreamingReplyModel.FULL_REPLY, joined,
"reasoning deltas should reconstruct the model's streamed text");
// 所有 reasoning 事件的 nodeId 都应为该 agent nodeId且非 last
for (FlowEvent e : reasoning) {
Assertions.assertEquals("streamingBridgeAgent", e.getNodeId(),
"reasoning event nodeId must be the agent's nodeId");
Assertions.assertFalse(e.isLast(),
"reasoning deltas must not be marked last");
}
// 2) 末尾收到一条 agent.result(last=true)
List<FlowEvent> results = events.stream()
.filter(e -> ReActAgentComponent.FLOW_EVENT_TYPE_RESULT.equals(e.getType()))
.collect(Collectors.toList());
Assertions.assertEquals(1, results.size(),
"exactly one final agent.result event expected");
FlowEvent finalEvent = results.get(0);
Assertions.assertTrue(finalEvent.isLast(),
"final agent.result must be marked last=true");
Assertions.assertEquals("streamingBridgeAgent", finalEvent.getNodeId(),
"final result event nodeId must be the agent's nodeId");
// 3) handleReply 被调用一次证明流式路径末尾 Msg 也走了 handleReply
Assertions.assertEquals(1, StreamingBridgeCmp.HANDLE_REPLY_COUNT.get(),
"handleReply must be called once with the final streamed Msg");
// 4) 流式路径下 ctx.getChatUsage() 正确累加R-stream 方案 b 生效
io.agentscope.core.model.ChatUsage usage = StreamingBridgeCmp.LAST_CHAT_USAGE;
Assertions.assertNotNull(usage,
"streamed path must accumulate ChatUsage (reactor Context propagation)");
Assertions.assertEquals(StreamingReplyModel.USAGE_INPUT, usage.getInputTokens(),
"accumulated inputTokens must match the model's final chunk usage");
Assertions.assertEquals(StreamingReplyModel.USAGE_OUTPUT, usage.getOutputTokens(),
"accumulated outputTokens must match the model's final chunk usage");
}
private static String toString(Throwable t) {
StringBuilder sb = new StringBuilder();
Throwable cur = t;
while (cur != null) {
sb.append(cur.getClass().getSimpleName()).append(": ").append(cur.getMessage());
cur = cur.getCause();
if (cur != null) {
sb.append(" || caused by: ");
}
}
return sb.toString();
}
}

View File

@ -0,0 +1,73 @@
package com.yomahub.liteflow.test.agent.v2;
import io.agentscope.core.message.ContentBlock;
import io.agentscope.core.message.Msg;
import io.agentscope.core.message.TextBlock;
import io.agentscope.core.model.ChatResponse;
import io.agentscope.core.model.ChatUsage;
import io.agentscope.core.model.GenerateOptions;
import io.agentscope.core.model.ToolSchema;
import reactor.core.publisher.Flux;
import java.util.List;
/**
* 确定性的流式回放{@link io.agentscope.core.model.Model} 实现专供
* {@code AgentEventBridge}Task 6.1流式集成测试在<b>无网络无真实 LLM</b> 的前提下端到端跑通
*
* <p>{@link #stream(List, List, GenerateOptions)} 发出 <b>两个</b> {@link ChatResponse}
* <ol>
* <li>第一个 {@link TextBlock}"Hello "不带 usage</li>
* <li>第二个 {@link TextBlock}"world"+ {@link ChatUsage}inputTokens=100,
* outputTokens=40+ {@code finishReason="stop"}</li>
* </ol>
*
* <p>ReActAgent 把每个 chunk {@link TextBlock} 转成 {@code TextBlockDeltaEvent}
* {@code agent.reasoning}并在模型调用结束时发 {@code ModelCallEndEvent}
* 携带最后 chunk usage+ 最终 {@code AgentResultEvent}携带聚合 Msg
* 故订阅 {@code ExecuteOption.eventListener} 应能收到 2 reasoning 增量 + 1
* 末尾 {@code agent.result}(last=true)且流式路径下 {@code ctx.getChatUsage()} 应为
* 100/40验证 reactor Context 把累加器传到 middleware 调度线程
*
* <p>无状态线程安全不依赖任何凭据
*/
final class StreamingReplyModel implements io.agentscope.core.model.Model {
static final String DELTA_1 = "Hello ";
static final String DELTA_2 = "world";
static final String FULL_REPLY = DELTA_1 + DELTA_2;
static final int USAGE_INPUT = 100;
static final int USAGE_OUTPUT = 40;
private final String modelName;
StreamingReplyModel(String modelName) {
this.modelName = modelName;
}
@Override
public Flux<ChatResponse> stream(
List<Msg> messages, List<ToolSchema> tools, GenerateOptions options) {
ChatResponse chunk1 = ChatResponse.builder()
.id("stream-1-" + System.nanoTime())
.content(List.<ContentBlock>of(TextBlock.builder().text(DELTA_1).build()))
.finishReason("stop")
.build();
ChatResponse chunk2 = ChatResponse.builder()
.id("stream-2-" + System.nanoTime())
.content(List.<ContentBlock>of(TextBlock.builder().text(DELTA_2).build()))
.usage(ChatUsage.builder()
.inputTokens(USAGE_INPUT)
.outputTokens(USAGE_OUTPUT)
.time(0.5)
.build())
.finishReason("stop")
.build();
return Flux.just(chunk1, chunk2);
}
@Override
public String getModelName() {
return modelName;
}
}

View File

@ -0,0 +1,11 @@
liteflow.rule-source=feature/streamingbridge/flow.el.xml
liteflow.print-banner=false
# 最小 agent 配置workspace 指向 tmpStreamingReplyModel 不实际写盘memory=NONE
# shell=disabledskills 关闭。本测试不接触真实 LLM。
liteflow.agent.workspace.root=target/wk/v2_streaming_bridge_test
liteflow.agent.workspace.auto-create=true
liteflow.agent.shell.mode=disabled
liteflow.agent.defaults.max-iterations=3
liteflow.agent.logging.react-enabled=false
liteflow.agent.skills.enabled=false

View File

@ -0,0 +1,6 @@
<?xml version="1.0" encoding="UTF-8"?>
<flow>
<chain name="streamingBridgeChain">
THEN(streamingBridgeAgent);
</chain>
</flow>