Solon v4.1.0

harness - 调用与流式请求

</> markdown
2026年9月7日 上午11:28:13

HarnessEngine.prompt(...) 返回 ReActRequest。在 4.1 API 中,同步调用返回完整 ReActResponse,流式调用返回 Flux<AgentEvent>;Harness 沿用 ReAct 事件模型,并用 TaskWrapEvent 包装 task / multitask 子代理事件。

1、同步调用

HarnessEngine engine = HarnessEngine.of(workspace, harnessHome)
        .sessionProvider(sessionProvider)
        .build();

ReActResponse response = engine
        .prompt("帮我整理今天的代码改动")
        .call();

System.out.println(response.getContent());

call()callAsync() 直接返回完整响应,不向调用方追加顶层 End 事件。只有 stream() 会在正常构造完整响应后追加 RunEndEvent

2、4.1 流式调用

下面示例只使用当前事件的真实 API:

engine.prompt("帮我整理今天的代码改动")
        .stream()
        .subscribe(
                event -> {
                    if (event instanceof RunStartEvent) {
                        RunStartEvent e = (RunStartEvent) event;
                        System.out.println("运行开始: " + e.getRunId());

                    } else if (event instanceof PlanEvent) {
                        PlanEvent e = (PlanEvent) event;
                        System.out.println("计划阶段: " + e.getPlanStage());
                        System.out.println("计划索引: " + e.getPlanIndex());
                        System.out.println("计划列表: " + e.getPlans());

                    } else if (event instanceof ReasonDeltaEvent) {
                        ReasonDeltaEvent e = (ReasonDeltaEvent) event;
                        ChatEvent chatEvent = e.getChatEvent();
                        if (chatEvent == null || !e.hasText()) {
                            return;
                        }

                        if (chatEvent.is(ChatEventType.TEXT_DELTA)) {
                            System.out.print(e.getText());
                        } else if (chatEvent.is(ChatEventType.THINKING_DELTA)) {
                            System.out.print("[思考] " + e.getText());
                        }

                    } else if (event instanceof ReasonEndEvent) {
                        ReasonEndEvent e = (ReasonEndEvent) event;
                        System.out.println("Reason 耗时: " + e.getDurationMs());
                        System.out.println("包含工具调用: " + e.isToolCalls());

                    } else if (event instanceof ToolCallStartEvent) {
                        ToolCallStartEvent e = (ToolCallStartEvent) event;
                        System.out.println("调用工具: " + e.getToolName());
                        System.out.println("参数: " + e.getArgs());

                    } else if (event instanceof ToolCallEndEvent) {
                        ToolCallEndEvent e = (ToolCallEndEvent) event;
                        System.out.println("Observation: " + e.getText());
                        System.out.println("结果消息: " + e.getResult());
                        if (e.getError() != null) {
                            System.err.println("工具异常: " + e.getError().getMessage());
                        }

                    } else if (event instanceof HITLPendingEvent) {
                        HITLPendingEvent e = (HITLPendingEvent) event;
                        for (HITLTask task : e.getPendingTasks()) {
                            System.out.println("待审批: " + task.getCallUuid());
                            System.out.println("工具: " + task.getToolName());
                            System.out.println("参数: " + task.getArgs());
                            System.out.println("原因: " + task.getComment());
                        }

                    } else if (event instanceof HITLDecidedEvent) {
                        HITLDecidedEvent e = (HITLDecidedEvent) event;
                        System.out.println("调用: " + e.getCallId());
                        System.out.println("工具: " + e.getToolName());
                        System.out.println("生效参数: " + e.getArgs());
                        System.out.println("备注: " + e.getComment());
                        System.out.println("批准: " + e.isApproved());
                        System.out.println("拒绝: " + e.isRejected());
                        System.out.println("跳过: " + e.isSkipped());

                    } else if (event instanceof TaskWrapEvent) {
                        TaskWrapEvent e = (TaskWrapEvent) event;
                        AgentEvent child = e.getRealEvent();
                        System.out.println("子任务: " + e.getTaskId());
                        System.out.println("父 runId: " + e.getParentRunId());
                        System.out.println("子事件: " + child.getClass().getSimpleName());

                    } else if (event instanceof RunEndEvent) {
                        RunEndEvent e = (RunEndEvent) event;
                        System.out.println("最终答案: " + e.getText());
                        System.out.println("异常终态: " + e.isAbnormal());
                        System.out.println("总耗时: "
                                + e.getMetrics().getTotalDuration() + "ms");
                    }
                },
                error -> System.err.println("流以 onError 结束: " + error.getMessage()),
                () -> System.out.println("流已完成"));

