912 lines
39 KiB
Java
912 lines
39 KiB
Java
package com.superbiz.agent.service;
|
||
|
||
import com.alibaba.cloud.ai.graph.OverAllState;
|
||
import com.alibaba.cloud.ai.graph.RunnableConfig;
|
||
import com.alibaba.cloud.ai.graph.agent.ReactAgent;
|
||
import com.alibaba.cloud.ai.graph.agent.flow.agent.SequentialAgent;
|
||
import com.alibaba.cloud.ai.graph.agent.hook.Hook;
|
||
import com.alibaba.cloud.ai.graph.agent.hook.skills.SkillsAgentHook;
|
||
import com.alibaba.cloud.ai.graph.exception.GraphRunnerException;
|
||
import com.alibaba.cloud.ai.graph.skills.registry.SkillRegistry;
|
||
import com.fasterxml.jackson.databind.JsonNode;
|
||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||
import com.superbiz.agent.agent.tool.DateTimeTools;
|
||
import com.superbiz.agent.agent.tool.InternalDocsTools;
|
||
import com.superbiz.agent.agent.tool.QueryLogsTools;
|
||
import com.superbiz.agent.agent.tool.QueryMetricsTools;
|
||
import com.superbiz.agent.domain.entity.DiagnosisSession;
|
||
import com.superbiz.agent.hook.AgentLoggingHook;
|
||
import com.superbiz.agent.hook.PlannerSkillMetadataHook;
|
||
import com.superbiz.agent.hook.TokenTrackingChatModel;
|
||
import com.superbiz.agent.hook.TokenUsageHolder;
|
||
import com.superbiz.agent.hook.VerifierInputHook;
|
||
import com.superbiz.agent.repository.AgentStepRepository;
|
||
import com.superbiz.agent.repository.DiagnosisSessionRepository;
|
||
import com.superbiz.agent.repository.ToolInvocationRepository;
|
||
import com.superbiz.agent.tool.LookupKnowledgeTool;
|
||
import com.superbiz.agent.tool.RetrievedDocTracker;
|
||
import com.superbiz.agent.util.QuestionComplexity;
|
||
import com.superbiz.agent.util.SessionContextHolder;
|
||
import com.superbiz.agent.util.VerifierContextHolder;
|
||
|
||
import jakarta.annotation.PostConstruct;
|
||
import org.slf4j.Logger;
|
||
import org.slf4j.LoggerFactory;
|
||
import org.springframework.ai.chat.messages.AssistantMessage;
|
||
import org.springframework.ai.chat.model.ChatModel;
|
||
import org.springframework.ai.tool.ToolCallback;
|
||
import org.springframework.ai.tool.ToolCallbackProvider;
|
||
import org.springframework.beans.factory.annotation.Autowired;
|
||
import org.springframework.beans.factory.annotation.Value;
|
||
import org.springframework.core.io.ClassPathResource;
|
||
import org.springframework.stereotype.Service;
|
||
|
||
import java.io.IOException;
|
||
import java.nio.charset.StandardCharsets;
|
||
import java.util.ArrayList;
|
||
import java.util.LinkedHashMap;
|
||
import java.util.List;
|
||
import java.util.Map;
|
||
import java.util.Optional;
|
||
import java.util.UUID;
|
||
|
||
/**
|
||
* 聊天服务
|
||
* 封装 ReactAgent 对话的公共逻辑,包括模型创建、系统提示词构建、Agent 配置等
|
||
*/
|
||
@Service
|
||
public class ChatService {
|
||
|
||
private static final Logger logger = LoggerFactory.getLogger(ChatService.class);
|
||
private static final String LOW_CONFID_DISCLAIMER = "以下结论基于当前已获取证据,仍存在部分证据缺口,请谨慎参考。";
|
||
private static final String DEGRADED_PREFIX = "当前无法基于已获取证据生成可靠结论,建议人工介入。";
|
||
|
||
/** 封装 answer + 后端生成的 sessionId,用于 feedback 关联 */
|
||
public record ChatResult(String answer, String sessionId) {}
|
||
|
||
@Autowired
|
||
private InternalDocsTools internalDocsTools;
|
||
|
||
@Autowired
|
||
private DateTimeTools dateTimeTools;
|
||
|
||
@Autowired
|
||
private QueryMetricsTools queryMetricsTools;
|
||
|
||
@Autowired(required = false) // Mock 模式下才注册,所以设置为 optional,真实环境通过mcp配置注入
|
||
private QueryLogsTools queryLogsTools;
|
||
|
||
@Autowired(required = false)
|
||
private ToolCallbackProvider tools;
|
||
|
||
@Autowired
|
||
private ChatModel chatModel;
|
||
|
||
@Autowired
|
||
private LookupKnowledgeTool lookupKnowledgeTool;
|
||
|
||
@Autowired
|
||
private DiagnosisSessionRepository diagnosisSessionRepository;
|
||
|
||
@Autowired
|
||
private AgentStepRepository agentStepRepository;
|
||
|
||
@Autowired
|
||
private ToolInvocationRepository toolInvocationRepository;
|
||
|
||
@Autowired
|
||
private EvaluationService evaluationService;
|
||
|
||
@Autowired
|
||
private RetrievedDocTracker retrievedDocTracker;
|
||
|
||
@Autowired
|
||
private KnowledgeDomainService knowledgeDomainService;
|
||
|
||
@Autowired(required = false)
|
||
private SkillRegistry skillRegistry;
|
||
|
||
@Autowired
|
||
private ToolTraceSummaryService toolTraceSummaryService;
|
||
|
||
@Autowired
|
||
private SelfEvaluationMergeService selfEvaluationMergeService;
|
||
|
||
@Value("${verifier.low-confidence-threshold:0.5}")
|
||
private double verifierLowConfidenceThreshold;
|
||
|
||
@Value("${chat.complex.retry-on-low-confidence:false}")
|
||
private boolean retryOnLowConfidence;
|
||
|
||
/** 多 Agent Chat 的 Prompt */
|
||
private String chatPlannerPrompt;
|
||
private String chatExecutorPrompt;
|
||
private String chatVerifierPrompt;
|
||
private final ObjectMapper objectMapper = new ObjectMapper();
|
||
|
||
@PostConstruct
|
||
public void init() {
|
||
// 加载 Prompt
|
||
try {
|
||
chatPlannerPrompt = new String(
|
||
new ClassPathResource("prompts/chat-planner-prompt.md").getInputStream().readAllBytes(),
|
||
StandardCharsets.UTF_8);
|
||
chatExecutorPrompt = new String(
|
||
new ClassPathResource("prompts/chat-executor-prompt.md").getInputStream().readAllBytes(),
|
||
StandardCharsets.UTF_8);
|
||
chatVerifierPrompt = new String(
|
||
new ClassPathResource("prompts/chat-verifier-prompt.md").getInputStream().readAllBytes(),
|
||
StandardCharsets.UTF_8);
|
||
logger.info("Chat 多 Agent Prompts 加载成功");
|
||
} catch (IOException e) {
|
||
logger.error("加载 Chat Prompt 文件失败", e);
|
||
throw new RuntimeException("Failed to load chat prompts", e);
|
||
}
|
||
|
||
// 包装 ChatModel 以捕获 token 用量
|
||
chatModel = new TokenTrackingChatModel(chatModel);
|
||
logger.info("ChatModel 已包装 TokenTrackingChatModel");
|
||
}
|
||
|
||
/**
|
||
* 获取注入的 ChatModel
|
||
*/
|
||
public ChatModel getChatModel() {
|
||
return chatModel;
|
||
}
|
||
|
||
/**
|
||
* 构建系统提示词(包含历史消息)
|
||
* @param history 历史消息列表
|
||
* @return 完整的系统提示词
|
||
*/
|
||
public String buildSystemPrompt(List<Map<String, String>> history) {
|
||
StringBuilder systemPromptBuilder = new StringBuilder();
|
||
|
||
// 基础系统提示
|
||
systemPromptBuilder.append("你是一个专业的智能助手,可以获取当前时间、查询天气信息、搜索内部文档知识库,以及查询 Prometheus 告警信息。\n");
|
||
systemPromptBuilder.append("当用户询问时间相关问题时,**必须每次都调用 getCurrentDateTime 工具**,因为时间会不断变化。即使历史消息中有时间信息,也不要直接复用,必须重新查询最新时间。\n");
|
||
systemPromptBuilder.append("当用户需要查询公司内部文档、流程、最佳实践或技术指南时,使用 lookupKnowledgeTool 工具。\n");
|
||
systemPromptBuilder.append("当用户的问题匹配某个诊断 Skill 时,先调用 read_skill 读取对应流程,再按流程调用证据工具。\n");
|
||
systemPromptBuilder.append("当用户需要查询 Prometheus 告警、监控指标或系统告警状态时,使用 queryPrometheusAlerts 工具。\n");
|
||
systemPromptBuilder.append("当用户需要查询腾讯云日志时,请调用腾讯云mcp服务查询,默认查询地域ap-guangzhou,查询时间范围为近一个月。\n\n");
|
||
|
||
// 添加历史消息(过滤时间查询相关内容)
|
||
if (!history.isEmpty()) {
|
||
systemPromptBuilder.append("--- 对话历史 ---\n");
|
||
for (Map<String, String> msg : history) {
|
||
String role = msg.get("role");
|
||
String content = msg.get("content");
|
||
|
||
// 🔧 过滤时间查询相关的历史消息,避免 LLM 复用旧的时间信息
|
||
if ("user".equals(role) && isTimeQuery(content)) {
|
||
continue; // 跳过时间查询问题
|
||
}
|
||
if ("assistant".equals(role) && containsTimeInfo(content)) {
|
||
continue; // 跳过包含时间信息的回答
|
||
}
|
||
|
||
if ("user".equals(role)) {
|
||
systemPromptBuilder.append("用户: ").append(content).append("\n");
|
||
} else if ("assistant".equals(role)) {
|
||
systemPromptBuilder.append("助手: ").append(content).append("\n");
|
||
}
|
||
}
|
||
systemPromptBuilder.append("--- 对话历史结束 ---\n\n");
|
||
}
|
||
|
||
systemPromptBuilder.append("请基于以上对话历史,回答用户的新问题。");
|
||
|
||
return systemPromptBuilder.toString();
|
||
}
|
||
|
||
/**
|
||
* 判断是否为时间查询问题
|
||
*/
|
||
private boolean isTimeQuery(String content) {
|
||
if (content == null) {
|
||
return false;
|
||
}
|
||
// 匹配常见的时间查询模式
|
||
return content.matches(".*(现在|当前|此时).*(几点|时间).*") ||
|
||
content.matches(".*(几点|时间).*(了|呢|[??]).*") ||
|
||
content.toLowerCase().matches(".*(what.*time|current.*time).*");
|
||
}
|
||
|
||
/**
|
||
* 判断是否包含时间信息
|
||
*/
|
||
private boolean containsTimeInfo(String content) {
|
||
if (content == null) {
|
||
return false;
|
||
}
|
||
// 匹配日期时间格式:2026年5月31日、15:57、下午3点 等
|
||
return content.matches(".*(\\d{4}年\\d{1,2}月\\d{1,2}日|\\d{1,2}:\\d{2}|[上下午]+\\d{1,2}[点时]).*");
|
||
}
|
||
|
||
/**
|
||
* 动态构建方法工具数组
|
||
* 根据已注入的 Bean 暴露本地工具,避免 mock/真实模式下漏注入。
|
||
*/
|
||
public Object[] buildMethodToolsArray() {
|
||
List<Object> methodTools = new ArrayList<>();
|
||
methodTools.add(dateTimeTools);
|
||
methodTools.add(lookupKnowledgeTool);
|
||
if (queryLogsTools != null) {
|
||
methodTools.add(queryLogsTools);
|
||
}
|
||
if (queryMetricsTools != null) {
|
||
methodTools.add(queryMetricsTools);
|
||
}
|
||
return methodTools.toArray();
|
||
}
|
||
|
||
/**
|
||
* 获取工具回调列表,mcp服务提供的工具
|
||
*/
|
||
public ToolCallback[] getToolCallbacks() {
|
||
if (tools == null) {
|
||
return new ToolCallback[0];
|
||
}
|
||
return tools.getToolCallbacks();
|
||
}
|
||
|
||
/**
|
||
* 记录可用工具列表:mcp服务提供的工具
|
||
*/
|
||
public void logAvailableTools() {
|
||
if (tools == null) {
|
||
logger.info("MCP 未启用,无远程工具");
|
||
return;
|
||
}
|
||
ToolCallback[] toolCallbacks = tools.getToolCallbacks();
|
||
logger.info("可用工具列表:");
|
||
for (ToolCallback toolCallback : toolCallbacks) {
|
||
logger.info(">>> {}", toolCallback.getToolDefinition().name());
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 创建 ReactAgent
|
||
* @param chatModel 聊天模型
|
||
* @param systemPrompt 系统提示词
|
||
* @return 配置好的 ReactAgent
|
||
*/
|
||
public ReactAgent createReactAgent(ChatModel chatModel, String systemPrompt) {
|
||
return ReactAgent.builder()
|
||
.name("intelligent_assistant")
|
||
.model(chatModel)
|
||
.systemPrompt(systemPrompt)
|
||
.methodTools(buildMethodToolsArray())
|
||
.tools(getToolCallbacks())
|
||
.hooks(buildHooks("intelligent_assistant"))
|
||
.build();
|
||
}
|
||
|
||
/**
|
||
* 执行 ReactAgent 对话(非流式)
|
||
* @param agent ReactAgent 实例
|
||
* @param question 用户问题
|
||
* @return ChatResult(answer + sessionId)
|
||
*/
|
||
public ChatResult executeChat(ReactAgent agent, String question) throws GraphRunnerException {
|
||
return executeChat(agent, question, null);
|
||
}
|
||
|
||
public ChatResult executeChat(ReactAgent agent, String question, String requestedSessionId) throws GraphRunnerException {
|
||
logger.info("========================================");
|
||
logger.info("📝 用户问题: {}", question);
|
||
|
||
String sessionId = resolveSessionId(requestedSessionId);
|
||
long startTime = System.currentTimeMillis();
|
||
|
||
// 创建或更新诊断会话
|
||
DiagnosisSession session = startDiagnosisSession(sessionId, question);
|
||
diagnosisSessionRepository.save(session);
|
||
|
||
// 设置 ThreadLocal 上下文(LookupKnowledgeTool 通过此获取 sessionId)
|
||
SessionContextHolder.setSessionId(sessionId);
|
||
|
||
try {
|
||
// 通过 RunnableConfig 将 sessionId 传入 Hook(线程安全,异步也兼容)
|
||
var config = RunnableConfig.builder()
|
||
.addMetadata("sessionId", sessionId)
|
||
.build();
|
||
|
||
var response = agent.call(question, config);
|
||
long duration = System.currentTimeMillis() - startTime;
|
||
|
||
String answer = response.getText();
|
||
|
||
// 更新诊断会话
|
||
session.setStatus("SUCCESS");
|
||
session.setAnswer(answer);
|
||
session.setTotalDurationMs((int) duration);
|
||
backfillSessionMetrics(session);
|
||
diagnosisSessionRepository.save(session);
|
||
|
||
evaluationService.evaluate(sessionId, answer);
|
||
|
||
logger.info("⏱️ 总耗时: {} ms", duration);
|
||
logger.info("📏 输出长度: {} 字符", answer.length());
|
||
logger.info("========================================");
|
||
|
||
return new ChatResult(answer, sessionId);
|
||
} catch (Exception e) {
|
||
session.setStatus("FAILED");
|
||
diagnosisSessionRepository.save(session);
|
||
throw e;
|
||
} finally {
|
||
retrievedDocTracker.clearSession(sessionId);
|
||
SessionContextHolder.clear();
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 根据问题复杂度自动选择执行策略
|
||
* @param chatModel 聊天模型
|
||
* @param toolCallbacks 工具回调
|
||
* @param question 用户问题
|
||
* @param history 历史消息
|
||
* @return ChatResult(answer + sessionId)
|
||
*/
|
||
public ChatResult executeChatWithStrategy(ChatModel chatModel, ToolCallback[] toolCallbacks,
|
||
String question, List<Map<String, String>> history) throws GraphRunnerException {
|
||
return executeChatWithStrategy(chatModel, toolCallbacks, question, history, null);
|
||
}
|
||
|
||
public ChatResult executeChatWithStrategy(ChatModel chatModel, ToolCallback[] toolCallbacks,
|
||
String question, List<Map<String, String>> history,
|
||
String requestedSessionId) throws GraphRunnerException {
|
||
if (QuestionComplexity.isComplex(question)) {
|
||
logger.info("📊 问题判定为复杂,使用多 Agent(Planner + Executor)执行");
|
||
return executeChatComplex(chatModel, toolCallbacks, question, history, requestedSessionId);
|
||
} else {
|
||
logger.info("📊 问题判定为简单,使用单 Agent 执行");
|
||
String systemPrompt = buildSystemPrompt(history);
|
||
ReactAgent agent = createReactAgent(chatModel, systemPrompt);
|
||
return executeChat(agent, question, requestedSessionId);
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 多 Agent 复杂对话执行(Planner -> Executor -> Verifier)
|
||
*/
|
||
public ChatResult executeChatComplex(ChatModel chatModel, ToolCallback[] toolCallbacks,
|
||
String question, List<Map<String, String>> history) throws GraphRunnerException {
|
||
return executeChatComplex(chatModel, toolCallbacks, question, history, null);
|
||
}
|
||
|
||
public ChatResult executeChatComplex(ChatModel chatModel, ToolCallback[] toolCallbacks,
|
||
String question, List<Map<String, String>> history,
|
||
String requestedSessionId) throws GraphRunnerException {
|
||
String sessionId = resolveSessionId(requestedSessionId);
|
||
long startTime = System.currentTimeMillis();
|
||
|
||
DiagnosisSession session = startDiagnosisSession(sessionId, question);
|
||
diagnosisSessionRepository.save(session);
|
||
|
||
SessionContextHolder.setSessionId(sessionId);
|
||
VerifierContextHolder.setOriginalQuery(question);
|
||
VerifierContextHolder.setRetryContext(null);
|
||
VerifierContextHolder.setExecutorFinalAnswer(null);
|
||
|
||
try {
|
||
VerifierDecision finalDecision = null;
|
||
String retryContext = null;
|
||
String answer = null;
|
||
RunnableConfig config = RunnableConfig.builder()
|
||
.addMetadata("sessionId", sessionId)
|
||
.build();
|
||
|
||
for (int round = 1; round <= 2; round++) {
|
||
VerifierContextHolder.setRetryContext(retryContext);
|
||
VerifierContextHolder.setToolTraceSummary(null);
|
||
|
||
ReactAgent planner = buildChatPlannerAgent(chatModel, history, retryContext);
|
||
ReactAgent executor = buildChatExecutorAgent(chatModel, toolCallbacks, history, retryContext);
|
||
ReactAgent verifier = buildChatVerifierAgent(chatModel);
|
||
|
||
SequentialAgent workflow = SequentialAgent.builder()
|
||
.name("chat_workflow")
|
||
.description("按固定顺序执行 Planner、Executor、Verifier 的多 Agent 工作流")
|
||
.subAgents(List.of(planner, executor, verifier))
|
||
.build();
|
||
|
||
String workflowInput = buildWorkflowInput(question, retryContext);
|
||
Optional<OverAllState> stateOptional = workflow.invoke(workflowInput, config);
|
||
if (stateOptional.isEmpty()) {
|
||
finalDecision = buildVerifierFallbackDecision(round, "workflow 未返回有效状态");
|
||
answer = buildLowConfidenceOutput(answer, finalDecision);
|
||
persistVerifierEvaluation(session, finalDecision, round);
|
||
break;
|
||
}
|
||
|
||
String plannerPlan = extractStateText(stateOptional, "planner_plan");
|
||
answer = extractStateText(stateOptional, "executor_feedback");
|
||
VerifierContextHolder.setExecutorFinalAnswer(answer);
|
||
String verifierOutput = extractStateText(stateOptional, "verifier_output");
|
||
if ((verifierOutput == null || verifierOutput.isBlank()) && answer != null && !answer.isBlank()) {
|
||
verifierOutput = invokeVerifierFallback(verifier, question, round, config);
|
||
}
|
||
finalDecision = parseVerifierDecision(verifierOutput, round);
|
||
logger.debug("Sequential workflow round {} finished: plannerPlanLength={}, answerLength={}, verifierOutputLength={}",
|
||
round,
|
||
plannerPlan != null ? plannerPlan.length() : 0,
|
||
answer != null ? answer.length() : 0,
|
||
verifierOutput != null ? verifierOutput.length() : 0);
|
||
|
||
if (finalDecision == null) {
|
||
finalDecision = buildVerifierFallbackDecision(round, "verifier_output 缺失或无法解析");
|
||
answer = buildLowConfidenceOutput(answer, finalDecision);
|
||
persistVerifierEvaluation(session, finalDecision, round);
|
||
break;
|
||
}
|
||
|
||
if ("PASS".equals(finalDecision.verdict())) {
|
||
answer = answer == null || answer.isBlank() ? "抱歉,多 Agent 分析未能生成有效结论。" : answer;
|
||
persistVerifierEvaluation(session, finalDecision, round);
|
||
break;
|
||
}
|
||
|
||
if ("REJECT".equals(finalDecision.verdict())) {
|
||
answer = buildDegradedOutput(finalDecision);
|
||
persistVerifierEvaluation(session, finalDecision, round);
|
||
break;
|
||
}
|
||
|
||
boolean shouldRetry = retryOnLowConfidence
|
||
&& finalDecision.groundednessScore() < verifierLowConfidenceThreshold
|
||
&& round < 2;
|
||
if (!shouldRetry) {
|
||
answer = buildLowConfidenceOutput(answer, finalDecision);
|
||
persistVerifierEvaluation(session, finalDecision, round);
|
||
break;
|
||
}
|
||
|
||
retryContext = buildRetryContext(finalDecision);
|
||
persistVerifierEvaluation(session, finalDecision, round);
|
||
}
|
||
|
||
long duration = System.currentTimeMillis() - startTime;
|
||
|
||
if (answer == null || answer.isBlank()) {
|
||
answer = "抱歉,多 Agent 分析未能生成有效结论。";
|
||
}
|
||
|
||
session.setStatus("SUCCESS");
|
||
session.setAnswer(answer);
|
||
session.setTotalDurationMs((int) duration);
|
||
backfillSessionMetrics(session);
|
||
diagnosisSessionRepository.save(session);
|
||
|
||
evaluationService.evaluate(sessionId, answer);
|
||
|
||
logger.info("⏱️ 多 Agent 总耗时: {} ms", duration);
|
||
logger.info("📏 输出长度: {} 字符", answer.length());
|
||
|
||
return new ChatResult(answer, sessionId);
|
||
|
||
} catch (Exception e) {
|
||
session.setStatus("FAILED");
|
||
diagnosisSessionRepository.save(session);
|
||
logger.error("多 Agent 执行失败", e);
|
||
return new ChatResult("执行失败: " + e.getMessage(), sessionId);
|
||
} finally {
|
||
retrievedDocTracker.clearSession(sessionId);
|
||
SessionContextHolder.clear();
|
||
VerifierContextHolder.clear();
|
||
}
|
||
}
|
||
|
||
private ReactAgent buildChatPlannerAgent(ChatModel chatModel, List<Map<String, String>> history,
|
||
String retryContext) {
|
||
StringBuilder prompt = new StringBuilder(chatPlannerPrompt);
|
||
|
||
// 注入 knowledge map
|
||
String knowledgeMap = knowledgeDomainService.buildKnowledgeMap();
|
||
if (!knowledgeMap.isBlank()) {
|
||
prompt.append("\n\n## 可用知识库\n\n").append(knowledgeMap);
|
||
}
|
||
|
||
if (!history.isEmpty()) {
|
||
prompt.append("\n\n--- 对话历史 ---\n");
|
||
for (Map<String, String> msg : history) {
|
||
prompt.append(msg.get("role")).append(": ").append(msg.get("content")).append("\n");
|
||
}
|
||
prompt.append("--- 对话历史结束 ---\n");
|
||
}
|
||
if (retryContext != null && !retryContext.isBlank()) {
|
||
prompt.append("\n\n--- 本轮补证据约束 ---\n").append(retryContext).append("\n");
|
||
}
|
||
return ReactAgent.builder()
|
||
.name("chat_planner")
|
||
.description("负责拆解问题、规划步骤")
|
||
.model(chatModel)
|
||
.systemPrompt(prompt.toString())
|
||
.hooks(buildHooks("planner"))
|
||
.outputKey("planner_plan")
|
||
.build();
|
||
}
|
||
|
||
private ReactAgent buildChatVerifierAgent(ChatModel chatModel) {
|
||
return ReactAgent.builder()
|
||
.name("chat_verifier")
|
||
.description("负责验证 Executor 答案的事实准确性")
|
||
.model(chatModel)
|
||
.systemPrompt(chatVerifierPrompt)
|
||
.hooks(new AgentLoggingHook(agentStepRepository, "verifier"),
|
||
new VerifierInputHook(toolTraceSummaryService))
|
||
.outputKey("verifier_output")
|
||
.build();
|
||
}
|
||
|
||
private ReactAgent buildChatExecutorAgent(ChatModel chatModel, ToolCallback[] toolCallbacks,
|
||
List<Map<String, String>> history, String retryContext) {
|
||
StringBuilder prompt = new StringBuilder(chatExecutorPrompt);
|
||
prompt.append("\n\n--- Skill 读取约束 ---\n")
|
||
.append("如果 planner_plan 已给出 selected_skill,本轮 Executor 只允许对该 skill 调用一次 read_skill。")
|
||
.append("读取后必须复用已加载的 playbook 指令继续执行证据工具,不要为了检查支持文件、确认流程或生成报告再次读取同一个 skill。")
|
||
.append("只有 ChatService 启动新的补证据 retry round 时,才可以重新读取 selected_skill。\n");
|
||
if (!history.isEmpty()) {
|
||
prompt.append("\n\n--- 对话历史 ---\n");
|
||
for (Map<String, String> msg : history) {
|
||
prompt.append(msg.get("role")).append(": ").append(msg.get("content")).append("\n");
|
||
}
|
||
prompt.append("--- 对话历史结束 ---\n");
|
||
}
|
||
if (retryContext != null && !retryContext.isBlank()) {
|
||
prompt.append("\n\n--- 本轮补证据约束 ---\n").append(retryContext).append("\n");
|
||
}
|
||
return ReactAgent.builder()
|
||
.name("chat_executor")
|
||
.description("负责执行具体步骤并及时反馈")
|
||
.model(chatModel)
|
||
.systemPrompt(prompt.toString())
|
||
.methodTools(buildMethodToolsArray())
|
||
.tools(toolCallbacks)
|
||
.hooks(buildHooks("executor"))
|
||
.outputKey("executor_feedback")
|
||
.build();
|
||
}
|
||
|
||
private Hook[] buildHooks(String agentName) {
|
||
AgentLoggingHook loggingHook = new AgentLoggingHook(agentStepRepository, agentName);
|
||
if (skillRegistry == null || skillRegistry.size() == 0) {
|
||
return new Hook[]{loggingHook};
|
||
}
|
||
if ("planner".equals(agentName)) {
|
||
return new Hook[]{
|
||
new PlannerSkillMetadataHook(skillRegistry),
|
||
loggingHook
|
||
};
|
||
}
|
||
return new Hook[]{
|
||
loggingHook,
|
||
SkillsAgentHook.builder()
|
||
.skillRegistry(skillRegistry)
|
||
.build()
|
||
};
|
||
}
|
||
|
||
private String resolveSessionId(String requestedSessionId) {
|
||
if (requestedSessionId != null && !requestedSessionId.isBlank()) {
|
||
return requestedSessionId;
|
||
}
|
||
return UUID.randomUUID().toString().substring(0, 8);
|
||
}
|
||
|
||
private DiagnosisSession startDiagnosisSession(String sessionId, String question) {
|
||
DiagnosisSession session = diagnosisSessionRepository.findBySessionId(sessionId)
|
||
.orElseGet(() -> DiagnosisSession.builder()
|
||
.sessionId(sessionId)
|
||
.agentFlow("CHAT")
|
||
.build());
|
||
session.setQuery(question);
|
||
session.setStatus("RUNNING");
|
||
session.setAgentFlow("CHAT");
|
||
session.setAnswer(null);
|
||
session.setTotalDurationMs(null);
|
||
session.setTotalTokenCount(null);
|
||
session.setStepCount(null);
|
||
session.setToolCallCount(null);
|
||
return session;
|
||
}
|
||
|
||
private String buildWorkflowInput(String question, String retryContext) {
|
||
StringBuilder input = new StringBuilder();
|
||
input.append("请按固定工作流完成本轮 Planner -> Executor -> Verifier。\n\n");
|
||
input.append("--- 用户问题 ---\n").append(question);
|
||
if (retryContext != null && !retryContext.isBlank()) {
|
||
input.append("\n\n--- retry_context ---\n").append(retryContext);
|
||
}
|
||
input.append("\n\nVerifier 完成后由外层代码读取 verifier_output 并决定最终用户输出。");
|
||
return input.toString();
|
||
}
|
||
|
||
private String invokeVerifierFallback(ReactAgent verifier, String question, int round, RunnableConfig config) {
|
||
try {
|
||
logger.warn("Sequential workflow round {} finished without verifier_output, invoking chat_verifier fallback", round);
|
||
return verifier.call("请基于 executor_final_answer 和 tool_trace_summary 输出 verifier JSON。原始问题:" + question, config)
|
||
.getText();
|
||
} catch (Exception e) {
|
||
logger.error("chat_verifier fallback 执行失败", e);
|
||
return null;
|
||
}
|
||
}
|
||
|
||
private VerifierDecision parseVerifierDecision(String verifierOutput, int round) {
|
||
if (verifierOutput == null || verifierOutput.isBlank()) {
|
||
return null;
|
||
}
|
||
|
||
try {
|
||
JsonNode root = objectMapper.readTree(sanitizeJsonPayload(verifierOutput));
|
||
List<Map<String, Object>> factsChecked = parseFactsChecked(root.path("facts_checked"));
|
||
|
||
return new VerifierDecision(
|
||
root.path("verdict").asText("LOW_CONFID"),
|
||
root.path("groundedness_score").asDouble(0.0),
|
||
root.path("critical_fact_count").asInt(0),
|
||
factsChecked,
|
||
root.path("rationale").asText(""),
|
||
round
|
||
);
|
||
} catch (Exception e) {
|
||
logger.error("解析 verifier_output 失败: {}", verifierOutput, e);
|
||
return null;
|
||
}
|
||
}
|
||
|
||
private String sanitizeJsonPayload(String raw) {
|
||
String trimmed = raw.trim();
|
||
if (trimmed.startsWith("```")) {
|
||
int firstNewline = trimmed.indexOf('\n');
|
||
int lastFence = trimmed.lastIndexOf("```");
|
||
if (firstNewline >= 0 && lastFence > firstNewline) {
|
||
return trimmed.substring(firstNewline + 1, lastFence).trim();
|
||
}
|
||
}
|
||
return trimmed;
|
||
}
|
||
|
||
private List<Map<String, Object>> parseFactsChecked(JsonNode factsNode) {
|
||
List<Map<String, Object>> factsChecked = new ArrayList<>();
|
||
if (!factsNode.isArray()) {
|
||
return factsChecked;
|
||
}
|
||
for (JsonNode factNode : factsNode) {
|
||
Map<String, Object> fact = new LinkedHashMap<>();
|
||
fact.put("fact", factNode.path("fact").asText(""));
|
||
fact.put("is_critical", factNode.path("is_critical").asBoolean(false));
|
||
fact.put("verification", factNode.path("verification").asText(""));
|
||
fact.put("detail", factNode.path("detail").asText(""));
|
||
fact.put("evidence_refs", parseEvidenceRefs(factNode.path("evidence_refs")));
|
||
factsChecked.add(fact);
|
||
}
|
||
return factsChecked;
|
||
}
|
||
|
||
private List<Map<String, Object>> parseEvidenceRefs(JsonNode evidenceRefsNode) {
|
||
List<Map<String, Object>> evidenceRefs = new ArrayList<>();
|
||
if (!evidenceRefsNode.isArray()) {
|
||
return evidenceRefs;
|
||
}
|
||
for (JsonNode refNode : evidenceRefsNode) {
|
||
Map<String, Object> evidenceRef = new LinkedHashMap<>();
|
||
evidenceRef.put("trace_ref", refNode.path("trace_ref").asText(""));
|
||
evidenceRef.put("tool_name", refNode.path("tool_name").asText(""));
|
||
evidenceRef.put("topic_domain", refNode.path("topic_domain").asText(""));
|
||
evidenceRef.put("note", refNode.path("note").asText(""));
|
||
|
||
List<Long> sourceInvocationIds = new ArrayList<>();
|
||
JsonNode idsNode = refNode.path("source_invocation_ids");
|
||
if (idsNode.isArray()) {
|
||
for (JsonNode idNode : idsNode) {
|
||
if (idNode.canConvertToLong()) {
|
||
sourceInvocationIds.add(idNode.asLong());
|
||
}
|
||
}
|
||
}
|
||
evidenceRef.put("source_invocation_ids", sourceInvocationIds);
|
||
evidenceRefs.add(evidenceRef);
|
||
}
|
||
return evidenceRefs;
|
||
}
|
||
|
||
private VerifierDecision buildVerifierFallbackDecision(int round, String rationale) {
|
||
return new VerifierDecision("LOW_CONFID", 0.0, 0, List.of(), rationale, round);
|
||
}
|
||
|
||
private String extractStateText(Optional<OverAllState> stateOptional, String key) {
|
||
if (stateOptional.isEmpty()) {
|
||
return null;
|
||
}
|
||
return stateOptional.get().value(key)
|
||
.map(value -> {
|
||
if (value instanceof AssistantMessage assistantMessage) {
|
||
return assistantMessage.getText();
|
||
}
|
||
return String.valueOf(value);
|
||
})
|
||
.orElse(null);
|
||
}
|
||
|
||
private void persistVerifierEvaluation(DiagnosisSession session, VerifierDecision decision, int round) {
|
||
if (decision == null) {
|
||
return;
|
||
}
|
||
Map<String, Object> verifierEvaluation = new LinkedHashMap<>();
|
||
verifierEvaluation.put("verdict", decision.verdict());
|
||
verifierEvaluation.put("groundedness_score", decision.groundednessScore());
|
||
verifierEvaluation.put("critical_fact_count", decision.criticalFactCount());
|
||
verifierEvaluation.put("facts_checked", decision.factsChecked());
|
||
verifierEvaluation.put("rationale", decision.rationale());
|
||
verifierEvaluation.put("round", round);
|
||
verifierEvaluation.put("traceability_version", "v1");
|
||
verifierEvaluation.put("tool_trace_summary",
|
||
Optional.ofNullable(VerifierContextHolder.getToolTraceSummary()).orElse(List.of()));
|
||
|
||
String merged = selfEvaluationMergeService.mergeVerifierEvaluation(session.getSelfEvaluation(), verifierEvaluation);
|
||
session.setSelfEvaluation(merged);
|
||
diagnosisSessionRepository.save(session);
|
||
}
|
||
|
||
private String buildRetryContext(VerifierDecision decision) {
|
||
try {
|
||
List<String> missingFacts = extractEvidenceGaps(decision);
|
||
Map<String, Object> retryContext = new LinkedHashMap<>();
|
||
retryContext.put("round", decision.round());
|
||
retryContext.put("missing_evidence_facts", missingFacts);
|
||
retryContext.put("instruction", "仅补充以上断言相关证据,不要重复已完成检索");
|
||
return objectMapper.writeValueAsString(retryContext);
|
||
} catch (Exception e) {
|
||
logger.error("构造 retry_context 失败", e);
|
||
return "{\"round\":1,\"missing_evidence_facts\":[],\"instruction\":\"仅补充缺失证据\"}";
|
||
}
|
||
}
|
||
|
||
private String buildLowConfidenceOutput(String executorAnswer, VerifierDecision decision) {
|
||
StringBuilder output = new StringBuilder(LOW_CONFID_DISCLAIMER);
|
||
List<String> confirmedFacts = extractConfirmedFacts(decision);
|
||
output.append("\n\n已确认信息:");
|
||
if (confirmedFacts.isEmpty()) {
|
||
output.append("\n- 暂无可稳定确认的信息");
|
||
} else {
|
||
for (String fact : confirmedFacts) {
|
||
output.append("\n- ").append(fact);
|
||
}
|
||
}
|
||
|
||
List<String> gaps = extractEvidenceGaps(decision);
|
||
output.append("\n\n当前缺口:");
|
||
if (!gaps.isEmpty()) {
|
||
for (String gap : gaps) {
|
||
output.append("\n- ").append(gap);
|
||
}
|
||
} else {
|
||
output.append("\n- 当前缺少足够的直接证据支撑核心结论");
|
||
}
|
||
|
||
output.append("\n\n建议下一步:");
|
||
for (String suggestion : buildNextStepSuggestions(decision)) {
|
||
output.append("\n- ").append(suggestion);
|
||
}
|
||
return output.toString();
|
||
}
|
||
|
||
private String buildDegradedOutput(VerifierDecision decision) {
|
||
StringBuilder output = new StringBuilder(DEGRADED_PREFIX);
|
||
|
||
List<String> confirmedFacts = extractConfirmedFacts(decision);
|
||
List<String> gaps = extractEvidenceGaps(decision);
|
||
List<String> suggestions = buildNextStepSuggestions(decision);
|
||
|
||
output.append("\n\n已确认信息:");
|
||
if (confirmedFacts.isEmpty()) {
|
||
output.append("\n- 暂无可稳定确认的信息");
|
||
} else {
|
||
for (String fact : confirmedFacts) {
|
||
output.append("\n- ").append(fact);
|
||
}
|
||
}
|
||
|
||
output.append("\n\n证据缺口:");
|
||
if (gaps.isEmpty()) {
|
||
output.append("\n- 当前缺少足够的直接证据支撑核心结论");
|
||
} else {
|
||
for (String gap : gaps) {
|
||
output.append("\n- ").append(gap);
|
||
}
|
||
}
|
||
|
||
output.append("\n\n建议下一步:");
|
||
for (String suggestion : suggestions) {
|
||
output.append("\n- ").append(suggestion);
|
||
}
|
||
return output.toString();
|
||
}
|
||
|
||
private List<String> extractConfirmedFacts(VerifierDecision decision) {
|
||
List<String> confirmedFacts = new ArrayList<>();
|
||
for (Map<String, Object> fact : decision.factsChecked()) {
|
||
String verification = String.valueOf(fact.get("verification"));
|
||
boolean critical = Boolean.TRUE.equals(fact.get("is_critical"));
|
||
if (critical && "direct_evidence".equals(verification)) {
|
||
confirmedFacts.add(String.valueOf(fact.get("fact")));
|
||
}
|
||
}
|
||
return confirmedFacts;
|
||
}
|
||
|
||
private List<String> extractEvidenceGaps(VerifierDecision decision) {
|
||
List<String> gaps = new ArrayList<>();
|
||
for (Map<String, Object> fact : decision.factsChecked()) {
|
||
String verification = String.valueOf(fact.get("verification"));
|
||
boolean critical = Boolean.TRUE.equals(fact.get("is_critical"));
|
||
if (critical && ("no_evidence".equals(verification) || "contradicted".equals(verification))) {
|
||
gaps.add(String.valueOf(fact.get("fact")) + ":" + String.valueOf(fact.get("detail")));
|
||
}
|
||
}
|
||
if (gaps.isEmpty() && "LOW_CONFID".equals(decision.verdict())) {
|
||
for (Map<String, Object> fact : decision.factsChecked()) {
|
||
String verification = String.valueOf(fact.get("verification"));
|
||
boolean critical = Boolean.TRUE.equals(fact.get("is_critical"));
|
||
if (critical && "indirect_support".equals(verification)) {
|
||
gaps.add(String.valueOf(fact.get("fact")) + ":缺少直接证据锚点");
|
||
}
|
||
}
|
||
}
|
||
return gaps;
|
||
}
|
||
|
||
private List<String> buildNextStepSuggestions(VerifierDecision decision) {
|
||
List<String> suggestions = new ArrayList<>();
|
||
List<Map<String, Object>> toolSummary = toolTraceSummaryService.buildVerifierTraceSummary(SessionContextHolder.getSessionId(), null);
|
||
boolean hasKnowledgeTool = toolSummary.stream().anyMatch(item -> "lookup_knowledge".equals(item.get("tool_name")));
|
||
boolean hasFailedEvidence = toolSummary.stream().anyMatch(item -> !Boolean.TRUE.equals(item.get("success")));
|
||
|
||
if (!hasKnowledgeTool) {
|
||
suggestions.add("补充知识库或业务文档检索结果,建立可引用的证据锚点");
|
||
}
|
||
if (hasFailedEvidence) {
|
||
suggestions.add("优先重试失败的证据型查询,补齐日志、指标或知识库侧证据");
|
||
}
|
||
if (suggestions.isEmpty()) {
|
||
suggestions.add("围绕上述证据缺口补充只读查询,再由人工复核最终结论");
|
||
}
|
||
return suggestions;
|
||
}
|
||
|
||
private record VerifierDecision(
|
||
String verdict,
|
||
double groundednessScore,
|
||
int criticalFactCount,
|
||
List<Map<String, Object>> factsChecked,
|
||
String rationale,
|
||
int round
|
||
) {
|
||
}
|
||
|
||
/** 从 agent_step 和 tool_invocation 汇总指标回填 diagnosis_session */
|
||
private void backfillSessionMetrics(DiagnosisSession session) {
|
||
try {
|
||
List<com.superbiz.agent.domain.entity.AgentStep> steps =
|
||
agentStepRepository.findBySessionIdOrderByStepIndex(session.getSessionId());
|
||
|
||
int totalTokens = 0;
|
||
int stepCount = 0;
|
||
for (var s : steps) {
|
||
stepCount++;
|
||
if (s.getTokenCount() != null) totalTokens += s.getTokenCount();
|
||
}
|
||
long toolCallCount = toolInvocationRepository.countBySessionId(session.getSessionId());
|
||
session.setTotalTokenCount(totalTokens);
|
||
session.setStepCount(stepCount);
|
||
session.setToolCallCount(Math.toIntExact(toolCallCount));
|
||
} catch (Exception e) {
|
||
logger.warn("回填会话指标失败: sessionId={}", session.getSessionId(), e);
|
||
}
|
||
}
|
||
}
|