344 lines
16 KiB
Java
344 lines
16 KiB
Java
package com.superbiz.agent.harness.application;
|
||
|
||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||
import com.superbiz.agent.harness.application.persistence.ChatRunStore;
|
||
import com.superbiz.agent.harness.audit.DiagnosisTraceRecorder;
|
||
import com.superbiz.agent.harness.audit.TraceAuditEvents;
|
||
import com.superbiz.agent.harness.application.persistence.RoutingHistory;
|
||
import com.superbiz.agent.harness.application.routing.IntentRouterInput;
|
||
import com.superbiz.agent.harness.application.routing.IntentRoutingException;
|
||
import com.superbiz.agent.harness.contract.IntentType;
|
||
import com.superbiz.agent.harness.contract.PreviousTurn;
|
||
import com.superbiz.agent.harness.contract.PublishedResult;
|
||
import com.superbiz.agent.harness.contract.ReleaseOutcome;
|
||
import com.superbiz.agent.harness.core.DiagnosisHarnessCore;
|
||
import com.superbiz.agent.harness.core.RunCancellationReason;
|
||
import com.superbiz.agent.harness.core.RunContext;
|
||
import com.superbiz.agent.harness.core.RunState;
|
||
import com.superbiz.agent.harness.retry.RetryExecutionException;
|
||
|
||
import java.util.Objects;
|
||
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;
|
||
|
||
public ChatApplicationUseCase(DiagnosisHarnessCore core,
|
||
Supplier<String> sessionIdSupplier,
|
||
ChatRunStore runStore,
|
||
IntentRouting router,
|
||
SystemChatOperation systemChat,
|
||
KnowledgeQueryOperation knowledgeQuery,
|
||
DiagnosisOperation diagnosis,
|
||
ObjectMapper objectMapper) {
|
||
this(core, sessionIdSupplier, runStore, router, systemChat, knowledgeQuery,
|
||
diagnosis, objectMapper, DiagnosisTraceRecorder.noop());
|
||
}
|
||
|
||
public ChatApplicationUseCase(DiagnosisHarnessCore core,
|
||
Supplier<String> sessionIdSupplier,
|
||
ChatRunStore runStore,
|
||
IntentRouting router,
|
||
SystemChatOperation systemChat,
|
||
KnowledgeQueryOperation knowledgeQuery,
|
||
DiagnosisOperation diagnosis,
|
||
ObjectMapper objectMapper,
|
||
DiagnosisTraceRecorder traceRecorder) {
|
||
this.core = Objects.requireNonNull(core, "core must not be null");
|
||
this.sessionIdSupplier = Objects.requireNonNull(
|
||
sessionIdSupplier, "sessionIdSupplier must not be null");
|
||
this.runStore = Objects.requireNonNull(runStore, "runStore must not be null");
|
||
this.router = Objects.requireNonNull(router, "router must not be null");
|
||
this.systemChat = Objects.requireNonNull(systemChat, "systemChat must not be null");
|
||
this.knowledgeQuery = Objects.requireNonNull(knowledgeQuery, "knowledgeQuery must not be null");
|
||
this.diagnosis = Objects.requireNonNull(diagnosis, "diagnosis must not be null");
|
||
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
|
||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||
}
|
||
|
||
public ChatApplicationResult execute(ChatApplicationRequest request) {
|
||
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 {
|
||
history = runStore.findLatestRoutingHistory(sessionId);
|
||
previousTurn = runStore.findPreviousTurn(sessionId);
|
||
} catch (RuntimeException exception) {
|
||
throw new ChatApplicationException(
|
||
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)); // SSE metadata + 取消句柄
|
||
observer.onStatus(ChatApplicationStatus.ROUTING);
|
||
|
||
// ★ 步骤3:意图路由——只回答「走哪条应用分支」,不调业务工具
|
||
intent = router.route(context, new IntentRouterInput(
|
||
request.query(),
|
||
history.map(RoutingHistory::intent).orElse(null),
|
||
history.map(RoutingHistory::userQuery).orElse(null)));
|
||
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,
|
||
path.publishedResult(), durationMillis(startedNanos));
|
||
traceRecorder.record(TraceAuditEvents.runFinished(
|
||
context, intent, path.outcome(), durationMillis(startedNanos)));
|
||
return new ChatApplicationResult(
|
||
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,
|
||
durationMillis(startedNanos));
|
||
} catch (RuntimeException persistenceFailure) {
|
||
exception.addSuppressed(persistenceFailure);
|
||
}
|
||
traceRecorder.record(TraceAuditEvents.runFinished(
|
||
context, intent, terminal, durationMillis(startedNanos)));
|
||
throw safeFailure(exception, terminal);
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 路径执行完成后的 Run 终态处理。
|
||
*
|
||
* <p>诊断预算耗尽且已由子路径产出 FALLBACK 时,不二次 completeSuccess
|
||
*(lifecycle 已是 BUDGET_EXHAUSTED)。
|
||
*/
|
||
private void completePath(RunContext context, IntentType intent, PathResult path) {
|
||
if (path.handledBudgetTermination()) {
|
||
if (intent != IntentType.DIAGNOSIS
|
||
|| path.outcome() != ReleaseOutcome.FALLBACK
|
||
|| path.publishedResult() != null
|
||
|| context.lifecycle().state() != RunState.BUDGET_EXHAUSTED) {
|
||
throw new IllegalStateException("Invalid handled Diagnosis budget fallback");
|
||
}
|
||
return;
|
||
}
|
||
core.checkActive(context);
|
||
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,
|
||
PreviousTurn previousTurn,
|
||
ChatApplicationObserver observer) {
|
||
return switch (intent) {
|
||
case SYSTEM_CHAT -> {
|
||
observer.onStatus(ChatApplicationStatus.SYSTEM_RESPONDING);
|
||
yield new PathResult(
|
||
ReleaseOutcome.SUCCESS, systemChat.execute(context, query), null, false);
|
||
}
|
||
case KNOWLEDGE_QUERY -> {
|
||
observer.onStatus(ChatApplicationStatus.KNOWLEDGE_SEARCHING);
|
||
observer.onStatus(ChatApplicationStatus.KNOWLEDGE_ANSWERING);
|
||
yield new PathResult(
|
||
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(),
|
||
result.handledBudgetTermination());
|
||
}
|
||
};
|
||
}
|
||
|
||
private void persistStart(RunContext context, String query) {
|
||
try {
|
||
runStore.start(context, query);
|
||
} catch (RuntimeException exception) {
|
||
throw new ChatApplicationException(
|
||
ChatFailureCode.RUN_PERSISTENCE_FAILED,
|
||
"无法创建诊断运行记录,请稍后重试", exception);
|
||
}
|
||
}
|
||
|
||
private void persistIntent(String runId, IntentType intent) {
|
||
try {
|
||
runStore.markIntent(runId, intent);
|
||
} catch (RuntimeException exception) {
|
||
throw new ChatApplicationException(
|
||
ChatFailureCode.RUN_PERSISTENCE_FAILED,
|
||
"无法记录请求类型,请稍后重试", exception);
|
||
}
|
||
}
|
||
|
||
private void persistFinish(RunContext context,
|
||
IntentType intent,
|
||
ReleaseOutcome outcome,
|
||
String safeJson,
|
||
PublishedResult publishedResult,
|
||
int durationMs) {
|
||
try {
|
||
runStore.finish(context, intent, outcome, safeJson, publishedResult, durationMs);
|
||
} catch (RuntimeException exception) {
|
||
throw new ChatApplicationException(
|
||
ChatFailureCode.RUN_PERSISTENCE_FAILED,
|
||
"无法记录运行终态,请稍后重试", exception);
|
||
}
|
||
}
|
||
|
||
private ReleaseOutcome terminalOutcome(RunContext context) {
|
||
RunState state = context.lifecycle().state();
|
||
if (state == RunState.CANCELLED) {
|
||
return ReleaseOutcome.CANCELLED;
|
||
}
|
||
if (!state.isTerminal()) {
|
||
core.completeFailure(context, "APPLICATION_EXECUTION_FAILED");
|
||
}
|
||
return context.lifecycle().state() == RunState.CANCELLED
|
||
? ReleaseOutcome.CANCELLED : ReleaseOutcome.FAILED;
|
||
}
|
||
|
||
private ChatApplicationException safeFailure(RuntimeException exception,
|
||
ReleaseOutcome terminal) {
|
||
if (terminal == ReleaseOutcome.CANCELLED) {
|
||
return new ChatApplicationException(
|
||
ChatFailureCode.RUN_CANCELLED, "Chat Run was cancelled", exception);
|
||
}
|
||
if (exception instanceof IntentRoutingException) {
|
||
return new ChatApplicationException(
|
||
ChatFailureCode.ROUTING_UNAVAILABLE,
|
||
"当前暂时无法识别请求类型,请稍后重试", exception);
|
||
}
|
||
if (exception instanceof ChatApplicationException applicationFailure) {
|
||
return applicationFailure;
|
||
}
|
||
if (exception instanceof RetryExecutionException retry
|
||
&& retry.failure() == com.superbiz.agent.harness.retry.RetryFailure.CANCELLED) {
|
||
return new ChatApplicationException(
|
||
ChatFailureCode.RUN_CANCELLED, "Chat Run was cancelled", exception);
|
||
}
|
||
return new ChatApplicationException(
|
||
ChatFailureCode.INTERNAL_FAILURE, "当前暂时无法处理该请求,请稍后重试", exception);
|
||
}
|
||
|
||
private String resolveSessionId(String requested) {
|
||
String value = requested == null ? sessionIdSupplier.get() : requested;
|
||
if (value == null || !SAFE_ID.matcher(value).matches()) {
|
||
throw new IllegalArgumentException("sessionId is invalid");
|
||
}
|
||
return value;
|
||
}
|
||
|
||
private String write(ChatApplicationContent content) {
|
||
try {
|
||
return objectMapper.writeValueAsString(content);
|
||
} catch (JsonProcessingException exception) {
|
||
throw new ChatApplicationException(
|
||
ChatFailureCode.INTERNAL_FAILURE, "Public content is not serializable", exception);
|
||
}
|
||
}
|
||
|
||
private static int durationMillis(long startedNanos) {
|
||
long millis = Math.max(0L, (System.nanoTime() - startedNanos) / 1_000_000L);
|
||
return millis >= Integer.MAX_VALUE ? Integer.MAX_VALUE : (int) millis;
|
||
}
|
||
|
||
private record PathResult(
|
||
ReleaseOutcome outcome,
|
||
ChatApplicationContent content,
|
||
PublishedResult publishedResult,
|
||
boolean handledBudgetTermination) {
|
||
}
|
||
|
||
private static final class CoreRunControl implements ChatRunControl {
|
||
|
||
private final DiagnosisHarnessCore core;
|
||
private final RunContext context;
|
||
|
||
private CoreRunControl(DiagnosisHarnessCore core, RunContext context) {
|
||
this.core = core;
|
||
this.context = context;
|
||
}
|
||
|
||
@Override
|
||
public String sessionId() {
|
||
return context.sessionId();
|
||
}
|
||
|
||
@Override
|
||
public String runId() {
|
||
return context.runId();
|
||
}
|
||
|
||
@Override
|
||
public boolean cancelClientDisconnect() {
|
||
return core.cancel(context, RunCancellationReason.CLIENT_DISCONNECTED);
|
||
}
|
||
}
|
||
}
|