ReasonDeltaEvent 包装底层 ChatEvent。准确区分正文和思考时,必须检查 getChatEvent()TEXT_DELTA / THINKING_DELTA,并在读取文本前检查 hasText()isThinking() 可用于快速分组,但不能把边界事件当成文本增量。

3、指定 Session 与请求选项

AgentSession session = engine.getSession("user-123");

engine.prompt("帮我整理今天的代码改动")
        .session(session)
        .options(options -> {
            options.chatModel(engine.getModelOrDefInstance("deepseek-v4-flash"));
            options.toolContextPut(HarnessEngine.ATTR_CWD, "work/project-a");
        })
        .stream()
        .subscribe(this::handleEvent, this::handleError);

ReActRequest 是可变且非线程安全的;每次请求应从新的 engine.prompt(...) 开始,不要跨线程复用同一个 Request。

4、AgentEvent 公共 API

所有 Agent 事件都实现 AgentEvent。公共 API 只有:

API含义
getRunId()当前逻辑任务标识
getAgentName()产生事件的 Agent 名称
getSession()当前 AgentSession
getMeta()事件元数据 Map
hasMeta(name)是否存在指定元数据
hasText()当前事件是否有非空文本
getText()事件的文本投影;无文本事件默认返回空串

公共接口没有统一的 getMessage()getContent()、事件类型枚举或事件分组。消息、Trace、Response 等载荷只存在于特定事件类中。消费 Agent 事件应使用 instanceof 或 Reactor 的 ofType(...)

事件、Session 和 Trace 都是运行态对象。向 SSE、WebSocket 或消息队列发送时,应投影为自己的 DTO,而不是直接把事件对象序列化成稳定协议。

5、Harness/ReAct 主要事件

事件主要 API语义
RunStartEventgetTrace()本次 ReAct 请求开始
ContextSizeEventgetContextLength()getMessageCount()getTokenCount()isCompressed()、压缩前后统计上下文检查或压缩结果;可能先于 Reason
ReasonStartEventgetSystemPrompt()getWorkingMemory()getReasonId()一次 Reason 开始;没有 getTurnCount()
ReasonDeltaEventgetChatEvent()getReasonId()isThinking()getText()被选择转发的底层 Chat 事件
ReasonEndEventgetResponse()getMessage()getThinking()getText()isToolCalls()getToolCalls()getDurationMs()getReasonId()一次模型调用结束,不是整个请求终态
ActionStartEventgetToolCalls()一批本地工具处理开始
ToolCallStartEventgetCallId()getToolName()getArgs()getReasonId()单个调用开始
ToolCallEndEventStart 载荷,加 getResult()getText()getError()getDurationMs()单个调用处理结束
ActionEndEventgetTrace()一批工具处理结束
PlanEventgetPlans()getPlanStage()getPlanIndex()getReasonId()计划创建、推进或修订
HITLPendingEventgetPendingTasks()getTrace()一批调用等待审核
HITLDecidedEventgetCallId()getToolName()getArgs()getComment()getDecision()、三个决策判断方法人工决策已经生效
RunEndEventgetResponse()getMessage()getTrace()getMetrics()getText()isAbnormal()本次请求的顶层终态
TaskWrapEventgetRealEvent()getParentRunId()、任务信息Harness 子代理事件包装

ToolCallEndEvent.getResult() 返回写入 ReAct 上下文的 ChatMessage,可能为 nullgetText() 是其文本投影。普通工具函数的参数错误或执行异常常被转换为 observation 文本,此时 getError() 仍可能为 null,所以审计端还要观察 observation 和整个流的 onError

Metrics 的耗时 API 是 getTotalDuration(),不是 getDurationMs()

