docs(harness): annotate entry orchestration, enums, and exception recovery paths

This commit is contained in:
zhuyongxin
2026-07-30 19:08:03 +08:00
parent b39a625e5b
commit 5b2fb985d9
26 changed files with 639 additions and 6 deletions
@@ -1,11 +1,28 @@
package com.superbiz.agent.harness.application;
/**
* Chat 应用进度状态(SSE status 事件),不是错误码。
*
* <p>只表示「当前走到哪一阶段」,成功/失败结局看 ReleaseOutcome / ChatFailureCode。
*/
public enum ChatApplicationStatus {
/** Intent Router 识别请求类型。 */
ROUTING("正在识别请求类型"),
/** 系统闲聊路径生成回答。 */
SYSTEM_RESPONDING("正在生成回答"),
/** 知识库检索中。 */
KNOWLEDGE_SEARCHING("正在查询知识库"),
/** 知识答案整理中。 */
KNOWLEDGE_ANSWERING("正在整理知识答案"),
/** 诊断 Agent 收集证据 / 写草稿(ReAct 循环中)。 */
DIAGNOSIS_RUNNING("正在收集诊断证据"),
/** Evidence / Semantic 门控或无效 draft 的安全发布阶段。 */
SAFETY_VALIDATING("正在进行安全校验");
private final String message;
@@ -14,6 +31,7 @@ public enum ChatApplicationStatus {
this.message = message;
}
/** 面向用户的简短进度文案。 */
public String message() {
return message;
}
@@ -23,16 +23,34 @@ import java.util.Optional;
import java.util.function.Supplier;
import java.util.regex.Pattern;
/**
* Chat 应用编排(Harness Application 层):一次请求从创建到公开结果的负责人。
*
* <p>调用链:
* <pre>
* Controller → execute()
* → core.startRun() // 建 RunContext 边界
* → router.route() // Intent Router(Spring AI 单次模型调用)
* → executePath(intent) // 按意图分叉
* DIAGNOSIS → DiagnosisChatExecutor(Agent → Release)
* → completePath / persistFinish
* </pre>
*
* <p>它编排路径,但不做业务根因判断;也不把 HTTP/SSE 细节塞进 Core。
*/
public final class ChatApplicationUseCase {
private static final Pattern SAFE_ID = Pattern.compile("[A-Za-z0-9][A-Za-z0-9._-]{0,63}");
/** 执行规则:预算、deadline、取消、终态(first-terminal-wins)。 */
private final DiagnosisHarnessCore core;
private final Supplier<String> sessionIdSupplier;
private final ChatRunStore runStore;
/** 意图路由:只产出 IntentType,不执行诊断。 */
private final IntentRouting router;
private final SystemChatOperation systemChat;
private final KnowledgeQueryOperation knowledgeQuery;
/** 诊断子路径:Agent 收集证据写草稿 + Release 门控发布。 */
private final DiagnosisOperation diagnosis;
private final ObjectMapper objectMapper;
private final DiagnosisTraceRecorder traceRecorder;
@@ -74,12 +92,18 @@ public final class ChatApplicationUseCase {
return execute(request, ChatApplicationObserver.noop());
}
/**
* 一次 Chat 请求的主编排。
*
* <p>observer 通常是 SSE session:onStarted 推 metadata,onStatus 推进度。
*/
public ChatApplicationResult execute(ChatApplicationRequest request,
ChatApplicationObserver observer) {
Objects.requireNonNull(request, "request must not be null");
Objects.requireNonNull(observer, "observer must not be null");
String sessionId = resolveSessionId(request.sessionId());
// 会话级上下文:上一路由结果 + 上一轮诊断摘要(给 Router / Diagnosis 用)
Optional<RoutingHistory> history;
Optional<PreviousTurn> previousTurn;
try {
@@ -90,14 +114,19 @@ public final class ChatApplicationUseCase {
ChatFailureCode.RUN_PERSISTENCE_FAILED,
"无法读取会话上下文,请稍后重试", exception);
}
// ★ 步骤1:创建 Run 边界(runId/deadline/budget/cancel/lifecycle),显式向下传递
RunContext context = core.startRun(sessionId);
long startedNanos = System.nanoTime();
IntentType intent = null;
try {
// ★ 步骤2:落库 RUN 开始 + 对外/对内可观测
persistStart(context, request.query());
traceRecorder.record(TraceAuditEvents.runStarted(context));
observer.onStarted(new CoreRunControl(core, context));
observer.onStarted(new CoreRunControl(core, context)); // SSE metadata + 取消句柄
observer.onStatus(ChatApplicationStatus.ROUTING);
// ★ 步骤3:意图路由——只回答「走哪条应用分支」,不调业务工具
intent = router.route(context, new IntentRouterInput(
request.query(),
history.map(RoutingHistory::intent).orElse(null),
@@ -105,8 +134,11 @@ public final class ChatApplicationUseCase {
traceRecorder.record(TraceAuditEvents.routingDecision(context, intent));
persistIntent(context.runId(), intent);
// ★ 步骤4:按 intent 分叉执行(编排决策点)
PathResult path = executePath(
intent, context, request.query(), previousTurn.orElse(null), observer);
// ★ 步骤5:写入 Run 终态(成功)并持久化公开结果
completePath(context, intent, path);
String safeJson = write(path.content());
persistFinish(context, intent, path.outcome(), safeJson,
@@ -117,6 +149,7 @@ public final class ChatApplicationUseCase {
context.sessionId(), context.runId(), intent, path.outcome(),
path.content().contentType(), path.content());
} catch (RuntimeException exception) {
// 统一失败出口:尽量落终态,再映射成安全的对外失败码
ReleaseOutcome terminal = terminalOutcome(context);
try {
runStore.finish(context, intent, terminal, null, null,
@@ -130,6 +163,12 @@ public final class ChatApplicationUseCase {
}
}
/**
* 路径执行完成后的 Run 终态处理。
*
* <p>诊断预算耗尽且已由子路径产出 FALLBACK 时,不二次 completeSuccess
*(lifecycle 已是 BUDGET_EXHAUSTED)。
*/
private void completePath(RunContext context, IntentType intent, PathResult path) {
if (path.handledBudgetTermination()) {
if (intent != IntentType.DIAGNOSIS
@@ -144,6 +183,12 @@ public final class ChatApplicationUseCase {
core.completeSuccess(context);
}
/**
* 根据 Router 产出的 intent 选择执行分支。
*
* <p>这是 Application 的核心编排决策:Router 只给枚举,分支执行权在这里。
* DIAGNOSIS 继续进入 {@link com.superbiz.agent.harness.application.executor.DiagnosisChatExecutor}。
*/
private PathResult executePath(IntentType intent,
RunContext context,
String query,
@@ -162,6 +207,7 @@ public final class ChatApplicationUseCase {
ReleaseOutcome.SUCCESS, knowledgeQuery.execute(context, query), null, false);
}
case DIAGNOSIS -> {
// 诊断子编排:Agent(ReAct) → Guard/Release,不在本类展开
DiagnosisExecutionResult result = diagnosis.execute(
context, query, previousTurn, observer::onStatus);
yield new PathResult(result.outcome(), result.content(), result.publishedResult(),
@@ -1,11 +1,38 @@
package com.superbiz.agent.harness.application;
/**
* Chat 应用层对外失败码(给 SSE failure / 客户端看的粗粒度原因)。
*
* <p>层级:Application 出口。不要和下列内部码混淆:
* <ul>
* <li>{@code RunState}:Run 内存生命周期终态</li>
* <li>{@code ReleaseOutcome}:诊断发布裁决(SUCCESS/FALLBACK/...)</li>
* <li>{@code FallbackType}:FALLBACK 时的细分原因</li>
* <li>{@code ToolBoundaryErrorCode}:单次工具边界错误</li>
* </ul>
*
* <p>由 {@link ChatApplicationUseCase} 在 catch 中映射,文案对用户安全,不暴露内部堆栈。
*/
public enum ChatFailureCode {
/** Intent Router 不可用或输出无法解析,无法决定走哪条路径。 */
ROUTING_UNAVAILABLE,
/** 系统闲聊分支暂时无法回答。 */
SYSTEM_CHAT_UNAVAILABLE,
/** 知识问答分支暂时无法完成检索/作答。 */
KNOWLEDGE_UNAVAILABLE,
/** 诊断分支整体不可用(非具体 FALLBACK 细分)。 */
DIAGNOSIS_UNAVAILABLE,
/** session/run 读写落库失败(start/intent/finish 等)。 */
RUN_PERSISTENCE_FAILED,
/** Run 被取消(客户端断开、用户取消等),对应 RunState.CANCELLED。 */
RUN_CANCELLED,
/** 未归类的内部失败;兜底码,应尽量少用、并靠 Trace 排查。 */
INTERNAL_FAILURE
}
@@ -23,9 +23,22 @@ import com.superbiz.agent.harness.release.DiagnosisReleaseUseCase;
import java.util.Objects;
import java.util.function.Consumer;
/**
* 诊断路径子编排(Application 在 intent=DIAGNOSIS 时的执行器)。
*
* <p>两段式流水线,本身不实现 ReAct 循环:
* <ol>
* <li>{@link DiagnosisAgentUseCase}:Spring AI Alibaba ReactAgent 收集证据并写 Draft</li>
* <li>{@link DiagnosisReleaseUseCase}:Evidence/Semantic 门控 + 发布 SUCCESS 或 FALLBACK</li>
* </ol>
*
* <p>由 {@code ChatApplicationUseCase.executePath} 在路由完成后调用。
*/
public final class DiagnosisChatExecutor implements DiagnosisOperation {
/** Agent 接入:内部创建 ReactAgent 并 agent.call。 */
private final DiagnosisAgentUseCase diagnosisAgent;
/** 验证与发布:未证明的结论不能 SUCCESS 出口。 */
private final DiagnosisReleaseUseCase releaseUseCase;
private final PublishedResultPolicy publishedPolicy;
private final DiagnosisTraceRecorder traceRecorder;
@@ -46,27 +59,38 @@ public final class DiagnosisChatExecutor implements DiagnosisOperation {
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
}
/**
* 诊断子路径:Agent 出草稿 → Release 裁决对外形态。
*
* @param statusSink 回写 SSE 进度(DIAGNOSIS_RUNNING / SAFETY_VALIDATING)
*/
@Override
public DiagnosisExecutionResult execute(RunContext context,
String query,
PreviousTurn previousTurn,
Consumer<ChatApplicationStatus> statusSink) {
Objects.requireNonNull(statusSink, "statusSink must not be null");
// SSE:诊断 Agent 运行中(内部可能多轮模型 + 工具)
statusSink.accept(ChatApplicationStatus.DIAGNOSIS_RUNNING);
DiagnosisAgentExecution execution;
try {
// ★ 段1:ReAct Agent——规划 tool_call、收证据、产出 DiagnosisDraft
execution = diagnosisAgent.execute(
context, new DiagnosisAgentInput(query, previousTurn));
} catch (DiagnosisAgentOutputException exception) {
// Draft 契约失败:有观察事实可降级 FALLBACK,否则上抛
return recoverInvalidDraft(context, exception, statusSink);
}
// SSE:进入门控(Evidence 结构 + Semantic 语义)
statusSink.accept(ChatApplicationStatus.SAFETY_VALIDATING);
// ★ 段2:验证与发布——决定 SUCCESS 报告还是 SAFE_FALLBACK
DiagnosisReleaseResult released = releaseUseCase.execute(context, query, execution);
if (released.outcome() == ReleaseOutcome.FALLBACK) {
return new DiagnosisExecutionResult(
ReleaseOutcome.FALLBACK,
new FallbackContent(released.fallback()),
null,
// 预算耗尽导致的 FALLBACK:上层 completePath 不再 completeSuccess
execution.stopReason()
== com.superbiz.agent.harness.progress.DiagnosisStopReason.BUDGET_LIMIT_REACHED);
}
@@ -80,22 +104,52 @@ public final class DiagnosisChatExecutor implements DiagnosisOperation {
published);
}
/**
* Agent「已经跑完」但最终文本不是合法 {@code DiagnosisDraft} 时的恢复路径。
*
* <p>触发条件(见 {@link DiagnosisAgentOutputException#isDraftContractFailure()}):
* 空输出 / 非法 JSON / schema 不符。真正的执行崩溃({@code EXECUTION_FAILED})不进恢复,直接再抛。
*
* <p>除「记审计 + 包装 FALLBACK 返回」外,还承担:
* <ol>
* <li><b>失败分类闸门</b>:只有 Draft 契约失败可恢复;其它 Agent 异常保持失败语义上抛</li>
* <li><b>fail-closed 安全门</b>:必须 {@code progress.hasObservedFacts()},
* 否则再抛——禁止在「什么都没查到」时用模板假装一次有依据的降级</li>
* <li><b>丢弃非法 draft 正文</b>:不把空串/烂 JSON/错 schema 送进 Evidence/Semantic Guard,
* 也不可能走出 SUCCESS;发布只基于已验真 progress(canonical 工具事实)</li>
* <li><b>走专用 Release 入口</b>:{@link DiagnosisReleaseUseCase#releaseInvalidDraft},
* 而不是 {@code execute(context, query, execution)}——因为没有可校验的 draft</li>
* <li><b>SSE 阶段对齐</b>:推 {@code SAFETY_VALIDATING},与正常门控路径对外进度一致</li>
* <li><b>固定对外形态</b>:{@code FALLBACK + FallbackContent},{@code publishedResult=null}
* (不落可回放的成功发布快照)</li>
* </ol>
*
* <p>与 {@code DiagnosisAgentUseCase#controlledExecution} 的分工:
* controlledExecution 处理 loop <b>中途</b>被预算/收集收敛打断(常无 draft,转 stopped);
* 本方法处理 loop <b>结束后</b>输出契约失败(有/无 progress 决定 FALLBACK 还是失败)。
*/
private DiagnosisExecutionResult recoverInvalidDraft(
RunContext context,
DiagnosisAgentOutputException exception,
Consumer<ChatApplicationStatus> statusSink) {
// 1) 仅 Draft 契约失败可恢复;EXECUTION_FAILED 等保持原异常
if (!exception.isDraftContractFailure()) {
throw exception;
}
// 2) 是否已有可发布的观察事实(来自工具 canonical,不是模型胡写的 draft)
boolean hasProgress = exception.progress().hasObservedFacts();
// 3) 审计:记录 kind / 输出字节 / 有无 progress,便于区分「模型格式烂」vs「彻底空跑」
traceRecorder.record(TraceAuditEvents.agentDraftInvalid(
context, exception.kind(), exception.outputBytes(), hasProgress));
// 4) 无安全事实 → fail closed,交给 Application 写 FAILED
if (!hasProgress) {
throw exception;
}
// 5) 有事实:对齐 SSE 阶段,走「无合法 draft」专用发布(INSUFFICIENT_EVIDENCE 类 FALLBACK)
statusSink.accept(ChatApplicationStatus.SAFETY_VALIDATING);
DiagnosisReleaseResult released = releaseUseCase.releaseInvalidDraft(
context, exception.progress());
// 6) 对外只给安全 Fallback;非法 draft 正文永不出现在 content 里
return new DiagnosisExecutionResult(
ReleaseOutcome.FALLBACK,
new FallbackContent(released.fallback()),
@@ -29,12 +29,26 @@ import java.util.Objects;
import java.util.Set;
import java.util.function.Consumer;
/**
* 意图路由:判断本轮走 SYSTEM_CHAT / KNOWLEDGE_QUERY / DIAGNOSIS。
*
* <p>在 Harness 中的位置:Application 主编排里的「路由节点」,不是业务执行器。
* <ul>
* <li>使用 Spring AI 的 {@link Prompt} / Message 构造请求</li>
* <li>经 {@link GuardModelCall} 调 ChatModel(单次结构化输出,不是 ReactAgent)</li>
* <li>预算、超时、重试、审计由 Harness 包裹</li>
* </ul>
*
* <p>只返回 {@link IntentType};由 {@code ChatApplicationUseCase.executePath} 决定后续分支。
*/
public final class IntentRouter implements IntentRouting {
/** 输出契约:JSON 只能有 intent 一个字段。 */
private static final Set<String> OUTPUT_FIELDS = Set.of("intent");
private final DiagnosisHarnessCore core;
private final HarnessRetryExecutor retryExecutor;
/** 统一模型调用壳:入账 token、受 RunContext 约束。 */
private final GuardModelCall modelCall;
private final ObjectMapper objectMapper;
private final IntentRouterLimits limits;
@@ -69,6 +83,11 @@ public final class IntentRouter implements IntentRouting {
this.systemPrompt = IntentRouterPrompt.load();
}
/**
* 对当前 query(+ 可选历史 intent/query)做一次意图分类。
*
* @return 仅 IntentType;Application 据此 switch 到对应 Executor
*/
@Override
public IntentType route(RunContext context, IntentRouterInput input) {
Objects.requireNonNull(context, "context must not be null");
@@ -78,11 +97,14 @@ public final class IntentRouter implements IntentRouting {
if (bytes > limits.maxInputBytes()) {
throw new IntentRoutingException(new IllegalArgumentException("Router input exceeded limit"));
}
// 先占 Run 字节预算,再调模型
core.reserveRunBytes(context, bytes);
// Spring AI Prompt:system 规则 + user 侧结构化输入 JSON
Prompt prompt = new Prompt(List.of(
new SystemMessage(systemPrompt), new UserMessage(json)));
long started = System.nanoTime();
try {
// 技术失败可按 policy 重试;取消/预算耗尽不能被重试吞掉
return retryExecutor.execute(
context,
context.retryPolicies().intentRouter(),
@@ -103,6 +125,7 @@ public final class IntentRouter implements IntentRouting {
}
}
/** 严格解析:字段集合必须恰好为 {intent},值必须是 IntentType 枚举名。 */
private IntentType parse(String output) {
JsonNode root;
try {