---
title: "harness - 调用与流式请求"
---


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

### 1、同步调用

```java
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：

```java
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 与请求选项

```java
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 | 语义 |
|---|---|---|
| `RunStartEvent` | `getTrace()` | 本次 ReAct 请求开始 |
| `ContextSizeEvent` | `getContextLength()`、`getMessageCount()`、`getTokenCount()`、`isCompressed()`、压缩前后统计 | 上下文检查或压缩结果；可能先于 Reason |
| `ReasonStartEvent` | `getSystemPrompt()`、`getWorkingMemory()`、`getReasonId()` | 一次 Reason 开始；没有 `getTurnCount()` |
| `ReasonDeltaEvent` | `getChatEvent()`、`getReasonId()`、`isThinking()`、`getText()` | 被选择转发的底层 Chat 事件 |
| `ReasonEndEvent` | `getResponse()`、`getMessage()`、`getThinking()`、`getText()`、`isToolCalls()`、`getToolCalls()`、`getDurationMs()`、`getReasonId()` | 一次模型调用结束，不是整个请求终态 |
| `ActionStartEvent` | `getToolCalls()` | 一批本地工具处理开始 |
| `ToolCallStartEvent` | `getCallId()`、`getToolName()`、`getArgs()`、`getReasonId()` | 单个调用开始 |
| `ToolCallEndEvent` | Start 载荷，加 `getResult()`、`getText()`、`getError()`、`getDurationMs()` | 单个调用处理结束 |
| `ActionEndEvent` | `getTrace()` | 一批工具处理结束 |
| `PlanEvent` | `getPlans()`、`getPlanStage()`、`getPlanIndex()`、`getReasonId()` | 计划创建、推进或修订 |
| `HITLPendingEvent` | `getPendingTasks()`、`getTrace()` | 一批调用等待审核 |
| `HITLDecidedEvent` | `getCallId()`、`getToolName()`、`getArgs()`、`getComment()`、`getDecision()`、三个决策判断方法 | 人工决策已经生效 |
| `RunEndEvent` | `getResponse()`、`getMessage()`、`getTrace()`、`getMetrics()`、`getText()`、`isAbnormal()` | **本次请求的顶层终态** |
| `TaskWrapEvent` | `getRealEvent()`、`getParentRunId()`、任务信息 | Harness 子代理事件包装 |

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

`Metrics` 的耗时 API 是 `getTotalDuration()`，不是 `getDurationMs()`：

```java
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`。

首次挂起：

```text
ReasonEndEvent
  -> HITLPendingEvent
  -> RunEndEvent（当前实现 isAbnormal=true，Session 仍为 pending）
  -> onComplete
```

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

提交决策并恢复：

```java
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);
```

典型恢复顺序：

```text
RunStartEvent（复用原 runId）
  -> HITLDecidedEvent+
  -> ActionStartEvent
  -> ToolCallStartEvent
  -> ToolCallEndEvent
  -> ActionEndEvent
  -> [后续 Reason / Action ...]
  -> RunEndEvent
  -> onComplete
```

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

### 7、正常事件顺序

无工具调用：

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

包含本地工具调用：

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

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

子代理正常路径：

```text
父 ToolCallStartEvent（task / multitask）
  -> TaskWrapEvent（子 RunStartEvent）
  -> TaskWrapEvent（子 Reason / Action / ToolCall 事件）*
  -> TaskWrapEvent（子 RunEndEvent）
  -> 父 ToolCallEndEvent
  -> ...
  -> 父 RunEndEvent（未包装）
```

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

### 8、异常终态、onError 与取消

Agent 层没有统一的 `ErrorEvent` 或 `CancelEvent`，需要区分三种情况。

#### 8.1 Reactor onError

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

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

#### 8.2 RunEndEvent.isAbnormal()

部分 ReAct 故障会转成可展示的最终答案。此时仍会收到 `RunEndEvent` 和 `onComplete`，但 `isAbnormal()` 为 `true`：

```java
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`：

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

subscription.dispose();
```

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

### 9、获取完整流式响应

先过滤顶层 End，再取响应：

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

异步场景保留为 `Mono`：

```java
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(...)`。

```java
@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` 后应关闭连接；用户点击停止时也应关闭：

```javascript
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_DELTA` 和 `hasText()` |
| `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 订阅 |