RunEndEvent end = ...;
Metrics metrics = end.getMetrics();
System.out.println(metrics.getPromptTokens());
System.out.println(metrics.getCompletionTokens());
System.out.println(metrics.getTotalTokens());
System.out.println(metrics.getTotalDuration());

6、runId 与 HITL 恢复

runId 标识由一次非空 Prompt 创建的逻辑任务,不严格等于一次 stream() 订阅。HITL 挂起后,使用同一 Session 和空 Prompt 恢复会复用原 runId,所以同一 runId 可以跨两次请求出现两组 RunStartEvent / RunEndEvent

首次挂起:

ReasonEndEvent
  -> HITLPendingEvent
  -> RunEndEvent(当前实现 isAbnormal=true,Session 仍为 pending)
  -> onComplete

这次不会伪造尚未发生的 Action/Tool 事件。挂起不是 onError;该 RunEndEvent 只结束本次请求。

提交决策并恢复:

AgentSession session = engine.getSession(sessionId);
List<HITLTask> tasks = HITL.getPendingTasks(session);

for (HITLTask task : tasks) {
    HITL.submit(session, task, HITLDecision.approve());
}

// 同一个 Session,无参 prompt() 表示空 Prompt 恢复
engine.prompt()
        .session(session)
        .stream()
        .subscribe(this::handleEvent, this::handleError);

典型恢复顺序:

RunStartEvent(复用原 runId)
  -> HITLDecidedEvent+
  -> ActionStartEvent
  -> ToolCallStartEvent
  -> ToolCallEndEvent
  -> ActionEndEvent
  -> [后续 Reason / Action ...]
  -> RunEndEvent
  -> onComplete

恢复通常复用保存的 lastReasonMessage 并直接进入 Action。单个敏感调用被 reject 时可能直接进入结束,因此 Decided 后不保证一定出现 Action。

7、正常事件顺序

无工具调用:

RunStartEvent
  -> [ContextSizeEvent]
  -> ReasonStartEvent
  -> ReasonDeltaEvent*
  -> ReasonEndEvent
  -> RunEndEvent
  -> onComplete

包含本地工具调用:

RunStartEvent
  -> [ContextSizeEvent]
  -> ReasonStartEvent
  -> ReasonDeltaEvent*
  -> ReasonEndEvent
  -> ActionStartEvent
       -> ToolCallStartEvent
       -> [PlanEvent]              // 计划工具执行时通常位于工具 Start/End 之间
       -> ToolCallEndEvent
       -> ...更多工具(当前顺序处理)
  -> ActionEndEvent
  -> ...下一轮 Reason
  -> RunEndEvent
  -> onComplete

ContextSizeEvent 由压缩拦截器在 Reason 开始前产生,并非每个流都一定存在。ReasonEndEventActionEndEventToolCallEndEvent 都是局部结束事件,只有未包装的父 RunEndEvent 才是 Harness 主请求终态。

子代理正常路径:

父 ToolCallStartEvent(task / multitask)
  -> TaskWrapEvent(子 RunStartEvent)
  -> TaskWrapEvent(子 Reason / Action / ToolCall 事件)*
  -> TaskWrapEvent(子 RunEndEvent)
  -> 父 ToolCallEndEvent
  -> ...
  -> 父 RunEndEvent(未包装)

multitask 的不同 taskId 事件可以交错。关联父子任务应使用 getParentRunId()getTaskId(),不能仅靠 Agent 名称或“最近一次工具事件”推断。

8、异常终态、onError 与取消

Agent 层没有统一的 ErrorEventCancelEvent,需要区分三种情况。

8.1 Reactor onError

未被 Agent 吸收的异常通过 onError 终止,此后不会补发 RunEndEvent

engine.prompt(query).stream().subscribe(
        this::handleEvent,
        error -> retryOrRecover(error));

8.2 RunEndEvent.isAbnormal()

部分 ReAct 故障会转成可展示的最终答案。此时仍会收到 RunEndEventonComplete,但 isAbnormal()true

engine.prompt(query).stream()
        .ofType(RunEndEvent.class)
        .subscribe(end -> {
            if (end.isAbnormal()) {
                showAbnormal(end.getText());
            } else {
                showAnswer(end.getText());
            }
        });

HITL pending 当前也是这种可恢复的异常终态,必须结合 Pending 事件或 Session 状态解释。

