Solon v4.0.4

harness - 子代理流式输出适配

</> markdown
2026年7月30日 下午8:33:55

当主代理通过 task / multitask 委派子任务时,子代理内部产生的推理事件(思考、工具调用、结果)默认不会出现在父流中。本文从原理到实践,完整介绍如何让子代理的实时输出透传给调用方。

1、为什么子代理事件需要特殊处理

HarnessEngine.stream() 对外暴露的是一条 Flux<AgentEvent>,这条流由主代理驱动。当主代理调用 taskmultitask 工具时,TaskTalent 内部会为每个子任务创建独立的 ReActAgent 实例,并开启各自的子流。

子流默认是封闭的:它只负责推理,不主动向外暴露事件。若要让子代理的事件出现在父流中,需要满足两个条件:

  1. 父流处于流式模式:即父级 ReActTrace.getOptions().getStreamSink() 不为 null——只有在调用了 .stream() 且链路上注入了 FluxSink 时才成立。直接 .call() 同步调用时永远是同步模式,子代理事件不会传播。
  2. 消费者能识别 TaskWrapEvent:子代理事件被包装成 TaskWrapEvent 注入父流,消费者必须识别并解包这一包装层,才能还原出真实的子代理事件。

2、原理:TaskWrapEvent 与 FluxSink

流式模式下,子代理每产生一个 AgentEventTaskTalent 都会将它包裹在 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()子代理名称(如 explorebash
getTaskDescription()子任务摘要描述
isMultitask()是否来自 multitask 并行任务

并发安全: multitask 会并行运行多个子代理,它们共享同一个 FluxSinkTaskTalent 通过 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 的解包与分发逻辑作为完整示例。