harness - 调用与流式请求
engine.prompt(...) 返回的是一个 ReActRequest 接口(与 ReActAgent::prompt 完全一致)。HarnessEngine 内部基于 ReActAgent 驱动,因此流式输出的事件体系覆盖 ReAct 的全部事件,并额外扩展了 Harness 特有的 TaskWrapEvent。
详细事件总览请参考 《agent - 同步与流式响应(call 与 stream)》,本文侧重 Harness 场景的说明与完整示例。
1、同步调用
HarnessEngine engine = HarnessEngine.of(harnessHome)
.sessionProvider(sessionProvider)
.build();
// 同步调用,阻塞直到推理完成
AgentResponse resp = engine.prompt("帮我整理今天的代码改动").call();
System.out.println(resp.getContent());
2、流式调用
HarnessEngine engine = HarnessEngine.of(harnessHome)
.sessionProvider(sessionProvider)
.build();
engine.prompt("帮我整理今天的代码改动")
.stream()
.doOnNext(event -> {
// ================ 运行周期事件 ================
if (event instanceof RunStartEvent e) {
System.out.printf("[运行开始] agent=%s runId=%s%n", e.getAgentName(), e.getRunId());
} else if (event instanceof RunEndEvent e) {
System.out.printf("%n[最终答案] %s%n", e.getContent());
System.out.printf(" Token: %s", e.getMetrics() != null ? e.getMetrics().getTotalTokens() : "N/A");
// ================ 计划事件(PlanReAct 模式) ================
} else if (event instanceof PlanEvent e) {
System.out.printf("[计划] stage=%s 步骤数=%d 当前索引=%d%n",
e.getEvent(), e.getPlans().size(), e.getPlanIndex());
// ================ 思考事件 ================
} else if (event instanceof ReasonStartEvent e) {
System.out.printf("%n=== 第 %d 轮思考 ===%n", e.getTurnCount());
} else if (event instanceof ReasonDeltaEvent e) {
if (e.isThinking()) {
System.out.print("\033[90m" + e.getContent() + "\033[0m"); // 深度思考灰色
} else {
System.out.print(e.getContent());
}
} else if (event instanceof ReasonEndEvent e) {
System.out.printf("%n[思考完成] 耗时=%dms 含工具调用=%s%n",
e.getDurationMs(), e.isToolCalls());
// ================ 动作/工具事件 ================
} else if (event instanceof ActionStartEvent e) {
System.out.printf("%n[动作开始] 本次计划调用 %d 个工具%n", e.getToolCalls().size());
} else if (event instanceof ToolCallStartEvent e) {
System.out.printf(" → 调用工具 %s 参数: %s%n", e.getToolName(), e.getArgs());
} else if (event instanceof ToolCallEndEvent e) {
if (e.isError()) {
System.err.printf(" ✗ %s 失败: %s 耗时=%dms%n",
e.getToolName(), e.getErrorMessage(), e.getDurationMs());
} else {
System.out.printf(" ← %s 返回 耗时=%dms 结果长度=%d%n",
e.getToolName(), e.getDurationMs(),
e.getContent() != null ? e.getContent().length() : 0);
}
} else if (event instanceof ActionEndEvent e) {
System.out.printf("[动作结束]%n");
// ================ Harness 特有:子代理委派事件 ================
} else if (event instanceof TaskWrapEvent e) {
AgentEvent real = e.getRealEvent();
String taskTag = e.isMultitask() ? "multitask" : "task";
System.out.printf(" [子代理 %s] taskId=%s agent=%s desc=%s index=%d → %s%n",
taskTag, e.getTaskId(), e.getTaskAgentName(),
e.getTaskDescription(), e.getTaskIndex(),
real.getClass().getSimpleName());
// ================ 上下文管理 ================
} else if (event instanceof ContextSizeEvent e) {
System.out.printf("[上下文] 消息数=%d Token数=%d 已压缩=%s%n",
e.getMessageCount(), e.getTokenCount(), e.isCompressed());
// ================ HITL 人工审批 ================
} else if (event instanceof HITLPendingEvent e) {
System.out.println("⏸ 等待人工审批,挂起中...");
e.getPendingList().forEach(item ->
System.out.println(" 需审批: " + item.getToolName() + " 参数: " + item.getToolArgs()));
} else if (event instanceof HITLDecidedEvent e) {
System.out.printf("[HITL] 决策=%s 决策者=%s%n",
e.isApproved() ? "通过" : "拒绝", e.getDecidedBy());
}
})
.doOnError(err -> System.err.println("流式错误: " + err.getMessage()))
.blockLast();
3、指定会话与请求选项
通过 session(...) 绑定持久会话(不指定则为临时会话);通过 options(...) 可动态切换模型或指定工作区。
AgentSession session = engine.getSession("user-123");
engine.prompt("帮我整理今天的代码改动")
.session(session) // 绑定持久会话(不指定则为临时会话)
.options(o -> {
// 切换大模型(按模型名;不指定则使用主模型)
o.chatModel(engine.getModelOrMain("deepseek-v4-flash"));
// 动态指定工作区路径(不指定则使用默认工作区)
o.toolContextPut(HarnessEngine.ATTR_CWD, "/path/to/project");
})
.stream()
.doOnNext(event -> { /* 同上 */ })
.blockLast();
4、HarnessEngine 完整流式事件说明
所有事件均继承自 AgentEvent 接口(基类 AbsAgentEvent),提供以下公共字段:
| 公共方法 | 类型 | 说明 |
|---|---|---|
getRunId() | String | 本次推理运行的唯一 ID,一次 stream() 调用对应一个 runId |
getAgentName() | String | 产生该事件的智能体名称(子代理委派时可用于区分来源) |
getSession() | AgentSession | 关联的会话对象 |
getMessage() | ChatMessage / null | 事件携带的消息体,可能为 null(如起止信号) |
getContent() | String | 快捷获取消息体文本内容(无内容时返回空串) |
getMeta() | Map | 自定义元数据,可通过拦截器注入额外信息 |
完整事件表
| 事件类型 | 触发时机 | Harness 场景说明 |
|---|---|---|
RunStartEvent | 推理循环开始 | 每次 stream() 调用首先触发 |
PlanEvent | 计划生成/推进/修正 | 仅在 PlanReAct 模式下触发,用于实时展示 LLM 的任务拆解过程。关键信息:getEvent() → PlanStage(CREATE/PROGRESS/REVISE);getPlans() 获取完整计划列表;getPlanIndex() 当前执行索引 |
ReasonStartEvent | 每轮思考开始 | 每轮 LLM 调用前触发。关键信息:getTurnCount() 当前轮次计数 |
ReasonDeltaEvent | 思考增量输出 | 实时渲染推理过程;深度思考内容建议灰色显示。关键信息:isThinking() 深度思考标志;getReasonId() 推理 ID |
ReasonEndEvent | 每轮思考完成 | 可用于统计每轮思考耗时,或判断本轮是否触发了工具调用。关键信息:getDurationMs() 耗时;getResponse() ChatResponse;isToolCalls() 是否含工具调用;getToolCalls() 工具调用列表;getReasonId() |
ActionStartEvent | 一批工具开始 | Harness 内置大量工具(bash、read、write、glob 等),此事件表明该轮动作阶段开始。关键信息:getToolCalls() → Collection<ToolExchanger> |
ToolCallStartEvent | 单个工具调用前 | 可用于展示"正在执行 xxx"的实时状态。关键信息:getCallId() 调用 ID;getToolName() 工具名;getArgs() 参数 Map(不可修改);getReasonId() |
ToolCallEndEvent | 单个工具调用后 | 含执行结果和耗时;isError() 判断是否失败以便在 UI 上标记红色。关键信息:getToolName();getArgs();getDurationMs() 耗时;isError() 是否失败;getErrorMessage() 错误信息;getObservation() 结果消息 |
ActionEndEvent | 一批工具结束 | 可用于 UI 收起"执行中"动画 |
ContextSizeEvent | 上下文大小检查 | Harness 任务通常上下文较长,可监控是否触发压缩。关键信息:getMessageCount() 消息数;getTokenCount() Token 数;isCompressed() 是否已压缩 |
HITLPendingEvent | 人工审批挂起 | 配置了 HitlStrategy 时,危险操作在此暂停等待确认。关键信息:getPendingList() → List<PendingTask>(含工具名、参数、原因) |
HITLDecidedEvent | 人工审批决策返回 | 记录审批结果。关键信息:isApproved() 通过/拒绝;getDecidedBy() 决策者;getDecidedAt() 决策时间 |
TaskWrapEvent | 子代理委派时 | Harness 特有。当使用 task/multitask Talent 时,子代理内部的所有事件都会通过此事件包裹后传播到主代理流中。通过 getRealEvent() 获取子代理真实事件。关键信息:getRealEvent() 被包裹的真实事件;getTaskId() 任务 ID;getTaskIndex() 任务索引;getTaskAgentName() 子代理名;getTaskDescription() 任务描述;isMultitask() 是否并行 |
RunEndEvent | 推理完成 | 含最终答案、轨迹、Token 消耗指标;isAbnormal() 检查是否异常结束。关键信息:getMetrics() → Metrics(含 totalTokens、promptTokens、completionTokens、durationMs);getTrace() 完整轨迹 |
典型事件顺序
RunStartEvent
└─ PlanEvent (CREATE) ← PlanReAct 模式
└─ 循环:
ReasonStartEvent
├─ ReasonDeltaEvent (多个) ← 流式思考内容
└─ ReasonEndEvent ← 思考完成
ActionStartEvent
├─ ToolCallStartEvent (多个)
├─ ToolCallEndEvent (多个)
└─ ActionEndEvent
ContextSizeEvent (按策略触发)
HITLPendingEvent (按策略触发)
HITLDecidedEvent (按策略触发)
PlanEvent (PROGRESS/REVISE) ← PlanReAct 模式
└─ RunEndEvent
当子代理运行时,
ReasonStartEvent~ActionEndEvent阶段之间可能穿插TaskWrapEvent(包裹子代理的事件)。
关于已废弃的旧事件
4.0.4 版本后以下旧事件名已废弃(@Deprecated),请使用上表对应的新事件:
| 已废弃(别用) | 替代为 |
|---|---|
ThoughtChunk | ReasonEndEvent |
ReasonChunk | ReasonDeltaEvent |
ActionChunk | ToolCallStartEvent |
ObservationChunk | ToolCallEndEvent |
ReActChunk | RunEndEvent |
5、SSE(Server-Sent Events)集成示例
在 Solon Web 环境下,将 Flux<AgentEvent> 转为 SSE 流推送给前端(注意:必须使用 subscribe() 而非 blockLast(),避免阻塞 Web 线程):
@Controller
public class ChatController {
@Inject
private HarnessEngine engine;
@Mapping("/chat/stream")
public SseEmitter chat(String prompt, String sessionId) {
SseEmitter emitter = new SseEmitter(300_000L); // 5 分钟超时
AgentSession session = engine.getSession(sessionId);
engine.prompt(prompt)
.session(session)
.stream()
.doOnNext(event -> {
try {
// ===== 运行周期 =====
if (event instanceof RunStartEvent e) {
emitter.send(SseEmitter.event()
.name("run_start")
.data(Map.of("agentName", e.getAgentName(), "runId", e.getRunId())));
} else if (event instanceof RunEndEvent e) {
emitter.send(SseEmitter.event()
.name("done")
.data(Map.of("content", e.getContent(),
"totalTokens", e.getMetrics() != null
? e.getMetrics().getTotalTokens() : 0)));
emitter.complete();
// ===== 计划 =====
} else if (event instanceof PlanEvent e) {
emitter.send(SseEmitter.event()
.name("plan")
.data(Map.of("stage", e.getEvent().name(),
"plans", e.getPlans(),
"planIndex", e.getPlanIndex())));
// ===== 思考 =====
} else if (event instanceof ReasonDeltaEvent e) {
emitter.send(SseEmitter.event()
.name("delta")
.data(Map.of("content", e.getContent(),
"thinking", e.isThinking())));
} else if (event instanceof ReasonEndEvent e) {
emitter.send(SseEmitter.event()
.name("reason_end")
.data(Map.of("durationMs", e.getDurationMs(),
"hasToolCalls", e.isToolCalls())));
// ===== 工具调用 =====
} else if (event instanceof ToolCallStartEvent e) {
emitter.send(SseEmitter.event()
.name("tool_start")
.data(Map.of("tool", e.getToolName(),
"args", e.getArgs())));
} else if (event instanceof ToolCallEndEvent e) {
emitter.send(SseEmitter.event()
.name("tool_end")
.data(Map.of("tool", e.getToolName(),
"result", e.getContent(),
"error", e.isError(),
"durationMs", e.getDurationMs())));
// ===== Harness 子代理 =====
} else if (event instanceof TaskWrapEvent e) {
emitter.send(SseEmitter.event()
.name("task_wrap")
.data(Map.of("taskId", e.getTaskId(),
"agentName", e.getTaskAgentName(),
"description", e.getTaskDescription(),
"isMultitask", e.isMultitask(),
"realEventType", e.getRealEvent().getClass().getSimpleName())));
// ===== 上下文 =====
} else if (event instanceof ContextSizeEvent e) {
emitter.send(SseEmitter.event()
.name("context_size")
.data(Map.of("messageCount", e.getMessageCount(),
"tokenCount", e.getTokenCount(),
"compressed", e.isCompressed())));
// ===== HITL =====
} else if (event instanceof HITLPendingEvent e) {
emitter.send(SseEmitter.event()
.name("hitl_pending")
.data(Map.of("pendingList", e.getPendingList())));
} else if (event instanceof HITLDecidedEvent e) {
emitter.send(SseEmitter.event()
.name("hitl_decided")
.data(Map.of("approved", e.isApproved(),
"decidedBy", e.getDecidedBy())));
}
} catch (IOException ex) {
emitter.completeWithError(ex);
}
})
.doOnError(emitter::completeWithError)
.subscribe();
return emitter;
}
}
SSE 注意事项: - 必须使用
subscribe()非阻塞订阅,blockLast()会阻塞 Web 线程导致请求超时 - 使用SseEmitter.event().name("xxx")为每种事件指定事件名,前端可用EventSource.addEventListener("xxx", handler)按类型监听 - 设置合理的超时时间(如 300 秒)以适应长时间推理 - 前端收到done事件后应主动关闭 EventSource 连接