8.3 主动取消

取消订阅后不再投递事件,也不会补发 RunEndEvent

Disposable subscription = engine.prompt(query)
        .stream()
        .subscribe(this::handleEvent, this::handleError);

subscription.dispose();

ReActRequest 会在 Reactor cancel 时尝试中断执行线程,但底层阻塞调用何时停止取决于它对线程中断的响应。

9、获取完整流式响应

先过滤顶层 End,再取响应:

ReActResponse response = engine.prompt(query)
        .stream()
        .ofType(RunEndEvent.class)
        .map(RunEndEvent::getResponse)
        .blockFirst();

异步场景保留为 Mono

Mono<ReActResponse> response = engine.prompt(query)
        .stream()
        .ofType(RunEndEvent.class)
        .map(RunEndEvent::getResponse)
        .next();

不要对未过滤的事件流直接调用 blockFirst(),否则通常拿到 RunStartEvent。取消可能返回空结果;onError 会传播异常。

10、Solon SSE 集成

SSE 层应把事件投影为 DTO,并在连接完成、超时或出错时取消 Reactor 订阅。当前 Solon API 使用 new SseEvent()emitter.error(error),不存在 SseEmitter.event()completeWithError(...)

@Controller
public class ChatController {
    @Inject
    private HarnessEngine engine;

    @Mapping("/chat/stream")
    public SseEmitter chat(String prompt, String sessionId) {
        final SseEmitter emitter = new SseEmitter(300_000L);
        final AtomicReference<Disposable> subscription = new AtomicReference<>();

        Runnable cancel = () -> {
            Disposable disposable = subscription.getAndSet(null);
            if (disposable != null) {
                disposable.dispose();
            }
        };

        emitter.onCompletion(cancel);
        emitter.onTimeout(() -> {
            cancel.run();
            emitter.complete();
        });
        emitter.onError(error -> cancel.run());

        AgentSession session = engine.getSession(sessionId);
        Disposable disposable = engine.prompt(prompt)
                .session(session)
                .stream()
                .subscribe(
                        event -> sendEvent(emitter, event),
                        error -> {
                            cancel.run();
                            emitter.error(error);
                        });

        subscription.set(disposable);
        return emitter;
    }

    private void sendEvent(SseEmitter emitter, AgentEvent event) {
        try {
            if (event instanceof ReasonDeltaEvent) {
                ReasonDeltaEvent delta = (ReasonDeltaEvent) event;
                ChatEvent chatEvent = delta.getChatEvent();
                if (chatEvent == null || !delta.hasText()) {
                    return;
                }

                if (chatEvent.is(ChatEventType.TEXT_DELTA)) {
                    send(emitter, "text_delta", dto(
                            "runId", delta.getRunId(),
                            "reasonId", delta.getReasonId(),
                            "text", delta.getText()));
                } else if (chatEvent.is(ChatEventType.THINKING_DELTA)) {
                    send(emitter, "thinking_delta", dto(
                            "runId", delta.getRunId(),
                            "reasonId", delta.getReasonId(),
                            "text", delta.getText()));
                }

            } else if (event instanceof PlanEvent) {
                PlanEvent plan = (PlanEvent) event;
                send(emitter, "plan", dto(
                        "stage", plan.getPlanStage().name(),
                        "plans", plan.getPlans(),
                        "planIndex", plan.getPlanIndex()));

            } else if (event instanceof ToolCallStartEvent) {
                ToolCallStartEvent start = (ToolCallStartEvent) event;
                send(emitter, "tool_start", dto(
                        "callId", start.getCallId(),
                        "toolName", start.getToolName(),
                        "args", start.getArgs()));

            } else if (event instanceof ToolCallEndEvent) {
                ToolCallEndEvent end = (ToolCallEndEvent) event;
                send(emitter, "tool_end", dto(
                        "callId", end.getCallId(),
                        "toolName", end.getToolName(),
                        "text", end.getText(),
                        "error", end.getError() == null
                                ? null : end.getError().getMessage(),
                        "durationMs", end.getDurationMs()));

            } else if (event instanceof HITLPendingEvent) {
                HITLPendingEvent pending = (HITLPendingEvent) event;
                List<Map<String, Object>> tasks = new ArrayList<>();
                for (HITLTask task : pending.getPendingTasks()) {
                    tasks.add(dto(
                            "callUuid", task.getCallUuid(),
                            "toolName", task.getToolName(),
                            "args", task.getArgs(),
                            "comment", task.getComment()));
                }
                send(emitter, "hitl_pending", dto("tasks", tasks));

            } else if (event instanceof HITLDecidedEvent) {
                HITLDecidedEvent decided = (HITLDecidedEvent) event;
                send(emitter, "hitl_decided", dto(
                        "callId", decided.getCallId(),
                        "toolName", decided.getToolName(),
                        "args", decided.getArgs(),
                        "comment", decided.getComment(),
                        "approved", decided.isApproved(),
                        "rejected", decided.isRejected(),
                        "skipped", decided.isSkipped()));

            } else if (event instanceof TaskWrapEvent) {
                TaskWrapEvent wrapper = (TaskWrapEvent) event;
                AgentEvent child = wrapper.getRealEvent();
                send(emitter, "task_event", dto(
                        "parentRunId", wrapper.getParentRunId(),
                        "taskId", wrapper.getTaskId(),
                        "taskIndex", wrapper.getTaskIndex(),
                        "taskAgentName", wrapper.getTaskAgentName(),
                        "eventType", child.getClass().getSimpleName(),
                        "text", child.getText()));

            } else if (event instanceof RunEndEvent) {
                RunEndEvent end = (RunEndEvent) event;
                send(emitter, "done", dto(
                        "runId", end.getRunId(),
                        "text", end.getText(),
                        "abnormal", end.isAbnormal(),
                        "totalTokens", end.getMetrics().getTotalTokens(),
                        "totalDuration", end.getMetrics().getTotalDuration()));
                emitter.complete();
            }
        } catch (IOException error) {
            emitter.error(error);
        }
    }

