harness - 子代理流式输出适配
当主代理通过 task / multitask 委派子任务时,子代理内部产生的推理事件(思考、工具调用、结果)默认不会出现在父流中。本文从原理到实践,完整介绍如何让子代理的实时输出透传给调用方。
1、为什么子代理事件需要特殊处理
HarnessEngine.stream() 对外暴露的是一条 Flux<AgentEvent>,这条流由主代理驱动。当主代理调用 task 或 multitask 工具时,TaskTalent 内部会为每个子任务创建独立的 ReActAgent 实例,并开启各自的子流。
子流默认是封闭的:它只负责推理,不主动向外暴露事件。若要让子代理的事件出现在父流中,需要满足两个条件:
- 父流处于流式模式:即父级
ReActTrace.getOptions().getStreamSink()不为 null——只有在调用了.stream()且链路上注入了FluxSink时才成立。直接.call()同步调用时永远是同步模式,子代理事件不会传播。 - 消费者能识别
TaskWrapEvent:子代理事件被包装成TaskWrapEvent注入父流,消费者必须识别并解包这一包装层,才能还原出真实的子代理事件。
2、原理:TaskWrapEvent 与 FluxSink
流式模式下,子代理每产生一个 AgentEvent,TaskTalent 都会将它包裹在 TaskWrapEvent 中,通过父流的 FluxSink 注入主流:
主代理流(Flux<AgentEvent>)
├── RunStartEvent
├── ReasonDeltaEvent ← 主代理思考
├── ToolCallStartEvent → 触发 task / multitask 工具
│
│ TaskTalent 内部,子代理流运行:
│
│ ├── TaskWrapEvent(ReasonDeltaEvent) ← 子代理思考
│ ├── TaskWrapEvent(ToolCallStartEvent) ← 子代理工具开始
│ ├── TaskWrapEvent(ToolCallEndEvent) ← 子代理工具结束
│ └── TaskWrapEvent(RunEndEvent) ← 子代理推理完成
│
├── ToolCallEndEvent ← task/multitask 工具返回汇总
└── RunEndEvent ← 主代理最终结果
TaskWrapEvent 携带的信息:
| 字段 | 说明 |
|---|---|
getRealEvent() | 被包装的真实子代理事件 |
getParentRunId() | 父会话的 runId |
getTaskId() | 子任务唯一 ID(UUID) |
getTaskAgentName() | 子代理名称(如 explore、bash) |
getTaskDescription() | 子任务摘要描述 |
isMultitask() | 是否来自 multitask 并行任务 |
并发安全: multitask 会并行运行多个子代理,它们共享同一个 FluxSink。TaskTalent 通过 synchronized(sink) 保证多线程并发写入时的串行化:
private void emitTaskChunk(FluxSink<AgentEvent> sink, TaskWrapEvent wrapChunk) {
synchronized (sink) {
if (!sink.isCancelled()) {
sink.next(wrapChunk);
}
}
}
3、如何在消费端适配
消费 .stream() 时,流中会混入 TaskWrapEvent。适配的核心逻辑是:遇到 TaskWrapEvent 就解包,取出真实事件后与主代理事件按统一方式处理,同时保存子任务元信息(agentName、taskId 等)用于来源标识。
3.1 基本适配模式
engine.prompt(prompt)
.session(session)
.stream()
.flatMap(event -> {
// ① 来源元信息(仅子代理事件才有)
String taskAgentName = null;
String taskId = null;
String taskDescription = null;
// ② 解包 TaskWrapEvent
if (event instanceof TaskWrapEvent wrap) {
AgentEvent real = wrap.getRealEvent();
// 只透传感兴趣的事件类型,其余忽略
if (!(real instanceof ReasonDeltaEvent)
&& !(real instanceof ToolCallStartEvent)
&& !(real instanceof ToolCallEndEvent)
&& !(real instanceof RunEndEvent)) {
return Flux.empty();
}
taskAgentName = wrap.getTaskAgentName();
taskId = wrap.getTaskId();
taskDescription = wrap.getTaskDescription();
event = real; // 还原为真实事件
// 子代理流结束:可在此做一次"子任务完成"通知
if (event instanceof RunEndEvent runEnd) {
notifyTaskDone(taskAgentName, taskId, taskDescription, runEnd.isAbnormal());
return Flux.empty();
}
}
// ③ 统一处理主代理 + 子代理的真实事件
MyChunk chunk = null;
if (event instanceof ReasonDeltaEvent e && !e.isToolCalls()) {
chunk = e.getMessage().isThinking()
? MyChunk.ofThinking(e.getContent())
: MyChunk.ofText(e.getContent());
} else if (event instanceof ToolCallStartEvent e) {
chunk = MyChunk.ofToolStart(e.getToolName(), e.getArgs());
} else if (event instanceof ToolCallEndEvent e) {
chunk = MyChunk.ofToolEnd(e.getToolName(), e.getContent());
} else if (event instanceof RunEndEvent e) {
chunk = MyChunk.ofDone(e.getContent());
}
if (chunk == null) return Flux.empty();
// ④ 附加来源标记,供下游区分主代理/子代理
if (taskAgentName != null) {
chunk.setAgentName(taskAgentName);
chunk.setTaskId(taskId);
chunk.setTaskDescription(taskDescription);
}
return Flux.just(chunk);
})
.subscribe(chunk -> sseEmitter.send(chunk));
几个关键点:
task/multitask工具本身的ToolCallStartEvent/ToolCallEndEvent来自主代理,不会被包在TaskWrapEvent里——这两个工具事件通常应在 UI 上过滤掉(它们是调度指令,不是用户可见的操作)。- 子代理的
RunEndEvent包在TaskWrapEvent里,可借此在该子任务完成时立刻通知前端(不必等待整条主流结束),在并行任务场景下尤其有用。 agentName(子代理名,如explore)和taskId是区分多个并行子任务输出来源的关键字段,前端可据此将不同子任务的 chunk 分组展示。
3.2 同步模式下的行为差异
若父流未处于流式模式(直接 .call() 或未注入 FluxSink),TaskTalent 会退回到同步模式:
// 同步模式:blockLast 阻塞等待子代理完成,中间过程不向外暴露
ReActChunk agentChunk = agent.prompt(prompt)
.session(session)
.stream()
.blockLast();
result = agentChunk.getContent();
同步模式下,子代理完整执行后只把最终文本回传给主代理,不产生任何中间事件。主代理的调用方只能看到最终汇总结果,看不到子代理的工具调用过程。
选择依据: - 需要实时展示子代理操作过程 → 必须使用 .stream() 流式模式 - 只关心最终结果、不关心过程 → .call() 同步模式即可,更简单
4、参考示例:WebStreamBuilder
SolonCode 中的 WebStreamBuilder 是一个生产级的完整适配实现,将 AgentEvent(含 TaskWrapEvent)映射为 WebChunk 序列,通过 SSE 推送给 Web 前端。其对子代理流的处理逻辑与上文描述的模式一致,并在此基础上增加了:
- IM 通道同步转发(微信、飞书等)
- HITL(人机交互循环)挂起卡片推送
- 工具事件的来源前缀规则(
agentName/toolName) - 子代理
RunEndEvent时发出"任务完成"信号让前端立即结算 task-group
可直接参考其 buildStreamFlux 方法中对 TaskWrapEvent 的解包与分发逻辑作为完整示例。