harness - 调用与流式请求
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 | 语义 |
|---|---|---|
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():
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 开始前产生,并非每个流都一定存在。ReasonEndEvent、ActionEndEvent、ToolCallEndEvent 都是局部结束事件,只有未包装的父 RunEndEvent 才是 Harness 主请求终态。
子代理正常路径:
父 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:
engine.prompt(query).stream().subscribe(
this::handleEvent,
error -> retryOrRecover(error));
8.2 RunEndEvent.isAbnormal()
部分 ReAct 故障会转成可展示的最终答案。此时仍会收到 RunEndEvent 和 onComplete,但 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_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 订阅 |