    private void send(SseEmitter emitter, String name, Object data)
            throws IOException {
        emitter.send(new SseEvent()
                .name(name)
                .data(ONode.serialize(data)));
    }

    private Map<String, Object> dto(Object... values) {
        Map<String, Object> result = new LinkedHashMap<>();
        for (int i = 0; i + 1 < values.length; i += 2) {
            result.put(String.valueOf(values[i]), values[i + 1]);
        }
        return result;
    }
}

服务端 emitter.complete() 会触发完成回调并释放订阅;客户端断开、超时和 SSE 错误也必须调用 Disposable.dispose(),否则 Agent 可能继续运行。主动取消不会产生 RunEndEvent

前端收到 done 后应关闭连接;用户点击停止时也应关闭:

const source = new EventSource("/chat/stream?prompt=hello&sessionId=user-1");

source.addEventListener("text_delta", event => {
    const data = JSON.parse(event.data);
    appendText(data.text);
});

source.addEventListener("done", event => {
    renderDone(JSON.parse(event.data));
    source.close();
});

function cancel() {
    source.close();
}

11、4.1 迁移表

旧写法或旧思路4.1 写法
AgentChunk / 旧 Chunk 类型AgentEvent 与具体 Event 类
公共 getContent()公共文本投影使用 getText()
PlanEvent.getEvent()getPlanStage()
ReasonStartEvent.getTurnCount()已删除;用 reasonId 关联 Reason,不臆造轮次字段
所有 Reason Delta 都当正文检查 getChatEvent()TEXT_DELTA / THINKING_DELTAhasText()
ToolCallEndEvent.isError() / getErrorMessage() / getContent()getError() / getResult() / getText()
HITLPendingEvent.getPendingList()getPendingTasks()
HITLTask.getToolArgs()getArgs();调用主键为 getCallUuid()
Decided 的 getDecidedBy() / getDecidedAt()使用 getCallId()getToolName()getArgs()getComment()getDecision() 和三个状态判断方法
Metrics getDurationMs()getTotalDuration()
ReasonEndEvent 当最终结果使用 RunEndEvent.getResponse()
完成只看 onComplete同时处理 RunEndEvent.isAbnormal()、Reactor onError 与 cancel
SSE 直接序列化事件显式投影 DTO,并在连接结束时取消 Reactor 订阅