Merge branch 'emdash/afraid-geese-carry-h5718' into refactor/mvp1.0

# Conflicts:
#	devflow/index.md
This commit is contained in:
zhuyongxin
2026-06-26 17:36:59 +08:00
65 changed files with 4901 additions and 1034 deletions
@@ -15,7 +15,12 @@ import java.util.List;
/**
* 内部文档查询工具
* 使用 RAG (Retrieval-Augmented Generation) 从内部知识库检索相关文档
*
* @deprecated 请使用 {@link com.superbiz.agent.tool.LookupKnowledgeTool} 替代。
* lookup_knowledge 支持 L0 精确匹配 + L1 语义检索,性能更优且功能更全面。
* 计划在下一个版本中移除此工具。
*/
@Deprecated
@Component
public class InternalDocsTools {
@@ -45,7 +50,9 @@ public class InternalDocsTools {
*
* @param query 搜索查询,描述您要查找的信息
* @return JSON 格式的搜索结果,包含相关文档内容、相似度分数和元数据
* @deprecated 请使用 {@link com.superbiz.agent.tool.LookupKnowledgeTool#lookupKnowledge(String)} 替代
*/
@Deprecated
@Tool(description = "Use this tool to search internal documentation and knowledge base for relevant information. " +
"It performs RAG (Retrieval-Augmented Generation) to find similar documents and extract processing steps. " +
"This is useful when you need to understand internal procedures, best practices, or step-by-step guides " +
@@ -0,0 +1,56 @@
package com.superbiz.agent.config;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.ClassPathResource;
import jakarta.annotation.PostConstruct;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
/**
* AI Ops Agent Prompt 配置
* 从独立的 Markdown 文件加载 Prompt 模板
*/
@Slf4j
@Configuration
public class AiOpsPromptProperties {
private String planner;
private String executor;
private String supervisor;
@PostConstruct
public void loadPrompts() {
try {
planner = loadPromptFromFile("prompts/planner-prompt.md");
executor = loadPromptFromFile("prompts/executor-prompt.md");
supervisor = loadPromptFromFile("prompts/supervisor-prompt.md");
log.info("AI Ops Prompts 加载成功");
log.debug("Planner Prompt 长度: {} 字符", planner.length());
log.debug("Executor Prompt 长度: {} 字符", executor.length());
log.debug("Supervisor Prompt 长度: {} 字符", supervisor.length());
} catch (IOException e) {
log.error("加载 Prompt 文件失败", e);
throw new RuntimeException("Failed to load AI Ops prompts", e);
}
}
private String loadPromptFromFile(String path) throws IOException {
ClassPathResource resource = new ClassPathResource(path);
return new String(resource.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
}
public String getPlanner() {
return planner;
}
public String getExecutor() {
return executor;
}
public String getSupervisor() {
return supervisor;
}
}
@@ -83,16 +83,12 @@ public class ChatController {
// 记录可用工具
chatService.logAvailableTools();
ToolCallback[] toolCallbacks = tools != null ? tools.getToolCallbacks() : new ToolCallback[0];
// 根据问题复杂度自动选择单 Agent 或多 Agent
logger.info("开始 ReactAgent 对话(支持自动工具调用)");
// 构建系统提示词(包含历史消息)
String systemPrompt = chatService.buildSystemPrompt(history);
// 创建 ReactAgent
ReactAgent agent = chatService.createReactAgent(chatModel, systemPrompt);
// 执行对话
String fullAnswer = chatService.executeChat(agent, request.getQuestion());
String fullAnswer = chatService.executeChatWithStrategy(chatModel, toolCallbacks,
request.getQuestion(), history);
// 更新会话历史
session.addMessage(request.getQuestion(), fullAnswer);
@@ -0,0 +1,94 @@
package com.superbiz.agent.controller;
import com.superbiz.agent.service.KnowledgeBaseInitService;
import lombok.Data;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
import java.util.HashMap;
import java.util.Map;
/**
* 知识库管理控制器
* 提供知识库初始化、查询等接口
*/
@RestController
@RequestMapping("/api/knowledge")
public class KnowledgeBaseController {
private static final Logger logger = LoggerFactory.getLogger(KnowledgeBaseController.class);
@Autowired
private KnowledgeBaseInitService initService;
/**
* 初始化知识库
* 扫描 knowledge_base 目录下的所有文档,去重后批量导入到数据库和 Milvus
*
* @param force 是否强制重新导入(跳过去重检查)
* @return 初始化结果
*/
@PostMapping("/init")
public ResponseEntity<?> initKnowledgeBase(@RequestParam(defaultValue = "false") boolean force) {
logger.info("收到知识库初始化请求, force={}", force);
try {
KnowledgeBaseInitService.InitResult result = initService.initializeKnowledgeBase(force);
Map<String, Object> response = new HashMap<>();
response.put("success", true);
response.put("message", "知识库初始化完成");
response.put("scanned", result.getScanned());
response.put("skipped", result.getSkipped());
response.put("inserted", result.getInserted());
response.put("failed", result.getFailed());
response.put("details", result.getDetails());
logger.info("知识库初始化成功: 扫描={}, 跳过={}, 新增={}, 失败={}",
result.getScanned(), result.getSkipped(), result.getInserted(), result.getFailed());
return ResponseEntity.ok(response);
} catch (Exception e) {
logger.error("知识库初始化失败", e);
Map<String, Object> response = new HashMap<>();
response.put("success", false);
response.put("message", "初始化失败: " + e.getMessage());
return ResponseEntity.internalServerError().body(response);
}
}
/**
* 查询知识库统计信息
*
* @return 统计信息
*/
@GetMapping("/stats")
public ResponseEntity<?> getStats() {
try {
KnowledgeBaseInitService.Stats stats = initService.getStats();
Map<String, Object> response = new HashMap<>();
response.put("success", true);
response.put("totalDocuments", stats.getTotalDocuments());
response.put("totalVectors", stats.getTotalVectors());
response.put("categories", stats.getCategoryCount());
return ResponseEntity.ok(response);
} catch (Exception e) {
logger.error("查询统计信息失败", e);
Map<String, Object> response = new HashMap<>();
response.put("success", false);
response.put("message", "查询失败: " + e.getMessage());
return ResponseEntity.internalServerError().body(response);
}
}
}
@@ -0,0 +1,64 @@
package com.superbiz.agent.domain.entity;
import jakarta.persistence.*;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
/**
* Agent 决策步骤实体
* 对应表: agent_step
*/
@Entity
@Table(name = "agent_step", indexes = {
@Index(name = "idx_session_step", columnList = "session_id, step_index"),
@Index(name = "idx_agent_name", columnList = "agent_name")
})
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class AgentStep {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(name = "session_id", nullable = false, length = 64)
private String sessionId;
@Column(name = "step_index", nullable = false)
private Integer stepIndex;
@Column(name = "agent_name", nullable = false, length = 32)
private String agentName;
@Column(name = "model_input", columnDefinition = "TEXT")
private String modelInput;
@Column(name = "model_output", columnDefinition = "TEXT")
private String modelOutput;
@Column(name = "thought", columnDefinition = "TEXT")
private String thought;
@Column(name = "has_tool_call")
private Boolean hasToolCall;
@Column(name = "duration_ms")
private Integer durationMs;
@Column(name = "token_count")
private Integer tokenCount;
@Column(name = "created_at", nullable = false, updatable = false)
private LocalDateTime createdAt;
@PrePersist
protected void onCreate() {
createdAt = LocalDateTime.now();
}
}
@@ -38,7 +38,7 @@ public class ApiDocument {
// 文档分类
@Enumerated(EnumType.STRING)
@Column(name = "fault_category", length = 32, columnDefinition = "VARCHAR(32)")
private FaultCategory faultCategory = FaultCategory.EXTERNAL_API;
private FaultCategory faultCategory = FaultCategory.GENERAL;
@Column(name = "fault_source", length = 128)
private String faultSource;
@@ -1,128 +0,0 @@
package com.superbiz.agent.domain.entity;
import jakarta.persistence.*;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import com.superbiz.agent.domain.enums.DiagnosisStatus;
import com.superbiz.agent.domain.enums.FaultCategory;
import org.hibernate.annotations.JdbcTypeCode;
import org.hibernate.type.SqlTypes;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
/**
* 诊断记录实体
* 对应表: diagnosis_record
*/
@Entity
@Table(name = "diagnosis_record", indexes = {
@Index(name = "idx_business_id", columnList = "business_id"),
@Index(name = "idx_trace_id", columnList = "trace_id"),
@Index(name = "idx_session_id", columnList = "session_id"),
@Index(name = "idx_fault_category", columnList = "fault_category"),
@Index(name = "idx_error_code", columnList = "error_code"),
@Index(name = "idx_created_at", columnList = "created_at"),
@Index(name = "idx_status", columnList = "status")
})
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class DiagnosisRecord {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(name = "diagnosis_id", unique = true, nullable = false, length = 64)
private String diagnosisId;
// 关联信息
@Column(name = "session_id", length = 64)
private String sessionId;
@Column(name = "business_id", length = 128)
private String businessId;
@Column(name = "trace_id", length = 64)
private String traceId;
// 故障分类
@Enumerated(EnumType.STRING)
@Column(name = "fault_category", length = 32, columnDefinition = "VARCHAR(32)")
private FaultCategory faultCategory;
@Column(name = "fault_source", length = 128)
private String faultSource;
@Column(name = "fault_target", length = 256)
private String faultTarget;
// 错误信息
@Column(name = "error_code", length = 64)
private String errorCode;
@Column(name = "error_message", columnDefinition = "TEXT")
private String errorMessage;
@Column(name = "stack_trace", columnDefinition = "TEXT")
private String stackTrace;
// 诊断结果
@Column(name = "problem_type", length = 32)
private String problemType;
@Column(name = "root_cause", columnDefinition = "TEXT")
private String rootCause;
@Column(name = "solution", columnDefinition = "TEXT")
private String solution;
@Column(name = "report_markdown", columnDefinition = "TEXT")
private String reportMarkdown;
// 评估指标
@Enumerated(EnumType.STRING)
@Column(name = "status", length = 16, columnDefinition = "VARCHAR(16)")
private DiagnosisStatus status = DiagnosisStatus.PENDING;
@Column(name = "confidence")
private Integer confidence;
@Column(name = "duration")
private Integer duration;
// 用户反馈
@Column(name = "feedback", length = 16)
private String feedback;
// 调试字段 - JSON 类型
@JdbcTypeCode(SqlTypes.JSON)
@Column(name = "tool_calls", columnDefinition = "JSON")
private List<Map<String, Object>> toolCalls;
// 元数据
@Column(name = "created_by", length = 64)
private String createdBy;
@Column(name = "created_at", nullable = false, updatable = false)
private LocalDateTime createdAt;
@Column(name = "updated_at")
private LocalDateTime updatedAt;
@PrePersist
protected void onCreate() {
createdAt = LocalDateTime.now();
updatedAt = LocalDateTime.now();
}
@PreUpdate
protected void onUpdate() {
updatedAt = LocalDateTime.now();
}
}
@@ -0,0 +1,80 @@
package com.superbiz.agent.domain.entity;
import jakarta.persistence.*;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.hibernate.annotations.JdbcTypeCode;
import org.hibernate.type.SqlTypes;
import java.time.LocalDateTime;
/**
* 诊断会话实体
* 对应表: diagnosis_session
*/
@Entity
@Table(name = "diagnosis_session", indexes = {
@Index(name = "idx_created_at", columnList = "created_at"),
@Index(name = "idx_status", columnList = "status"),
@Index(name = "idx_agent_flow", columnList = "agent_flow")
})
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class DiagnosisSession {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(name = "session_id", unique = true, nullable = false, length = 64)
private String sessionId;
@Column(name = "query", nullable = false, columnDefinition = "TEXT")
private String query;
@Column(name = "status", length = 16)
private String status = "PENDING";
@Column(name = "agent_flow", length = 32)
private String agentFlow;
@Column(name = "total_duration_ms")
private Integer totalDurationMs;
@Column(name = "total_token_count")
private Integer totalTokenCount;
@Column(name = "step_count")
private Integer stepCount;
@Column(name = "tool_call_count")
private Integer toolCallCount;
@JdbcTypeCode(SqlTypes.JSON)
@Column(name = "self_evaluation", columnDefinition = "JSON")
private String selfEvaluation;
@Column(name = "feedback", length = 16)
private String feedback;
@Column(name = "created_at", nullable = false, updatable = false)
private LocalDateTime createdAt;
@Column(name = "updated_at")
private LocalDateTime updatedAt;
@PrePersist
protected void onCreate() {
createdAt = LocalDateTime.now();
updatedAt = LocalDateTime.now();
}
@PreUpdate
protected void onUpdate() {
updatedAt = LocalDateTime.now();
}
}
@@ -0,0 +1,84 @@
package com.superbiz.agent.domain.entity;
import jakarta.persistence.*;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.hibernate.annotations.JdbcTypeCode;
import org.hibernate.type.SqlTypes;
import java.time.LocalDateTime;
/**
* 工具调用明细实体
* 对应表: tool_invocation
*/
@Entity
@Table(name = "tool_invocation", indexes = {
@Index(name = "idx_session_id", columnList = "session_id"),
@Index(name = "idx_tool_name", columnList = "tool_name"),
@Index(name = "idx_retrieval_layer", columnList = "retrieval_layer")
})
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class ToolInvocation {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(name = "session_id", nullable = false, length = 64)
private String sessionId;
@Column(name = "step_id")
private Long stepId;
@Column(name = "tool_name", nullable = false, length = 64)
private String toolName;
@JdbcTypeCode(SqlTypes.JSON)
@Column(name = "input_params", nullable = false, columnDefinition = "JSON")
private String inputParams;
@Column(name = "output_preview", columnDefinition = "TEXT")
private String outputPreview;
@Column(name = "output_length")
private Integer outputLength;
@Column(name = "retrieval_layer", length = 8)
private String retrievalLayer;
@Column(name = "l0_match_count")
private Integer l0MatchCount;
@Column(name = "l1_match_count")
private Integer l1MatchCount;
@Column(name = "is_truncated")
private Boolean isTruncated;
@JdbcTypeCode(SqlTypes.JSON)
@Column(name = "retrieval_details", columnDefinition = "JSON")
private String retrievalDetails;
@Column(name = "duration_ms")
private Integer durationMs;
@Column(name = "success")
private Boolean success;
@Column(name = "error_message", columnDefinition = "TEXT")
private String errorMessage;
@Column(name = "created_at", nullable = false, updatable = false)
private LocalDateTime createdAt;
@PrePersist
protected void onCreate() {
createdAt = LocalDateTime.now();
}
}
@@ -1,21 +0,0 @@
package com.superbiz.agent.domain.enums;
/**
* 诊断状态枚举
*/
public enum DiagnosisStatus {
PENDING("待处理"),
RUNNING("诊断中"),
SUCCESS("成功"),
FAILED("失败");
private final String description;
DiagnosisStatus(String description) {
this.description = description;
}
public String getDescription() {
return description;
}
}
@@ -1,17 +1,14 @@
package com.superbiz.agent.domain.enums;
/**
* 故障类别枚举
* 文档分类枚举
*/
public enum FaultCategory {
EXTERNAL_API("外部接口调用失败"),
INTERNAL_ERROR("系统内部错误"),
DATABASE("数据库问题"),
CACHE("缓存问题"),
NETWORK("网络问题"),
THREAD("线程问题"),
MEMORY("内存问题"),
CONFIG("配置问题");
API("API 接口文档"),
INFRASTRUCTURE("基础设施文档"),
DOMAIN("领域业务文档"),
TROUBLESHOOTING("故障排查文档"),
GENERAL("通用文档");
private final String description;
@@ -22,4 +19,26 @@ public enum FaultCategory {
public String getDescription() {
return description;
}
/**
* 从字符串映射到枚举
*/
public static FaultCategory fromString(String category) {
if (category == null || category.isEmpty()) {
return GENERAL;
}
switch (category.toLowerCase()) {
case "api":
return API;
case "infrastructure":
return INFRASTRUCTURE;
case "domain":
return DOMAIN;
case "troubleshooting":
return TROUBLESHOOTING;
default:
return GENERAL;
}
}
}
@@ -38,4 +38,10 @@ public class DocumentChunk {
* 分片标题或上下文信息
*/
private String title;
/**
* 面包屑导航(完整标题层级路径)
* 例如: "故障诊断流程规范 > 应急响应流程 > 1. 初步评估"
*/
private String breadcrumb;
}
@@ -0,0 +1,356 @@
package com.superbiz.agent.hook;
import com.alibaba.cloud.ai.graph.agent.hook.messages.MessagesModelHook;
import com.alibaba.cloud.ai.graph.agent.hook.messages.AgentCommand;
import com.alibaba.cloud.ai.graph.agent.hook.HookPosition;
import com.alibaba.cloud.ai.graph.agent.hook.HookPositions;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import com.superbiz.agent.domain.entity.AgentStep;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.util.SessionContextHolder;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.messages.Message;
import org.springframework.ai.chat.messages.AssistantMessage;
import org.springframework.ai.chat.messages.UserMessage;
import org.springframework.ai.chat.messages.ToolResponseMessage;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* Agent 日志 Hook
* 记录 Agent 的思考过程、消息流转 + 持久化 agent_step 到 DB
*/
@Slf4j
@HookPositions({HookPosition.BEFORE_MODEL, HookPosition.AFTER_MODEL})
public class AgentLoggingHook extends MessagesModelHook {
private final AgentStepRepository agentStepRepository;
private final String agentName;
/** 每个 session 的步数计数器:sessionId → stepIndex */
private final ConcurrentHashMap<String, Integer> stepCounters = new ConcurrentHashMap<>();
/** beforeModel → afterModel 中间状态:sessionId_stepIndex → {stepId, startTime} */
private final ConcurrentHashMap<String, Map<String, Object>> pendingSteps = new ConcurrentHashMap<>();
public AgentLoggingHook(AgentStepRepository agentStepRepository, String agentName) {
this.agentStepRepository = agentStepRepository;
this.agentName = agentName;
}
@Override
public String getName() {
return "agent_logging_hook";
}
@Override
public AgentCommand beforeModel(List<Message> previousMessages, RunnableConfig config) {
// 优先从 config.metadata 取 sessionId(线程安全),兜底 ThreadLocal
String sessionId = config.metadata("sessionId")
.map(Object::toString)
.orElseGet(SessionContextHolder::getSessionId);
boolean hasSession = (sessionId != null);
int stepIndex = 0;
if (hasSession) {
stepIndex = stepCounters.merge(sessionId, 0, (old, one) -> old + 1);
}
log.info("========================================");
log.info("*** [Agent 思考] 第 {} 轮思考开始", (hasSession ? stepCounters.get(sessionId) : 0) + 1);
log.info("*** [Agent 思考] 当前消息数量: {}", previousMessages.size());
// 打印最后几条消息
int lastN = Math.min(3, previousMessages.size());
if (lastN > 0) {
log.info("*** [Agent 思考] 最近 {} 条消息:", lastN);
List<Message> recentMessages = previousMessages.subList(previousMessages.size() - lastN, previousMessages.size());
for (int i = 0; i < recentMessages.size(); i++) {
Message msg = recentMessages.get(i);
String role = getMessageRole(msg);
log.info(" [{}] 角色: {}, 类型: {}", i + 1, role, msg.getClass().getSimpleName());
}
}
log.info("*** [Agent 思考] 准备调用模型...");
log.info("========================================");
// 持久化 agent_step(beforeModel:先创建,先记 model_input 摘要)
if (sessionId != null) {
try {
String modelInputSummary = buildModelInputSummary(previousMessages);
AgentStep step = AgentStep.builder()
.sessionId(sessionId)
.stepIndex(stepIndex)
.agentName(agentName)
.modelInput(modelInputSummary)
.build();
AgentStep saved = agentStepRepository.save(step);
// 记录中间状态供 afterModel 使用
pendingSteps.put(sessionId + "_" + stepIndex, Map.of(
"stepId", saved.getId(),
"startTime", System.currentTimeMillis()
));
log.debug("agent_step 已创建: sessionId={}, stepIndex={}, id={}", sessionId, stepIndex, saved.getId());
} catch (Exception e) {
log.error("保存 agent_step 失败", e);
// 不中断 Agent 执行
}
}
return new AgentCommand(previousMessages);
}
@Override
public AgentCommand afterModel(List<Message> previousMessages, RunnableConfig config) {
String sessionId = SessionContextHolder.getSessionId();
boolean hasSession = (sessionId != null);
log.info("========================================");
log.info("*** [Agent 思考] 第 {} 轮思考完成", (hasSession ? stepCounters.getOrDefault(sessionId, 0) : 0));
// 查找最后一条 AssistantMessage(模型的回复)
AssistantMessage lastAssistant = null;
for (int i = previousMessages.size() - 1; i >= 0; i--) {
if (previousMessages.get(i) instanceof AssistantMessage) {
lastAssistant = (AssistantMessage) previousMessages.get(i);
break;
}
}
boolean hasToolCall = false;
if (lastAssistant != null) {
// 打印模型返回的文本内容
String textContent = extractTextContent(lastAssistant);
if (textContent != null && !textContent.isEmpty()) {
log.info("*** [Agent 思考] 模型返回文本: {}",
textContent.length() > 500
? textContent.substring(0, 500) + "... (已截断,总长度: " + textContent.length() + ")"
: textContent);
}
// 检查是否有工具调用
if (lastAssistant.getToolCalls() != null && !lastAssistant.getToolCalls().isEmpty()) {
hasToolCall = true;
log.info("*** [Agent 思考] 模型决定调用 {} 个工具:",
lastAssistant.getToolCalls().size());
lastAssistant.getToolCalls().forEach(toolCall -> {
log.info(" - 工具: {}, 参数: {}",
toolCall.name(),
toolCall.arguments());
});
log.info("*** [Agent 思考] 等待工具执行结果...");
} else {
log.info("*** [Agent 思考] 模型决定不调用工具");
log.info("*** [Agent 思考] 这是最终答案,准备返回给用户");
}
}
log.info("========================================");
// 更新 agent_step(afterModel:补全 model_output、耗时等)
if (sessionId != null) {
int stepIndex = stepCounters.getOrDefault(sessionId, 0);
String stepKey = sessionId + "_" + stepIndex;
Map<String, Object> pending = pendingSteps.remove(stepKey);
if (pending != null) {
try {
Long stepId = (Long) pending.get("stepId");
long startTime = (long) pending.get("startTime");
int durationMs = (int) (System.currentTimeMillis() - startTime);
AgentStep step = agentStepRepository.findById(stepId).orElse(null);
if (step != null) {
String thought = extractTextContent(lastAssistant);
if (thought != null && thought.length() > 2000) {
thought = thought.substring(0, 2000);
}
step.setThought(thought);
step.setHasToolCall(hasToolCall);
step.setDurationMs(durationMs);
if (lastAssistant != null) {
String outputSummary = buildModelOutputSummary(lastAssistant);
step.setModelOutput(outputSummary);
// 读取实际 token 用量(由 TokenTrackingChatModel 写入)
Integer tokenCount = TokenUsageHolder.get();
if (tokenCount != null) {
step.setTokenCount(tokenCount);
}
}
agentStepRepository.save(step);
log.debug("agent_step 已更新: sessionId={}, stepIndex={}, duration={}ms",
sessionId, stepIndex, durationMs);
}
} catch (Exception e) {
log.error("更新 agent_step 失败", e);
}
}
}
// 清理 token 上下文
TokenUsageHolder.clear();
return new AgentCommand(previousMessages);
}
/**
* 构建模型输入摘要(前 N 条消息的 role + 截断内容)
*/
private String buildModelInputSummary(List<Message> messages) {
StringBuilder sb = new StringBuilder();
int maxMessages = Math.min(messages.size(), 5);
for (int i = messages.size() - maxMessages; i < messages.size(); i++) {
Message msg = messages.get(i);
String role = getMessageRole(msg);
String content = msg.toString();
if (content.length() > 200) {
content = content.substring(0, 200) + "...";
}
sb.append("[").append(role).append("] ").append(content).append("\n");
}
String result = sb.toString();
if (result.length() > 500) {
result = result.substring(0, 500) + "...";
}
return result;
}
/**
* 构建模型输出摘要
*/
private String buildModelOutputSummary(AssistantMessage message) {
String text = extractTextContent(message);
if (text == null) {
text = "";
}
if (text.length() > 500) {
text = text.substring(0, 500) + "...";
}
StringBuilder sb = new StringBuilder();
sb.append("{\"text\":\"").append(escapeJson(text)).append("\"");
if (message.getToolCalls() != null && !message.getToolCalls().isEmpty()) {
sb.append(",\"toolCalls\":[");
for (int i = 0; i < message.getToolCalls().size(); i++) {
if (i > 0) sb.append(",");
sb.append("{\"name\":\"").append(escapeJson(message.getToolCalls().get(i).name()))
.append("\",\"arguments\":").append(message.getToolCalls().get(i).arguments()).append("}");
}
sb.append("]");
}
sb.append("}");
return sb.toString();
}
private String escapeJson(String s) {
if (s == null) return "";
return s.replace("\\", "\\\\")
.replace("\"", "\\\"")
.replace("\n", "\\n")
.replace("\r", "\\r")
.replace("\t", "\\t");
}
/**
* 提取 AssistantMessage 的文本内容
*/
private String extractTextContent(AssistantMessage message) {
if (message == null) return null;
try {
// 方法 1: 反射获取 text 字段
try {
java.lang.reflect.Field textField = message.getClass().getDeclaredField("text");
textField.setAccessible(true);
Object value = textField.get(message);
if (value != null) {
log.debug("通过 text 字段提取成功");
return value.toString();
}
} catch (NoSuchFieldException e) {
// 尝试下一种方法
}
// 方法 2: 反射获取 content 字段
try {
java.lang.reflect.Field contentField = message.getClass().getDeclaredField("content");
contentField.setAccessible(true);
Object value = contentField.get(message);
if (value != null) {
log.debug("通过 content 字段提取成功");
return value.toString();
}
} catch (NoSuchFieldException e) {
// 尝试下一种方法
}
// 方法 3: 调用 getText() 方法
try {
java.lang.reflect.Method getTextMethod = message.getClass().getMethod("getText");
Object value = getTextMethod.invoke(message);
if (value != null) {
log.debug("通过 getText() 方法提取成功");
return value.toString();
}
} catch (NoSuchMethodException e) {
// 尝试下一种方法
}
// 方法 4: 调用 getContent() 方法
try {
java.lang.reflect.Method getContentMethod = message.getClass().getMethod("getContent");
Object value = getContentMethod.invoke(message);
if (value != null) {
log.debug("通过 getContent() 方法提取成功");
return value.toString();
}
} catch (NoSuchMethodException e) {
// 方法不存在
}
// 方法 5: 打印类结构信息
log.warn("无法提取 AssistantMessage 文本内容,打印类信息:");
log.warn("类名: {}", message.getClass().getName());
log.warn("字段列表:");
for (java.lang.reflect.Field field : message.getClass().getDeclaredFields()) {
log.warn(" - {}: {}", field.getName(), field.getType().getSimpleName());
}
// 方法 6: toString() 兜底
String toString = message.toString();
if (toString != null && !toString.startsWith("AssistantMessage@")) {
log.debug("通过 toString() 提取");
return toString;
}
return null;
} catch (Exception e) {
log.error("提取 AssistantMessage 文本内容时出错", e);
return null;
}
}
/**
* 获取消息角色
*/
private String getMessageRole(Message message) {
if (message instanceof UserMessage) {
return "User(用户)";
} else if (message instanceof AssistantMessage) {
return "Assistant(模型)";
} else if (message instanceof ToolResponseMessage) {
return "Tool(工具返回)";
} else {
return message.getClass().getSimpleName();
}
}
}
@@ -0,0 +1,50 @@
package com.superbiz.agent.hook;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.ai.chat.model.ChatModel;
import org.springframework.ai.chat.model.ChatResponse;
import org.springframework.ai.chat.prompt.Prompt;
import reactor.core.publisher.Flux;
/**
* ChatModel 包装器 — 捕获每次模型调用的实际 token 用量
* 通过 TokenUsageHolder 传递给 AgentLoggingHook
*/
public class TokenTrackingChatModel implements ChatModel {
private static final Logger log = LoggerFactory.getLogger(TokenTrackingChatModel.class);
private final ChatModel delegate;
public TokenTrackingChatModel(ChatModel delegate) {
this.delegate = delegate;
}
@Override
public ChatResponse call(Prompt prompt) {
ChatResponse response = delegate.call(prompt);
captureTokenUsage(response);
return response;
}
@Override
public Flux<ChatResponse> stream(Prompt prompt) {
return delegate.stream(prompt);
}
private void captureTokenUsage(ChatResponse response) {
try {
if (response.getMetadata() == null || response.getMetadata().getUsage() == null) {
return;
}
var usage = response.getMetadata().getUsage();
Integer total = usage.getTotalTokens();
if (total != null && total > 0) {
TokenUsageHolder.set(total);
}
} catch (Exception e) {
log.debug("捕获 token 用量失败", e);
}
}
}
@@ -0,0 +1,22 @@
package com.superbiz.agent.hook;
/**
* Token 用量持有者(基于 ThreadLocal)
* ChatModel 调用后写入实际 token 数,AgentLoggingHook 读取
*/
public class TokenUsageHolder {
private static final ThreadLocal<Integer> TOKEN_COUNT = new ThreadLocal<>();
public static void set(Integer count) {
TOKEN_COUNT.set(count);
}
public static Integer get() {
return TOKEN_COUNT.get();
}
public static void clear() {
TOKEN_COUNT.remove();
}
}
@@ -0,0 +1,24 @@
package com.superbiz.agent.repository;
import com.superbiz.agent.domain.entity.AgentStep;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.util.List;
/**
* Agent 决策步骤 Repository
*/
@Repository
public interface AgentStepRepository extends JpaRepository<AgentStep, Long> {
/**
* 根据会话ID查询所有步骤(按步骤号排序)
*/
List<AgentStep> findBySessionIdOrderByStepIndex(String sessionId);
/**
* 统计某个会话的步骤数
*/
int countBySessionId(String sessionId);
}
@@ -1,73 +0,0 @@
package com.superbiz.agent.repository;
import com.superbiz.agent.domain.enums.DiagnosisStatus;
import com.superbiz.agent.domain.enums.FaultCategory;
import com.superbiz.agent.domain.entity.DiagnosisRecord;
import org.springframework.data.domain.Page;
import org.springframework.data.domain.Pageable;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Optional;
/**
* 诊断记录 Repository
*/
@Repository
public interface DiagnosisRecordRepository extends JpaRepository<DiagnosisRecord, Long> {
/**
* 根据诊断ID查询
*/
Optional<DiagnosisRecord> findByDiagnosisId(String diagnosisId);
/**
* 根据业务ID查询
*/
Optional<DiagnosisRecord> findByBusinessId(String businessId);
/**
* 根据链路追踪ID查询
*/
Optional<DiagnosisRecord> findByTraceId(String traceId);
/**
* 根据会话ID查询所有记录
*/
List<DiagnosisRecord> findBySessionId(String sessionId);
/**
* 根据故障类别和错误码查询
*/
List<DiagnosisRecord> findByFaultCategoryAndErrorCode(FaultCategory category, String errorCode);
/**
* 根据故障类别、故障源和错误码查询
*/
List<DiagnosisRecord> findByFaultCategoryAndFaultSourceAndErrorCode(
FaultCategory category, String faultSource, String errorCode);
/**
* 根据状态查询
*/
List<DiagnosisRecord> findByStatus(DiagnosisStatus status);
/**
* 根据时间范围查询(分页)
*/
Page<DiagnosisRecord> findByCreatedAtBetween(
LocalDateTime start, LocalDateTime end, Pageable pageable);
/**
* 根据故障类别和时间范围查询(分页)
*/
Page<DiagnosisRecord> findByFaultCategoryAndCreatedAtBetween(
FaultCategory category, LocalDateTime start, LocalDateTime end, Pageable pageable);
/**
* 查询有用反馈的高置信度记录(用于生成案例)
*/
List<DiagnosisRecord> findByFeedbackAndConfidenceGreaterThanEqual(String feedback, Integer confidence);
}
@@ -0,0 +1,12 @@
package com.superbiz.agent.repository;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.util.Optional;
@Repository
public interface DiagnosisSessionRepository extends JpaRepository<DiagnosisSession, Long> {
Optional<DiagnosisSession> findBySessionId(String sessionId);
}
@@ -0,0 +1,29 @@
package com.superbiz.agent.repository;
import com.superbiz.agent.domain.entity.ToolInvocation;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.util.List;
/**
* 工具调用明细 Repository
*/
@Repository
public interface ToolInvocationRepository extends JpaRepository<ToolInvocation, Long> {
/**
* 根据会话ID查询所有工具调用
*/
List<ToolInvocation> findBySessionId(String sessionId);
/**
* 根据工具名查询所有调用
*/
List<ToolInvocation> findByToolName(String toolName);
/**
* 根据会话ID和工具名查询
*/
List<ToolInvocation> findBySessionIdAndToolName(String sessionId, String toolName);
}
@@ -9,15 +9,25 @@ 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.AgentStep;
import com.superbiz.agent.domain.entity.AgentStep;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import com.superbiz.agent.hook.AgentLoggingHook;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.util.SessionContextHolder;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.ai.chat.messages.AssistantMessage;
import org.springframework.ai.tool.ToolCallback;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import com.superbiz.agent.config.AiOpsPromptProperties;
import com.superbiz.agent.tool.LookupKnowledgeTool;
import java.util.List;
import java.util.Optional;
import java.util.UUID;
/**
* AI Ops 智能运维服务
@@ -40,6 +50,18 @@ public class AiOpsService {
@Autowired(required = false) // Mock 模式下才注册
private QueryLogsTools queryLogsTools;
@Autowired
private LookupKnowledgeTool lookupKnowledgeTool;
@Autowired
private AiOpsPromptProperties promptProperties;
@Autowired
private DiagnosisSessionRepository diagnosisSessionRepository;
@Autowired
private AgentStepRepository agentStepRepository;
/**
* 执行 AI Ops 告警分析流程
*
@@ -51,34 +73,65 @@ public class AiOpsService {
public Optional<OverAllState> executeAiOpsAnalysis(ChatModel chatModel, ToolCallback[] toolCallbacks) throws GraphRunnerException {
logger.info("开始执行 AI Ops 多 Agent 协作流程");
// 构建 Planner 和 Executor Agent
ReactAgent plannerAgent = buildPlannerAgent(chatModel, toolCallbacks);
ReactAgent executorAgent = buildExecutorAgent(chatModel, toolCallbacks);
String sessionId = UUID.randomUUID().toString().substring(0, 8);
long startTime = System.currentTimeMillis();
// 构建 Supervisor Agent
SupervisorAgent supervisorAgent = SupervisorAgent.builder()
.name("ai_ops_supervisor")
.description("负责调度 Planner 与 Executor 的多 Agent 控制器")
.model(chatModel)
.systemPrompt(buildSupervisorSystemPrompt())
.subAgents(List.of(plannerAgent, executorAgent))
// 创建诊断会话
DiagnosisSession session = DiagnosisSession.builder()
.sessionId(sessionId)
.query("AI Ops 告警分析")
.status("RUNNING")
.agentFlow("AI_OPS")
.build();
diagnosisSessionRepository.save(session);
String taskPrompt = "你是企业级 SRE,接到了自动化告警排查任务。请结合工具调用,执行**规划→执行→再规划**的闭环,并最终按照固定模板输出《告警分析报告》。禁止编造虚假数据,如连续多次查询失败需诚实反馈无法完成的原因。";
// 设置 ThreadLocal 上下文(LookupKnowledgeTool 通过此获取 sessionId)
SessionContextHolder.setSessionId(sessionId);
logger.info("调用 Supervisor Agent 开始编排...");
try {
// 构建 Planner 和 Executor Agent(每个 Agent 各自带 Hook)
ReactAgent plannerAgent = buildPlannerAgent(chatModel, toolCallbacks);
ReactAgent executorAgent = buildExecutorAgent(chatModel, toolCallbacks);
Optional<OverAllState> stateOptional = supervisorAgent.invoke(taskPrompt);
// 构建 Supervisor Agent(不加 Hook)
SupervisorAgent supervisorAgent = SupervisorAgent.builder()
.name("ai_ops_supervisor")
.description("负责调度 Planner 与 Executor 的多 Agent 控制器")
.model(chatModel)
.systemPrompt(promptProperties.getSupervisor())
.subAgents(List.of(plannerAgent, executorAgent))
.build();
// 添加调试代码
if (stateOptional.isPresent()) {
OverAllState state = stateOptional.get();
logger.debug("Final State Keys: {}", state.data().keySet()); // 打印所有 key
logger.debug("Planner Plan: {}", state.value("planner_plan"));
logger.debug("Executor Feedback: {}", state.value("executor_feedback"));
String taskPrompt = "你是企业级 SRE,接到了自动化告警排查任务。请结合工具调用,执行**规划→执行→再规划**的闭环,并最终按照固定模板输出《告警分析报告》。禁止编造虚假数据,如连续多次查询失败需诚实反馈无法完成的原因。";
logger.info("调用 Supervisor Agent 开始编排...");
Optional<OverAllState> stateOptional = supervisorAgent.invoke(taskPrompt);
long duration = System.currentTimeMillis() - startTime;
// 更新诊断会话
session.setStatus(stateOptional.isPresent() ? "SUCCESS" : "FAILED");
session.setTotalDurationMs((int) duration);
backfillSessionMetrics(session);
diagnosisSessionRepository.save(session);
// 添加调试代码
if (stateOptional.isPresent()) {
OverAllState state = stateOptional.get();
logger.debug("Final State Keys: {}", state.data().keySet());
logger.debug("Planner Plan: {}", state.value("planner_plan"));
logger.debug("Executor Feedback: {}", state.value("executor_feedback"));
}
return stateOptional;
} catch (Exception e) {
session.setStatus("FAILED");
diagnosisSessionRepository.save(session);
throw e;
} finally {
SessionContextHolder.clear();
}
return stateOptional;
}
/**
@@ -113,9 +166,10 @@ public class AiOpsService {
.name("planner_agent")
.description("负责拆解告警、规划与再规划步骤")
.model(chatModel)
.systemPrompt(buildPlannerPrompt())
.systemPrompt(promptProperties.getPlanner())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "planner"))
.outputKey("planner_plan")
.build();
}
@@ -128,9 +182,10 @@ public class AiOpsService {
.name("executor_agent")
.description("负责执行 Planner 的首个步骤并及时反馈")
.model(chatModel)
.systemPrompt(buildExecutorPrompt())
.systemPrompt(promptProperties.getExecutor())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "executor"))
.outputKey("executor_feedback")
.build();
}
@@ -138,152 +193,37 @@ public class AiOpsService {
/**
* 动态构建方法工具数组
* 根据 cls.mock-enabled 决定是否包含 QueryLogsTools
* 工具顺序:知识库查询优先,日志查询次之,弃用工具最后
*/
private Object[] buildMethodToolsArray() {
if (queryLogsTools != null) {
// Mock 模式:包含 QueryLogsTools
return new Object[]{dateTimeTools, internalDocsTools, queryMetricsTools, queryLogsTools};
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools, queryLogsTools};
} else {
// 真实模式:不包含 QueryLogsTools(由 MCP 提供日志查询功能)
return new Object[]{dateTimeTools, internalDocsTools, queryMetricsTools};
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools};
}
}
/**
* 构建 Planner Agent 系统提示词
*/
private String buildPlannerPrompt() {
return """
你是 Planner Agent,同时承担 Replanner 角色,负责:
1. 读取当前输入任务 {input} 以及 Executor 的最近反馈 {executor_feedback}。
2. 分析 Prometheus 告警、日志、内部文档等信息,制定可执行的下一步步骤。
3. 在执行阶段,输出 JSON,包含 decision (PLAN|EXECUTE|FINISH)、step 描述、预期要调用的工具、以及必要的上下文。
4. 调用任何腾讯云日志/主题相关工具时,region 参数必须使用连字符格式(如 ap-guangzhou),若不确定请省略以使用默认值。
5. 严格禁止编造数据,只能引用工具返回的真实内容;如果连续 3 次调用同一工具仍失败或返回空结果,需停止该方向并在最终报告的结论部分说明"无法完成"的原因。
## 最终报告输出要求(CRITICAL)
当 decision=FINISH 时,你必须:
1. **不要输出 JSON 格式**
2. **直接输出完整的 Markdown 格式报告文本**
3. **报告必须严格遵循以下模板**:
```
# 告警分析报告
---
## 📋 活跃告警清单
| 告警名称 | 级别 | 目标服务 | 首次触发时间 | 最新触发时间 | 状态 |
|---------|------|----------|-------------|-------------|------|
| [告警1名称] | [级别] | [服务名] | [时间] | [时间] | 活跃 |
| [告警2名称] | [级别] | [服务名] | [时间] | [时间] | 活跃 |
---
## 🔍 告警根因分析1 - [告警名称]
### 告警详情
- **告警级别**: [级别]
- **受影响服务**: [服务名]
- **持续时间**: [X分钟]
### 症状描述
[根据监控指标描述症状]
### 日志证据
[引用查询到的关键日志]
### 根因结论
[基于证据得出的根本原因]
---
## 🛠️ 处理方案执行1 - [告警名称]
### 已执行的排查步骤
1. [步骤1]
2. [步骤2]
### 处理建议
[给出具体的处理建议]
### 预期效果
[说明预期的效果]
---
## 🔍 告警根因分析2 - [告警名称]
[如果有第2个告警,重复上述格式]
---
## 📊 结论
### 整体评估
[总结所有告警的整体情况]
### 关键发现
- [发现1]
- [发现2]
### 后续建议
1. [建议1]
2. [建议2]
### 风险评估
[评估当前风险等级和影响范围]
```
**重要提醒**:
- 最终输出必须是纯 Markdown 文本,不要包含 JSON 结构
- 不要使用 "finalReport": "..." 这样的格式
- 直接从 "# 告警分析报告" 开始输出
- 所有内容必须基于工具查询的真实数据,严禁编造
- 如果某个步骤失败,在结论中如实说明,不要跳过
""";
}
/** 从 agent_step 汇总指标回填 diagnosis_session */
private void backfillSessionMetrics(DiagnosisSession session) {
try {
List<AgentStep> steps = agentStepRepository.findBySessionIdOrderByStepIndex(session.getSessionId());
if (steps.isEmpty()) return;
/**
* 构建 Executor Agent 系统提示词
*/
private String buildExecutorPrompt() {
return """
你是 Executor Agent,负责读取 Planner 最新输出 {planner_plan},只执行其中的第一步。
- 确认步骤所需的工具与参数,尤其是 region 参数要使用连字符格式(ap-guangzhou);若 Planner 未给出则使用默认区域。
- 调用相应的工具并收集结果,如工具返回错误或空数据,需要将失败原因、请求参数一并记录,并停止进一步调用该工具(同一工具失败达到 3 次时应直接返回 FAILED)。
- 将日志、指标、文档等证据整理成结构化摘要,标注对应的告警名称或资源,方便 Planner 填充"告警根因分析 / 处理方案执行"章节。
- 以 JSON 形式返回执行状态、证据以及给 Planner 的建议,写入 executor_feedback,严禁编造未实际查询到的内容。
输出示例:
{
"status": "SUCCESS",
"summary": "近1小时未见 error 日志,仅有 info",
"evidence": "...",
"nextHint": "建议转向高占用进程"
}
""";
}
/**
* 构建 Supervisor Agent 系统提示词
*/
private String buildSupervisorSystemPrompt() {
return """
你是 AI Ops Supervisor,负责调度 planner_agent 与 executor_agent:
1. 当需要拆解任务或重新制定策略时,调用 planner_agent。
2. 当 planner_agent 输出 decision=EXECUTE 时,调用 executor_agent 执行第一步。
3. 根据 executor_agent 的反馈,评估是否需要再次调用 planner_agent,直到 decision=FINISH。
4. FINISH 后,确保向最终用户输出完整的《告警分析报告》,格式必须严格为:
告警分析报告\n---\n# 告警处理详情\n## 活跃告警清单\n## 告警根因分析N\n## 处理方案执行N\n## 结论。
5. 若步骤涉及腾讯云日志/主题工具,请确保使用连字符区域 ID(ap-guangzhou 等),或省略 region 以采用默认值。
6. 如果发现 Planner/Executor 在同一方向连续 3 次调用工具仍失败或没有数据,必须终止流程,直接输出"任务无法完成"的报告,明确告知失败原因,严禁凭空编造结果。
只允许在 planner_agent、executor_agent 与 FINISH 之间做出选择。
""";
int totalTokens = 0;
int stepCount = 0;
int toolCallCount = 0;
for (AgentStep s : steps) {
stepCount++;
if (s.getTokenCount() != null) totalTokens += s.getTokenCount();
if (Boolean.TRUE.equals(s.getHasToolCall())) toolCallCount++;
}
session.setTotalTokenCount(totalTokens);
session.setStepCount(stepCount);
session.setToolCallCount(toolCallCount);
} catch (Exception e) {
logger.warn("回填会话指标失败: sessionId={}", session.getSessionId(), e);
}
}
}
@@ -1,21 +1,41 @@
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.SupervisorAgent;
import com.alibaba.cloud.ai.graph.exception.GraphRunnerException;
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.TokenTrackingChatModel;
import com.superbiz.agent.hook.TokenUsageHolder;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.tool.LookupKnowledgeTool;
import com.superbiz.agent.util.QuestionComplexity;
import com.superbiz.agent.util.SessionContextHolder;
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.core.io.ClassPathResource;
import org.springframework.stereotype.Service;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.UUID;
/**
* 聊天服务
@@ -44,6 +64,40 @@ public class ChatService {
@Autowired
private ChatModel chatModel;
@Autowired
private LookupKnowledgeTool lookupKnowledgeTool;
@Autowired
private DiagnosisSessionRepository diagnosisSessionRepository;
@Autowired
private AgentStepRepository agentStepRepository;
/** 多 Agent Chat 的 Prompt */
private String chatPlannerPrompt;
private String chatExecutorPrompt;
@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);
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
*/
@@ -62,7 +116,7 @@ public class ChatService {
// 基础系统提示
systemPromptBuilder.append("你是一个专业的智能助手,可以获取当前时间、查询天气信息、搜索内部文档知识库,以及查询 Prometheus 告警信息。\n");
systemPromptBuilder.append("当用户询问时间相关问题时,**必须每次都调用 getCurrentDateTime 工具**,因为时间会不断变化。即使历史消息中有时间信息,也不要直接复用,必须重新查询最新时间。\n");
systemPromptBuilder.append("当用户需要查询公司内部文档、流程、最佳实践或技术指南时,使用 queryInternalDocs 工具。\n");
systemPromptBuilder.append("当用户需要查询公司内部文档、流程、最佳实践或技术指南时,使用 lookupKnowledgeTool 工具。\n");
systemPromptBuilder.append("当用户需要查询 Prometheus 告警、监控指标或系统告警状态时,使用 queryPrometheusAlerts 工具。\n");
systemPromptBuilder.append("当用户需要查询腾讯云日志时,请调用腾讯云mcp服务查询,默认查询地域ap-guangzhou,查询时间范围为近一个月。\n\n");
@@ -126,10 +180,10 @@ public class ChatService {
public Object[] buildMethodToolsArray() {
if (queryLogsTools != null) {
// Mock 模式:包含 QueryLogsTools
return new Object[]{dateTimeTools, internalDocsTools, queryMetricsTools, queryLogsTools};
return new Object[]{dateTimeTools, lookupKnowledgeTool};
} else {
// 真实模式:不包含 QueryLogsTools(由 MCP 提供日志查询功能)
return new Object[]{dateTimeTools, internalDocsTools, queryMetricsTools};
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools};
}
}
@@ -171,6 +225,7 @@ public class ChatService {
.systemPrompt(systemPrompt)
.methodTools(buildMethodToolsArray())
.tools(getToolCallbacks())
.hooks(new AgentLoggingHook(agentStepRepository, "intelligent_assistant"))
.build();
}
@@ -181,10 +236,209 @@ public class ChatService {
* @return AI 回复
*/
public String executeChat(ReactAgent agent, String question) throws GraphRunnerException {
logger.info("执行 ReactAgent.call() - 自动处理工具调用");
var response = agent.call(question);
String answer = response.getText();
logger.info("ReactAgent 对话完成,答案长度: {}", answer.length());
return answer;
logger.info("========================================");
logger.info("📝 用户问题: {}", question);
String sessionId = UUID.randomUUID().toString().substring(0, 8);
long startTime = System.currentTimeMillis();
// 创建诊断会话
DiagnosisSession session = DiagnosisSession.builder()
.sessionId(sessionId)
.query(question)
.status("RUNNING")
.agentFlow("CHAT")
.build();
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.setTotalDurationMs((int) duration);
backfillSessionMetrics(session);
diagnosisSessionRepository.save(session);
logger.info("⏱️ 总耗时: {} ms", duration);
logger.info("📏 输出长度: {} 字符", answer.length());
logger.info("========================================");
return answer;
} catch (Exception e) {
session.setStatus("FAILED");
diagnosisSessionRepository.save(session);
throw e;
} finally {
SessionContextHolder.clear();
}
}
/**
* 根据问题复杂度自动选择执行策略
* @param chatModel 聊天模型
* @param toolCallbacks 工具回调
* @param question 用户问题
* @param history 历史消息
* @return AI 回复
*/
public String executeChatWithStrategy(ChatModel chatModel, ToolCallback[] toolCallbacks,
String question, List<Map<String, String>> history) throws GraphRunnerException {
if (QuestionComplexity.isComplex(question)) {
logger.info("📊 问题判定为复杂,使用多 Agent(Planner + Executor)执行");
return executeChatComplex(chatModel, toolCallbacks, question, history);
} else {
logger.info("📊 问题判定为简单,使用单 Agent 执行");
String systemPrompt = buildSystemPrompt(history);
ReactAgent agent = createReactAgent(chatModel, systemPrompt);
return executeChat(agent, question);
}
}
/**
* 多 Agent 复杂对话执行(Planner + Executor + Supervisor)
*/
public String executeChatComplex(ChatModel chatModel, ToolCallback[] toolCallbacks,
String question, List<Map<String, String>> history) throws GraphRunnerException {
String sessionId = UUID.randomUUID().toString().substring(0, 8);
long startTime = System.currentTimeMillis();
DiagnosisSession session = DiagnosisSession.builder()
.sessionId(sessionId)
.query(question)
.status("RUNNING")
.agentFlow("CHAT")
.build();
diagnosisSessionRepository.save(session);
SessionContextHolder.setSessionId(sessionId);
try {
ReactAgent planner = buildChatPlannerAgent(chatModel, toolCallbacks, history);
ReactAgent executor = buildChatExecutorAgent(chatModel, toolCallbacks, history);
SupervisorAgent supervisor = SupervisorAgent.builder()
.name("chat_supervisor")
.description("负责调度 Planner 与 Executor 的多 Agent 控制器")
.model(chatModel)
.systemPrompt("你是一个智能任务调度器。分析用户问题,调用 Planner 拆解步骤,调用 Executor 执行各步骤。")
.subAgents(List.of(planner, executor))
.build();
Optional<OverAllState> stateOptional = supervisor.invoke(question);
long duration = System.currentTimeMillis() - startTime;
String answer = null;
if (stateOptional.isPresent()) {
// 从 state 中提取 Executor 的最终输出
OverAllState state = stateOptional.get();
Optional<AssistantMessage> executorOutput = state.value("executor_feedback")
.filter(AssistantMessage.class::isInstance)
.map(AssistantMessage.class::cast);
if (executorOutput.isPresent()) {
answer = executorOutput.get().getText();
}
}
if (answer == null || answer.isBlank()) {
answer = "抱歉,多 Agent 分析未能生成有效结论。";
}
session.setStatus("SUCCESS");
session.setTotalDurationMs((int) duration);
backfillSessionMetrics(session);
diagnosisSessionRepository.save(session);
logger.info("⏱️ 多 Agent 总耗时: {} ms", duration);
logger.info("📏 输出长度: {} 字符", answer.length());
return answer;
} catch (Exception e) {
session.setStatus("FAILED");
diagnosisSessionRepository.save(session);
logger.error("多 Agent 执行失败", e);
return "执行失败: " + e.getMessage();
} finally {
SessionContextHolder.clear();
}
}
private ReactAgent buildChatPlannerAgent(ChatModel chatModel, ToolCallback[] toolCallbacks,
List<Map<String, String>> history) {
StringBuilder prompt = new StringBuilder(chatPlannerPrompt);
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");
}
return ReactAgent.builder()
.name("chat_planner")
.description("负责拆解问题、规划步骤")
.model(chatModel)
.systemPrompt(prompt.toString())
// Planner 不注入工具,只能规划不能执行
.hooks(new AgentLoggingHook(agentStepRepository, "planner"))
.outputKey("planner_plan")
.build();
}
private ReactAgent buildChatExecutorAgent(ChatModel chatModel, ToolCallback[] toolCallbacks,
List<Map<String, String>> history) {
StringBuilder prompt = new StringBuilder(chatExecutorPrompt);
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");
}
return ReactAgent.builder()
.name("chat_executor")
.description("负责执行具体步骤并及时反馈")
.model(chatModel)
.systemPrompt(prompt.toString())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "executor"))
.outputKey("executor_feedback")
.build();
}
/** 从 agent_step 汇总 token、步数等指标回填 diagnosis_session */
private void backfillSessionMetrics(DiagnosisSession session) {
try {
List<com.superbiz.agent.domain.entity.AgentStep> steps =
agentStepRepository.findBySessionIdOrderByStepIndex(session.getSessionId());
if (steps.isEmpty()) return;
int totalTokens = 0;
int stepCount = 0;
int toolCallCount = 0;
for (var s : steps) {
stepCount++;
if (s.getTokenCount() != null) totalTokens += s.getTokenCount();
if (Boolean.TRUE.equals(s.getHasToolCall())) toolCallCount++;
}
session.setTotalTokenCount(totalTokens);
session.setStepCount(stepCount);
session.setToolCallCount(toolCallCount);
} catch (Exception e) {
logger.warn("回填会话指标失败: sessionId={}", session.getSessionId(), e);
}
}
}
@@ -56,7 +56,7 @@ public class DocumentChunkService {
}
/**
* 按照 Markdown 标题分割文档
* 按照 Markdown 标题分割文档,同时构建面包屑层级路径
*/
private List<Section> splitByHeadings(String content) {
List<Section> sections = new ArrayList<>();
@@ -65,20 +65,34 @@ public class DocumentChunkService {
Pattern headingPattern = Pattern.compile("^(#{1,6})\\s+(.+)$", Pattern.MULTILINE);
Matcher matcher = headingPattern.matcher(content);
// 标题层级栈:维护当前标题的完整路径
List<String> headingStack = new ArrayList<>();
int lastEnd = 0;
String currentTitle = null;
String currentBreadcrumb = null;
while (matcher.find()) {
int level = matcher.group(1).length(); // #→1, ##→2, ###→3 ...
String title = matcher.group(2).trim();
// 保存上一个章节
if (lastEnd < matcher.start()) {
String sectionContent = content.substring(lastEnd, matcher.start()).trim();
if (!sectionContent.isEmpty()) {
sections.add(new Section(currentTitle, sectionContent, lastEnd));
sections.add(new Section(
headingStack.isEmpty() ? null : headingStack.get(headingStack.size() - 1),
level,
currentBreadcrumb,
sectionContent,
lastEnd));
}
}
// 更新当前标题
currentTitle = matcher.group(2).trim();
// 维护层级栈:同级别或更高级别 → 弹出,低级 → 追加
while (!headingStack.isEmpty() && headingStack.size() >= level) {
headingStack.remove(headingStack.size() - 1);
}
headingStack.add(title);
currentBreadcrumb = String.join(" > ", headingStack);
lastEnd = matcher.start();
}
@@ -86,13 +100,18 @@ public class DocumentChunkService {
if (lastEnd < content.length()) {
String sectionContent = content.substring(lastEnd).trim();
if (!sectionContent.isEmpty()) {
sections.add(new Section(currentTitle, sectionContent, lastEnd));
sections.add(new Section(
headingStack.isEmpty() ? null : headingStack.get(headingStack.size() - 1),
headingStack.size(),
currentBreadcrumb,
sectionContent,
lastEnd));
}
}
// 如果没有找到任何标题,将整个文档作为一个章节
if (sections.isEmpty()) {
sections.add(new Section(null, content, 0));
sections.add(new Section(null, 0, null, content, 0));
}
return sections;
@@ -111,6 +130,7 @@ public class DocumentChunkService {
List<DocumentChunk> chunks = new ArrayList<>();
String content = section.content;
String title = section.title;
String breadcrumb = section.breadcrumb;
// 短章节直接作为一个分片(用 token 估算替代字符数做短路判断)
if (content.length() <= chunkConfig.getMaxSize()
@@ -121,6 +141,7 @@ public class DocumentChunkService {
.endOffset(section.startIndex + content.length())
.chunkIndex(startChunkIndex)
.title(title)
.breadcrumb(breadcrumb)
.build();
chunks.add(chunk);
return chunks;
@@ -155,7 +176,7 @@ public class DocumentChunkService {
logger.debug(" 触及硬上限 ({} tokens),强制切分", tokenCount + paraTokens);
chunkParaStart = saveChunkAndGetNextStart(
chunks, section, paraPositions,
chunkParaStart, i, title, chunkIndex);
chunkParaStart, i, title, breadcrumb, chunkIndex);
chunkIndex++;
String prevChunkContent = chunks.get(chunks.size() - 1).getContent();
@@ -168,7 +189,7 @@ public class DocumentChunkService {
// 安全切点:段落边界
chunkParaStart = saveChunkAndGetNextStart(
chunks, section, paraPositions,
chunkParaStart, i, title, chunkIndex);
chunkParaStart, i, title, breadcrumb, chunkIndex);
chunkIndex++;
// 新分片以重叠文本开头
@@ -194,6 +215,7 @@ public class DocumentChunkService {
.endOffset(section.startIndex + actualEnd)
.chunkIndex(chunkIndex)
.title(title)
.breadcrumb(breadcrumb)
.build();
chunks.add(chunk);
}
@@ -213,6 +235,7 @@ public class DocumentChunkService {
int fromPara,
int toPara,
String title,
String breadcrumb,
int chunkIndex) {
int actualStart = paraPositions.get(fromPara).start;
@@ -225,6 +248,7 @@ public class DocumentChunkService {
.endOffset(section.startIndex + actualEnd)
.chunkIndex(chunkIndex)
.title(title)
.breadcrumb(breadcrumb)
.build();
chunks.add(chunk);
@@ -392,12 +416,16 @@ public class DocumentChunkService {
* 章节数据类
*/
private static class Section {
String title;
String content;
int startIndex;
String title; // 最近一级标题名称
int level; // 标题级别(1-6),0=无标题
String breadcrumb; // 完整面包屑路径
String content; // 章节内容
int startIndex; // 在原文中的起始偏移
Section(String title, String content, int startIndex) {
Section(String title, int level, String breadcrumb, String content, int startIndex) {
this.title = title;
this.level = level;
this.breadcrumb = breadcrumb;
this.content = content;
this.startIndex = startIndex;
}
@@ -287,12 +287,12 @@ public class DocumentManagementService {
*/
private FaultCategory parseFaultCategory(String category) {
if (category == null || category.isBlank()) {
return FaultCategory.EXTERNAL_API;
return FaultCategory.GENERAL;
}
try {
return FaultCategory.valueOf(category.toUpperCase());
} catch (IllegalArgumentException e) {
return FaultCategory.EXTERNAL_API;
return FaultCategory.GENERAL;
}
}
@@ -0,0 +1,347 @@
package com.superbiz.agent.service;
import com.superbiz.agent.domain.entity.ApiDocument;
import com.superbiz.agent.domain.enums.FaultCategory;
import com.superbiz.agent.repository.ApiDocumentRepository;
import com.superbiz.agent.dto.KnowledgeEntry;
import com.superbiz.agent.dto.Frontmatter;
import com.superbiz.agent.dto.DocumentChunk;
import lombok.Data;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.io.IOException;
import java.nio.file.*;
import java.nio.file.attribute.BasicFileAttributes;
import java.time.LocalDateTime;
import java.util.*;
import java.util.stream.Collectors;
import java.util.stream.Collectors;
/**
* 知识库初始化服务
* 负责批量导入 knowledge_base 目录下的文档到数据库和 Milvus
*/
@Service
public class KnowledgeBaseInitService {
private static final Logger logger = LoggerFactory.getLogger(KnowledgeBaseInitService.class);
@Value("${knowledge.base-path:knowledge_base}")
private String knowledgeBasePath;
@Autowired
private ApiDocumentRepository apiDocumentRepository;
@Autowired
private FrontmatterParser frontmatterParser;
@Autowired
private DocumentChunkService documentChunkService;
@Autowired
private VectorIndexService vectorIndexService;
@Autowired
private VectorEmbeddingService vectorEmbeddingService;
@Autowired
private KnowledgeIndexService knowledgeIndexService;
/**
* 初始化知识库
*
* @param force 是否强制重新导入(跳过去重检查)
* @return 初始化结果
*/
@Transactional(rollbackFor = Exception.class)
public InitResult initializeKnowledgeBase(boolean force) {
logger.info("开始初始化知识库: basePath={}, force={}", knowledgeBasePath, force);
InitResult result = new InitResult();
Path baseDir = Paths.get(knowledgeBasePath);
if (!Files.exists(baseDir)) {
logger.error("知识库目录不存在: {}", knowledgeBasePath);
throw new RuntimeException("知识库目录不存在: " + knowledgeBasePath);
}
// 1. 扫描所有 Markdown 文件
List<Path> markdownFiles = scanMarkdownFiles(baseDir);
result.setScanned(markdownFiles.size());
logger.info("扫描到 {} 个 Markdown 文件", markdownFiles.size());
// 2. 如果非强制模式,获取已存在的文档(用于去重)
Set<String> existingFilePaths = new HashSet<>();
if (!force) {
existingFilePaths = apiDocumentRepository.findAll().stream()
.map(ApiDocument::getFilePath)
.collect(Collectors.toSet());
logger.info("已存在 个文档记录", existingFilePaths.size());
}
// 3. 逐个处理文档
for (Path file : markdownFiles) {
String relativePath = baseDir.relativize(file).toString().replace("\\", "/");
try {
// 去重检查
if (!force && existingFilePaths.contains(relativePath)) {
logger.debug("跳过已存在的文档: {}", relativePath);
result.incrementSkipped();
result.addDetail(relativePath, "已存在,跳过");
continue;
}
// 解析文档
String content = Files.readString(file);
Frontmatter frontmatter = frontmatterParser.parse(content);
if (frontmatter == null) {
logger.warn("文档格式无效: {}, frontmatter 解析失败", relativePath);
result.incrementFailed();
result.addDetail(relativePath, "格式无效: frontmatter 解析失败");
continue;
}
// 提取字段
String title = frontmatter.getTitle();
String summary = frontmatter.getSummary();
String category = frontmatter.getCategory() != null ? frontmatter.getCategory() : "general";
List<String> keywords = frontmatter.getKeywords();
if (title == null || title.isBlank()) {
logger.warn("文档缺少标题: {}", relativePath);
result.incrementFailed();
result.addDetail(relativePath, "缺少标题");
continue;
}
// 保存到数据库
ApiDocument document = saveToDatabase(relativePath, title, summary, category, content, keywords);
// 提取文档正文(去除 frontmatter)
String body = extractBody(content);
// 文档分块
List<DocumentChunk> chunks = documentChunkService.chunkDocument(body, relativePath);
logger.debug("文档分块完成: {} -> {} 个 chunk", relativePath, chunks.size());
// 上传到 Milvus
try {
vectorIndexService.indexDocumentChunks(document.getDocId(), chunks, category);
document.setStatus("INDEXED");
document.setChunkCount(chunks.size());
document.setIndexedAt(LocalDateTime.now());
apiDocumentRepository.save(document);
logger.info("文档已索引到 Milvus: {} (docId={}, chunks={})",
title, document.getDocId(), chunks.size());
} catch (Exception e) {
logger.error("上传到 Milvus 失败: {}", relativePath, e);
document.setStatus("FAILED");
document.setErrorMessage(e.getMessage());
apiDocumentRepository.save(document);
result.incrementFailed();
result.addDetail(relativePath, "Milvus 索引失败: " + e.getMessage());
continue; // 跳过该文档,继续处理下一个
}
// 添加到 L0 内存索引
KnowledgeEntry entry = KnowledgeEntry.builder()
.filePath(relativePath)
.title(title)
.keywords(keywords)
.summary(summary)
.category(category)
.build();
knowledgeIndexService.addToIndex(entry);
result.incrementInserted();
result.addDetail(relativePath, "导入成功(L0+L1)");
logger.info("文档导入成功: {} -> {} (L0+L1 索引已更新)", relativePath, title);
} catch (Exception e) {
logger.error("处理文档失败: {}", relativePath, e);
result.incrementFailed();
result.addDetail(relativePath, "处理失败: " + e.getMessage());
}
}
logger.info("知识库初始化完成: 扫描={}, 跳过={}, 新增={}, 失败={}",
result.getScanned(), result.getSkipped(), result.getInserted(), result.getFailed());
return result;
}
/**
* 获取知识库统计信息
*/
public Stats getStats() {
Stats stats = new Stats();
// 数据库中的文档数量
long totalDocuments = apiDocumentRepository.count();
stats.setTotalDocuments(totalDocuments);
// L0 索引中的文档数量
int indexSize = knowledgeIndexService.getIndexSize();
logger.debug("L0 索引大小: {}", indexSize);
// 按分类统计(从 fault_category 字段读取)
Map<String, Long> categoryCount = apiDocumentRepository.findAll().stream()
.collect(Collectors.groupingBy(
doc -> doc.getFaultCategory() != null ? doc.getFaultCategory().name() : "GENERAL",
Collectors.counting()
));
stats.setCategoryCount(categoryCount);
// Milvus 中的向量数量(需要实现)
// TODO: 查询 Milvus collection 的实体数量
stats.setTotalVectors(0L);
return stats;
}
/**
* 扫描目录下所有 Markdown 文件
*/
private List<Path> scanMarkdownFiles(Path baseDir) {
List<Path> files = new ArrayList<>();
try {
Files.walkFileTree(baseDir, new SimpleFileVisitor<Path>() {
@Override
public FileVisitResult visitFile(Path file, BasicFileAttributes attrs) {
if (file.toString().endsWith(".md")) {
files.add(file);
}
return FileVisitResult.CONTINUE;
}
@Override
public FileVisitResult visitFileFailed(Path file, IOException exc) {
logger.warn("访问文件失败: {}", file, exc);
return FileVisitResult.CONTINUE;
}
});
} catch (IOException e) {
logger.error("扫描目录失败: {}", baseDir, e);
throw new RuntimeException("扫描目录失败", e);
}
return files;
}
/**
* 保存文档到数据库
*/
private ApiDocument saveToDatabase(String filePath, String title, String summary,
String category, String content, List<String> keywords) {
ApiDocument document = new ApiDocument();
document.setDocId(UUID.randomUUID().toString());
document.setFileName(Paths.get(filePath).getFileName().toString());
document.setFilePath(filePath);
document.setApiName(title); // 使用 title 作为 apiName
document.setStatus("PENDING"); // 初始状态为 PENDING,索引成功后更新为 INDEXED
// 映射 category 到 FaultCategory 枚举
FaultCategory faultCategory = FaultCategory.fromString(category);
document.setFaultCategory(faultCategory);
// 将 frontmatter 信息保存到 metadata(JSON 格式)
String metadataJson = String.format(
"{\"title\":\"%s\",\"summary\":\"%s\",\"category\":\"%s\",\"keywords\":%s}",
escapeJson(title),
escapeJson(summary),
escapeJson(category),
"[\"" + String.join("\",\"", keywords.stream().map(this::escapeJson).toArray(String[]::new)) + "\"]"
);
document.setMetadata(metadataJson);
document.setFileSize((long) content.length());
return apiDocumentRepository.save(document);
}
/**
* JSON 转义
*/
private String escapeJson(String str) {
if (str == null) {
return "";
}
return str.replace("\\", "\\\\")
.replace("\"", "\\\"")
.replace("\n", "\\n")
.replace("\r", "\\r");
}
/**
* 提取文档正文(去除 frontmatter)
*/
private String extractBody(String content) {
if (!content.trim().startsWith("---")) {
return content;
}
int firstEnd = content.indexOf("---", 3);
if (firstEnd == -1) {
return content;
}
int secondEnd = content.indexOf("---", firstEnd + 3);
if (secondEnd == -1) {
return content.substring(firstEnd + 3).trim();
}
return content.substring(secondEnd + 3).trim();
}
// ==================== 数据模型 ====================
/**
* 初始化结果
*/
@Data
public static class InitResult {
private int scanned; // 扫描到的文件数量
private int skipped; // 跳过的文件数量(已存在)
private int inserted; // 成功导入的文件数量
private int failed; // 失败的文件数量
private Map<String, String> details = new LinkedHashMap<>(); // 详细信息
public void incrementSkipped() {
this.skipped++;
}
public void incrementInserted() {
this.inserted++;
}
public void incrementFailed() {
this.failed++;
}
public void addDetail(String filePath, String message) {
this.details.put(filePath, message);
}
}
/**
* 统计信息
*/
@Data
public static class Stats {
private long totalDocuments; // 数据库中的文档总数
private long totalVectors; // Milvus 中的向量总数
private Map<String, Long> categoryCount; // 按分类统计
}
}
@@ -1,5 +1,7 @@
package com.superbiz.agent.service;
import com.superbiz.agent.domain.entity.ApiDocument;
import com.superbiz.agent.repository.ApiDocumentRepository;
import com.superbiz.agent.dto.Frontmatter;
import com.superbiz.agent.dto.KnowledgeEntry;
import lombok.extern.slf4j.Slf4j;
@@ -12,6 +14,8 @@ import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.stream.Collectors;
@@ -25,11 +29,11 @@ import java.util.stream.Stream;
@Service
public class KnowledgeIndexService {
@Value("${knowledge.base-path}")
@Value("${knowledge.base-path:knowledge_base}")
private String knowledgeBasePath;
@Autowired
private FrontmatterParser frontmatterParser;
private ApiDocumentRepository apiDocumentRepository;
/**
* 内存索引(线程安全)
@@ -37,88 +41,108 @@ public class KnowledgeIndexService {
private final List<KnowledgeEntry> knowledgeIndex = new CopyOnWriteArrayList<>();
/**
* 启动时扫描知识库目录,构建索引
* 启动时从数据库加载索引
*/
@PostConstruct
public void loadIndex() {
log.info("开始扫描知识库目录: {}", knowledgeBasePath);
log.info("开始从数据库加载知识库索引");
try {
Path basePath = Paths.get(knowledgeBasePath);
// 从数据库读取所有已索引的文档
List<ApiDocument> documents = apiDocumentRepository.findAll();
// 目录不存在时自动创建
if (!Files.exists(basePath)) {
Files.createDirectories(basePath);
log.info("知识库目录已创建: {}", basePath.toAbsolutePath());
int loaded = 0;
for (ApiDocument doc : documents) {
try {
// 从 metadata JSON 中提取信息
KnowledgeEntry entry = parseDocumentToEntry(doc);
if (entry != null) {
knowledgeIndex.add(entry);
loaded++;
}
} catch (Exception e) {
log.warn("解析文档失败: docId={}, error={}", doc.getDocId(), e.getMessage());
}
}
// 递归扫描 .md 文件
try (Stream<Path> paths = Files.walk(basePath)) {
paths.filter(p -> p.toString().endsWith(".md"))
.forEach(this::indexFile);
}
log.info("知识库索引加载完成,共 {} 个文档", loaded);
log.info("知识库索引加载完成,共 {} 个文档", knowledgeIndex.size());
} catch (IOException e) {
} catch (Exception e) {
log.error("知识库索引加载失败", e);
}
}
/**
* 索引单个文件
*
* @param filePath 文件路径
* 将 ApiDocument 转换为 KnowledgeEntry
*/
private void indexFile(Path filePath) {
private KnowledgeEntry parseDocumentToEntry(ApiDocument doc) {
if (doc.getMetadata() == null || doc.getMetadata().isEmpty()) {
return null;
}
try {
// 读取文件内容
String content = Files.readString(filePath);
// 简单的 JSON 解析
String metadata = doc.getMetadata();
// 解析 frontmatter
Frontmatter frontmatter = frontmatterParser.parse(content);
if (frontmatter == null) {
log.debug("跳过文件(无有效 frontmatter): {}", filePath);
return;
}
String title = extractJsonValue(metadata, "title");
String summary = extractJsonValue(metadata, "summary");
String category = extractJsonValue(metadata, "category");
List<String> keywords = extractJsonArray(metadata, "keywords");
// 提取 category(从路径中获取)
String category = extractCategoryFromPath(filePath.toString());
// 构建索引条目
KnowledgeEntry entry = KnowledgeEntry.builder()
.filePath(filePath.toString())
.title(frontmatter.getTitle())
.keywords(frontmatter.getKeywords())
.summary(frontmatter.getSummary())
return KnowledgeEntry.builder()
.filePath(doc.getFilePath())
.title(title != null ? title : doc.getApiName())
.keywords(keywords)
.summary(summary)
.category(category)
.sections(frontmatter.getSections())
.build();
knowledgeIndex.add(entry);
log.debug("文档已加入索引: title={}, filePath={}", entry.getTitle(), filePath);
} catch (IOException e) {
log.warn("读取文件失败: {}", filePath, e);
} catch (Exception e) {
log.warn("解析 metadata 失败: {}", doc.getDocId(), e);
return null;
}
}
/**
* 从文件路径中提取 category
* 例如:knowledge_base/api/test.md -> api
* 从 JSON 字符串中提取值
*/
private String extractCategoryFromPath(String filePath) {
String normalized = filePath.replace("\\", "/");
String[] parts = normalized.split("/");
// 查找 knowledge_base 后的第一个目录
for (int i = 0; i < parts.length - 1; i++) {
if (parts[i].equals("knowledge_base") && i + 1 < parts.length) {
return parts[i + 1];
}
private String extractJsonValue(String json, String key) {
String pattern = "\"" + key + "\":\"";
int startIndex = json.indexOf(pattern);
if (startIndex == -1) {
return null;
}
return "default";
startIndex += pattern.length();
int endIndex = json.indexOf("\"", startIndex);
if (endIndex == -1) {
return null;
}
return json.substring(startIndex, endIndex);
}
/**
* 从 JSON 字符串中提取数组
*/
private List<String> extractJsonArray(String json, String key) {
String pattern = "\"" + key + "\":[";
int startIndex = json.indexOf(pattern);
if (startIndex == -1) {
return Collections.emptyList();
}
startIndex += pattern.length();
int endIndex = json.indexOf("]", startIndex);
if (endIndex == -1) {
return Collections.emptyList();
}
String arrayContent = json.substring(startIndex, endIndex);
return Arrays.stream(arrayContent.split(","))
.map(s -> s.trim().replaceAll("^\"|\"$", ""))
.filter(s -> !s.isEmpty())
.collect(Collectors.toList());
}
/**
@@ -174,13 +198,15 @@ public class KnowledgeIndexService {
/**
* 读取文档内容
*
* @param filePath 文件路径
* @param filePath 文件相对路径(如 api/payment-errors.md)
* @param maxChars 最大字符数
* @return 文档内容(前 maxChars 字符),失败返回 null
*/
public String readDocument(String filePath, int maxChars) {
try {
String content = Files.readString(Paths.get(filePath));
// 拼接完整路径:knowledge_base + 相对路径
Path fullPath = Paths.get(knowledgeBasePath, filePath);
String content = Files.readString(fullPath);
if (content.length() > maxChars) {
return content.substring(0, maxChars) + "...";
@@ -189,7 +215,7 @@ public class KnowledgeIndexService {
return content;
} catch (IOException e) {
log.error("读取文档失败: {}", filePath, e);
log.error("读取文档失败: {}/{}", knowledgeBasePath, filePath, e);
return null;
}
}
@@ -269,6 +269,11 @@ public class VectorIndexService {
metadata.put("title", chunk.getTitle());
}
// 面包屑导航(完整标题层级路径)
if (chunk.getBreadcrumb() != null && !chunk.getBreadcrumb().isEmpty()) {
metadata.put("breadcrumb", chunk.getBreadcrumb());
}
// 文档类别
metadata.put("category", category != null && !category.isBlank() ? category : "upload");
@@ -360,6 +365,11 @@ public class VectorIndexService {
metadata.put("title", chunk.getTitle());
}
// 面包屑导航(完整标题层级路径)
if (chunk.getBreadcrumb() != null && !chunk.getBreadcrumb().isEmpty()) {
metadata.put("breadcrumb", chunk.getBreadcrumb());
}
return metadata;
}
@@ -1,14 +1,18 @@
package com.superbiz.agent.tool;
import com.superbiz.agent.domain.entity.ToolInvocation;
import com.superbiz.agent.dto.*;
import com.superbiz.agent.repository.ToolInvocationRepository;
import com.superbiz.agent.service.KnowledgeIndexService;
import com.superbiz.agent.service.VectorSearchService;
import com.superbiz.agent.util.SessionContextHolder;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.tool.annotation.Tool;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.stream.Collectors;
/**
* 知识库查询工具
@@ -24,61 +28,232 @@ public class LookupKnowledgeTool {
@Autowired
private VectorSearchService vectorSearchService;
@Autowired
private ToolInvocationRepository toolInvocationRepository;
/**
* 查询知识库文档
*
* @param query 查询关键词
* @return 查询结果
*/
@Tool(description = "查询知识库文档。优先精确匹配关键词,未命中或多个匹配时自动补充语义相关片段。" +
"参数 query: 查询关键词,例如 'ERR_TIMEOUT'、'支付网关超时'")
@Tool(description = "查询内部知识库文档,获取错误码定义、接口文档、排障步骤、配置说明等背景信息。" +
"采用两阶段检索:L0 精确匹配关键词(< 10ms),L1 语义检索补充(200-500ms)。" +
"IMPORTANT: 遇到错误码、接口名、配置项、排障问题时,优先使用此工具。" +
"支持的查询场景:" +
"1) 错误码定义 - 查询错误码的含义和处理方法,例如 'ERR_TIMEOUT'、'ERR_CONNECTION_REFUSED';" +
"2) 接口文档 - 查询 API 接口定义、参数说明、返回格式,例如 'payment-gateway'、'/api/v1/orders';" +
"3) 排障步骤 - 查询故障诊断流程、最佳实践,例如 '支付超时排查'、'数据库连接池配置';" +
"4) 配置说明 - 查询系统配置、中间件参数,例如 'HikariCP'、'Redis 集群配置'。" +
"参数 query: 查询关键词或描述")
public LookupResult lookupKnowledge(String query) {
// 生成请求ID用于追踪
String requestId = java.util.UUID.randomUUID().toString().substring(0, 8);
long startTime = System.currentTimeMillis();
log.info("[{}] 收到知识库查询请求: query={}", requestId, query);
log.info("========================================");
log.info(">>> [工具调用] lookup_knowledge");
log.info(">>> 参数: query = \"{}\"", query);
log.info(">>> RequestId: {}", requestId);
log.info("----------------------------------------");
// Step 1: L0 精确匹配
long l0Start = System.currentTimeMillis();
List<KnowledgeEntry> l0Matches = knowledgeIndexService.exactMatch(query);
long l0Time = System.currentTimeMillis() - l0Start;
log.info("[{}] L0精确匹配完成: matches={}, time={}ms", requestId, l0Matches.size(), l0Time);
log.info("[L0 精确匹配] 完成: matches={}, time={}ms", l0Matches.size(), l0Time);
if (!l0Matches.isEmpty()) {
log.info("[L0 精确匹配] 找到文档:");
for (int i = 0; i < Math.min(3, l0Matches.size()); i++) {
KnowledgeEntry entry = l0Matches.get(i);
log.info(" - [{}] 标题: {}, 路径: {}", i+1, entry.getTitle(), entry.getFilePath());
}
}
// Step 2: 判断是否高置信度(唯一匹配)
boolean highConfidence = (l0Matches.size() == 1);
log.debug("[{}] 置信度判断: highConfidence={}, reason={}",
requestId, highConfidence, highConfidence ? "唯一匹配" : "多个或零个匹配");
log.info("[置信度判断] highConfidence={}, reason={}",
highConfidence, highConfidence ? "唯一匹配" : "多个或零个匹配");
// Step 3: L1 条件调用
List<VectorSearchService.SearchResult> l1Results = null;
if (!highConfidence) {
log.info("[{}] L0非唯一匹配,触发L1语义检索", requestId);
log.info("[L1 语义检索] L0非唯一匹配,触发L1语义检索...");
long l1Start = System.currentTimeMillis();
l1Results = vectorSearchService.searchSimilarDocuments(query, 3, null);
long l1Time = System.currentTimeMillis() - l1Start;
log.info("[{}] L1语义检索完成: matches={}, time={}ms",
requestId, l1Results != null ? l1Results.size() : 0, l1Time);
log.info("[L1 语义检索] 完成: matches={}, time={}ms",
l1Results != null ? l1Results.size() : 0, l1Time);
if (l1Results != null && !l1Results.isEmpty()) {
log.info("[L1 语义检索] 找到文档:");
for (int i = 0; i < Math.min(3, l1Results.size()); i++) {
VectorSearchService.SearchResult result = l1Results.get(i);
log.info(" - [{}] 文档ID: {}, 相似度得分: {}", i+1, result.getId(), result.getScore());
}
}
} else {
log.debug("[{}] L0唯一匹配,跳过L1检索", requestId);
log.info("[L1 语义检索] L0唯一匹配,跳过L1检索");
}
// Step 4: 组装结果
LookupResult result = buildResult(l0Matches, l1Results, highConfidence);
// 记录完整结果
// 记录结构化结果摘要(替代原始 MD 内容预览)
long totalTime = System.currentTimeMillis() - startTime;
log.info("[{}] 查询完成: found={}, hasL0={}, hasL1={}, confidence={}, totalTime={}ms",
requestId,
result.isFound(),
result.getPrimary() != null,
result.getSupplement() != null,
result.getPrimary() != null ? result.getPrimary().getConfidence() : "N/A",
totalTime);
log.info("----------------------------------------");
log.info("<<< [工具返回] lookup_knowledge");
log.info("<<< 结果: found={}, 耗时: {}ms (L0={}ms, L1={}ms)",
result.isFound(), totalTime, l0Time,
l1Results != null ? System.currentTimeMillis() - startTime - l0Time : 0);
// L0 精确匹配摘要
if (!l0Matches.isEmpty()) {
KnowledgeEntry top = l0Matches.get(0);
log.info("<<< [L0 主结果] 标题: {}", top.getTitle());
log.info("<<< [L0 主结果] 来源: {}", top.getFilePath());
if (top.getSummary() != null) {
log.info("<<< [L0 主结果] 摘要: {}", top.getSummary());
}
if (top.getKeywords() != null && !top.getKeywords().isEmpty()) {
log.info("<<< [L0 主结果] 关键词: {}", String.join(", ", top.getKeywords()));
}
// 内容概况:长度 + 章节数
String content = result.getPrimary() != null ? result.getPrimary().getContent() : null;
if (content != null) {
int headingCount = countMdHeadings(content);
log.info("<<< [L0 主结果] 内容: {} 字符, {} 个章节",
content.length(), headingCount);
}
}
// L1 语义检索摘要
if (l1Results != null && !l1Results.isEmpty()) {
VectorSearchService.SearchResult topL1 = l1Results.get(0);
log.info("<<< [L1 补充] 来源: {}", topL1.getMetadata() != null ? topL1.getMetadata() : topL1.getId());
log.info("<<< [L1 补充] 相似度: {}", String.format("%.4f", topL1.getScore()));
if (topL1.getContent() != null) {
String snippet = extractFirstMeaningfulLine(topL1.getContent(), 120);
log.info("<<< [L1 补充] 内容片段: {}", snippet);
log.info("<<< [L1 补充] 片段长度: {} 字符", topL1.getContent().length());
}
}
log.info("========================================");
// 记录 tool_invocation(持久化检索明细)
saveToolInvocation(query, l0Matches, l1Results, highConfidence, startTime, result);
return result;
}
/**
* 保存工具调用明细到 tool_invocation 表
*/
private void saveToolInvocation(String query, List<KnowledgeEntry> l0Matches,
List<VectorSearchService.SearchResult> l1Results,
boolean highConfidence, long startTime, LookupResult result) {
try {
String sessionId = SessionContextHolder.getSessionId();
if (sessionId == null) return; // 非会话上下文不记录
boolean hasL0 = l0Matches != null && !l0Matches.isEmpty();
boolean hasL1 = l1Results != null && !l1Results.isEmpty();
long duration = System.currentTimeMillis() - startTime;
String layer;
String outputPreview = null;
int outputLength = 0;
int l0Count = 0;
int l1Count = 0;
boolean truncated = false;
if (hasL0 && !highConfidence) {
layer = "L0+L1";
l0Count = l0Matches.size();
l1Count = l1Results.size();
} else if (hasL0) {
layer = "L0";
l0Count = l0Matches.size();
} else if (hasL1) {
layer = "L1";
l1Count = l1Results.size();
} else {
layer = null;
}
// 拼接 output_preview(前500字符)
if (result != null && result.getPrimary() != null && result.getPrimary().getContent() != null) {
String content = result.getPrimary().getContent();
outputLength = content.length();
if (content.length() > 500) {
outputPreview = content.substring(0, 500) + "...";
truncated = true;
} else {
outputPreview = content;
}
} else if (l1Results != null && !l1Results.isEmpty() && l1Results.get(0).getContent() != null) {
String content = l1Results.get(0).getContent();
outputLength = content.length();
if (content.length() > 500) {
outputPreview = content.substring(0, 500) + "...";
truncated = true;
} else {
outputPreview = content;
}
}
// 构建检索明细 JSON
StringBuilder details = new StringBuilder("{");
if (hasL0) {
details.append("\"l0_titles\":[");
for (int i = 0; i < Math.min(3, l0Matches.size()); i++) {
if (i > 0) details.append(",");
details.append("\"").append(escapeJson(l0Matches.get(i).getTitle())).append("\"");
}
details.append("]");
}
if (hasL1) {
if (hasL0) details.append(",");
details.append("\"l1_scores\":[");
for (int i = 0; i < Math.min(3, l1Results.size()); i++) {
if (i > 0) details.append(",");
details.append(l1Results.get(i).getScore());
}
details.append("]");
}
details.append("}");
ToolInvocation inv = ToolInvocation.builder()
.sessionId(sessionId)
.toolName("lookup_knowledge")
.inputParams("{\"query\":\"" + escapeJson(query) + "\"}")
.outputPreview(outputPreview)
.outputLength(outputLength)
.retrievalLayer(layer)
.l0MatchCount(hasL0 ? l0Count : null)
.l1MatchCount(hasL1 ? l1Count : null)
.isTruncated(truncated)
.retrievalDetails(details.toString())
.durationMs((int) duration)
.success(true)
.build();
toolInvocationRepository.save(inv);
log.debug("tool_invocation 已保存: sessionId={}, layer={}, duration={}ms", sessionId, layer, duration);
} catch (Exception e) {
log.error("保存 tool_invocation 失败", e);
}
}
private String escapeJson(String s) {
if (s == null) return "";
return s.replace("\\", "\\\\")
.replace("\"", "\\\"")
.replace("\n", "\\n")
.replace("\r", "\\r")
.replace("\t", "\\t");
}
/**
* 组装查询结果
*
@@ -98,7 +273,13 @@ public class LookupKnowledgeTool {
PrimaryResult primary = null;
if (l0Matches != null && !l0Matches.isEmpty()) {
KnowledgeEntry first = l0Matches.get(0);
String content = knowledgeIndexService.readDocument(first.getFilePath(), 2000);
boolean hasL1 = l1Results != null && !l1Results.isEmpty();
// 场景决策:唯一匹配或 L1 无结果 → LLM 需要正文内容;多匹配且有 L1 → 只需元数据
boolean needFullContent = highConfidence || !hasL1;
String content = needFullContent
? buildCompactSummary(first)
: buildMetadataOnlySummary(first);
if (content != null) {
primary = PrimaryResult.builder()
@@ -135,4 +316,118 @@ public class LookupKnowledgeTool {
return builder.build();
}
/**
* 统计 MD 文档中的章节数(二级标题 ## 数量)
*/
private int countMdHeadings(String content) {
if (content == null) return 0;
return (int) content.lines()
.filter(l -> l.trim().startsWith("##"))
.count();
}
/**
* 构建紧凑文档摘要(替代原始 MD 全文,节省上下文窗口)
* 组合:title/summary + 章节结构 + 正文片段(~500 字符)
*/
private String buildCompactSummary(KnowledgeEntry entry) {
String rawContent = knowledgeIndexService.readDocument(entry.getFilePath(), 2000);
if (rawContent == null) return null;
// 跳过 YAML frontmatter 得到正文
String body = rawContent;
if (body.startsWith("---")) {
int end = body.indexOf("---", 3);
if (end != -1) {
body = body.substring(end + 3).trim();
}
}
StringBuilder sb = new StringBuilder();
// 1. 元数据头(始终包含)
sb.append("文档: ").append(entry.getTitle()).append("\n");
if (entry.getSummary() != null) {
sb.append("摘要: ").append(entry.getSummary()).append("\n");
}
// 2. 章节结构(## 标题列表)
String headings = body.lines()
.filter(l -> l.trim().startsWith("##"))
.map(l -> " - " + l.trim().replaceAll("^#+\\s*", ""))
.collect(Collectors.joining("\n"));
if (!headings.isEmpty()) {
sb.append("章节:\n").append(headings).append("\n");
}
sb.append("---\n");
// 3. 正文片段(去标题行、去空行,智能截断)
String textContent = body.lines()
.filter(l -> !l.trim().startsWith("#") && !l.trim().isEmpty())
.collect(Collectors.joining("\n"))
.trim();
// 短文档保留更多内容,长文档节省上下文
int maxBodyChars = body.length() < 500 ? 800 : 500;
if (textContent.length() > maxBodyChars) {
sb.append(textContent, 0, maxBodyChars).append("...");
} else {
sb.append(textContent);
}
return sb.toString();
}
/**
* 构建纯元数据摘要(不读文件,仅用内存索引信息)
* 多匹配且有 L1 补充时使用,L0 只需告知 LLM 命中了哪些文档
*/
private String buildMetadataOnlySummary(KnowledgeEntry entry) {
StringBuilder sb = new StringBuilder();
sb.append("文档: ").append(entry.getTitle()).append("\n");
if (entry.getSummary() != null) {
sb.append("摘要: ").append(entry.getSummary()).append("\n");
}
if (entry.getKeywords() != null && !entry.getKeywords().isEmpty()) {
sb.append("关键词: ").append(String.join(", ", entry.getKeywords())).append("\n");
}
sb.append("来源: ").append(entry.getFilePath()).append("\n");
return sb.toString();
}
/**
* 提取 MD 内容中第一个有意义的文本行(跳过 frontmatter 和标题行)
*/
private String extractFirstMeaningfulLine(String content, int maxLen) {
if (content == null || content.isBlank()) return "(空)";
String text = content.trim();
// 跳过 YAML frontmatter (--- ... ---)
if (text.startsWith("---")) {
int end = text.indexOf("---", 3);
if (end != -1) {
text = text.substring(end + 3);
}
}
// 查找第一个非空、非标题行
String[] lines = text.split("\n");
for (String line : lines) {
String tl = line.trim();
if (!tl.isEmpty() && !tl.startsWith("#")) {
return tl.length() <= maxLen ? tl : tl.substring(0, maxLen) + "...";
}
}
// 兜底:第一行非空行
for (String line : lines) {
if (!line.trim().isEmpty()) {
String tl = line.trim();
return tl.length() <= maxLen ? tl : tl.substring(0, maxLen) + "...";
}
}
return "(无有效内容)";
}
}
@@ -0,0 +1,44 @@
package com.superbiz.agent.util;
import java.util.List;
/**
* 问题复杂度判断
* 用于决定使用单 Agent 还是多 Agent(Planner + Executor)处理
*/
public class QuestionComplexity {
/** 复杂问题关键词 — 需要多步分析、排查、根因定位 */
private static final List<String> COMPLEX_KEYWORDS = List.of(
"排查", "分析", "为什么", "根因", "调查", "对比", "影响范围",
"原因", "故障", "告警", "诊断", "链路", "流程", "步骤",
"root cause", "troubleshoot", "investigate"
);
/** 极简问题关键词 — 快速回答,无需多 Agent */
private static final List<String> SIMPLE_KEYWORDS = List.of(
"是什么", "查一下", "什么是", "时间", "天气", "定义",
"查", "找", "what is", "define", "time"
);
/**
* 判断是否为复杂问题
*/
public static boolean isComplex(String question) {
if (question == null || question.isBlank()) return false;
String q = question.toLowerCase();
// 复杂关键词匹配 → 多 Agent
for (String kw : COMPLEX_KEYWORDS) {
if (q.contains(kw)) return true;
}
// 简单关键词匹配 → 单 Agent
for (String kw : SIMPLE_KEYWORDS) {
if (q.contains(kw)) return false;
}
// 默认:长问题(>30 字)视为复杂,短问题视为简单
return question.length() > 30;
}
}
@@ -0,0 +1,29 @@
package com.superbiz.agent.util;
/**
* 会话上下文持有者(基于 ThreadLocal)
* <p>
* 用于在同步调用链路中传递 sessionId,兜底 LookupKnowledgeTool 等
* 无法通过 RunnableConfig 获取上下文的组件。
* 优先使用 RunnableConfig.metadata 传递,ThreadLocal 作为同步路径的补充。
* <p>
* 使用规范:
* 1. 调用方在 Agent 执行前调用 setSessionId()
* 2. finally 块中调用 clear()
*/
public class SessionContextHolder {
private static final ThreadLocal<String> SESSION_ID = new ThreadLocal<>();
public static void setSessionId(String sessionId) {
SESSION_ID.set(sessionId);
}
public static String getSessionId() {
return SESSION_ID.get();
}
public static void clear() {
SESSION_ID.remove();
}
}
@@ -0,0 +1,75 @@
-- V005: 创建会话存储体系(diagnosis_session + agent_step + tool_invocation)
-- 设计文档:openspec/changes/session-storage/design.md
CREATE TABLE diagnosis_session (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
session_id VARCHAR(64) UNIQUE NOT NULL COMMENT '会话唯一 ID',
query TEXT NOT NULL COMMENT '用户原始问题',
status VARCHAR(16) DEFAULT 'PENDING' COMMENT 'PENDING/RUNNING/SUCCESS/FAILED',
agent_flow VARCHAR(32) COMMENT 'CHAT / AI_OPS',
total_duration_ms INT COMMENT '总耗时(毫秒)',
total_token_count INT COMMENT '总 Token 消耗',
step_count INT COMMENT 'Agent 步数',
tool_call_count INT COMMENT '工具调用次数',
self_evaluation JSON COMMENT '自评估信号:{"confidence":0-100,"reasoning":"..."}',
feedback VARCHAR(16) COMMENT '用户反馈:useful/not_useful/null',
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_created_at (created_at),
INDEX idx_status (status),
INDEX idx_agent_flow (agent_flow)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='诊断会话表';
CREATE TABLE agent_step (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
session_id VARCHAR(64) NOT NULL COMMENT '关联 diagnosis_session',
step_index INT NOT NULL COMMENT '当前 Agent 的第几步(从0开始)',
agent_name VARCHAR(32) NOT NULL COMMENT 'intelligent_assistant/planner/executor',
model_input JSON COMMENT '模型输入摘要',
model_output JSON COMMENT '模型输出摘要(含工具调用决策)',
thought TEXT COMMENT 'Agent 思考过程',
has_tool_call BOOLEAN DEFAULT FALSE COMMENT '本轮是否调用了工具',
duration_ms INT COMMENT '本轮耗时',
token_count INT COMMENT '本轮 Token 消耗',
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_session_step (session_id, step_index),
INDEX idx_agent_name (agent_name)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='Agent 决策步骤表';
CREATE TABLE tool_invocation (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
session_id VARCHAR(64) NOT NULL COMMENT '关联 diagnosis_session',
step_id BIGINT COMMENT '关联 agent_step.id(可为空,不强制外键)',
tool_name VARCHAR(64) NOT NULL COMMENT 'lookup_knowledge/queryPrometheusAlerts/等',
input_params JSON NOT NULL COMMENT '工具入参',
output_preview TEXT COMMENT '输出前500字符',
output_length INT COMMENT '输出总字符数',
retrieval_layer VARCHAR(8) COMMENT 'L0/L1/L0+L1',
l0_match_count INT COMMENT 'L0 匹配数',
l1_match_count INT COMMENT 'L1 匹配数',
is_truncated BOOLEAN DEFAULT FALSE COMMENT '内容是否被截断',
retrieval_details JSON COMMENT '检索明细:{l0_titles:[], l1_scores:[]}',
duration_ms INT COMMENT '工具执行耗时',
success BOOLEAN DEFAULT TRUE COMMENT '是否成功',
error_message TEXT COMMENT '失败原因',
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_session_id (session_id),
INDEX idx_tool_name (tool_name),
INDEX idx_retrieval_layer (retrieval_layer)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工具调用明细表';
@@ -0,0 +1,6 @@
-- V006: 将 agent_step 的 model_input / model_output 从 JSON 改为 TEXT
-- 原因:buildModelInputSummary() 输出的是纯文本摘要,不是合法 JSON
ALTER TABLE agent_step
MODIFY COLUMN model_input TEXT COMMENT '模型输入摘要',
MODIFY COLUMN model_output TEXT COMMENT '模型输出摘要(含工具调用决策)';
@@ -0,0 +1,2 @@
-- V007: 删除旧的 diagnosis_record 表(已被 diagnosis_session + agent_step + tool_invocation 替代)
DROP TABLE IF EXISTS diagnosis_record;
@@ -0,0 +1,12 @@
你是任务执行器。执行 Planner 分配给你的具体步骤,并及时反馈结果。
## 职责
- 按步骤执行具体的查询任务
- 使用知识库查询、日志查询等工具获取信息
- 将执行结果汇总,给出完整的最终答案
## 规则
- 按顺序执行,不可跳过步骤
- 所有需要外部信息的地方,都必须调用对应的工具
- 不要凭记忆回答,必须基于工具返回的真实数据
- 执行完成后,综合所有结果给出完整的答案
@@ -0,0 +1,20 @@
你是智能任务规划器。分析用户的问题,拆解为具体的执行步骤。
## 职责
- 分析用户问题,拆解为可执行的步骤列表
- **你不能调用任何工具**,你的职责是制定计划,不是执行
- 输出 JSON 格式的计划,不输出其他内容
## 输出格式
```json
{
"plan": ["步骤1描述", "步骤2描述", "步骤3描述"],
"reasoning": "规划思路说明"
}
```
## 规则
- 每个步骤应该是一个可以独立执行的任务
- 步骤要具体可操作,不要模糊
- 如果问题需要查知识库,明确在步骤中说明要查什么
@@ -0,0 +1,76 @@
# 执行者 System Prompt
## 角色定位
你是诊断流程的**执行者**。你的任务非常明确:严格遵循规划者下发的任务清单,按步骤调用工具完成任务,并输出最终结果。
---
## 核心行为准则
### 1. 严格按步执行
- 规划者下发的是**有序的任务列表**(如 Step 1 → Step 2 → Step 3)
- 你必须按顺序执行,不可跳过、合并或重排步骤
- 每个步骤完成后,记录该步骤的产出,再进入下一步
### 2. 调用工具而不是凭记忆回答
- 所有需要外部信息的地方,都必须调用对应的工具
- 尤其注意:永远不要凭记忆回答错误码含义、接口定义、排障步骤
- 知识库查询:必须通过 `lookup_knowledge` 工具完成
### 3. 工具调用完毕后,必须结合日志、订单数据等证据综合分析
- 不要把工具的返回结果直接当作最终答案输出
- 你的结论必须基于**至少两个独立证据源**(如错误码+日志、接口文档+实际返回值)
---
## 可用工具
### lookup_knowledge(知识库查询)
用于查询内部知识库,获取错误码定义、接口文档、排障步骤等背景信息。
| 参数 | 说明 |
|------|------|
| `query` | 查询关键词或描述。例如:`ERR_TIMEOUT`、`payment-gateway`、`支付为什么失败` |
**内部机制**:
工具内部自动执行「先精确匹配(L0),未命中则语义检索(L1)」的两阶段检索逻辑,你无需关心哪一层。
**返回结果**:包含 `found`(是否找到)、`primary.content`(文档内容)、`primary.match_type`(来源标记:`exact_L0` 或 `semantic_L1`)等字段。
**使用规则**:
- 当你查到了错误码、接口名、服务名时:**必须**调用此工具
- 当需要查排障步骤、业务流程、最佳实践时:**必须**调用此工具
- 对当前结果没有十足把握时:**建议**调用此工具验证
---
## 任务执行规范
### 1. 每个步骤的产出要求
每完成一个工具调用后,你应该:
- 记录工具返回的关键信息
- 将新信息与已有上下文(日志、订单数据等)进行交叉验证
- 输出该步骤的阶段性结论
### 2. 最终输出的报告格式
```yaml
## 诊断结论
**问题根因**:XXX
**证据链**:
1. 订单状态返回错误码 ERR_TIMEOUT
2. 知识库 lookup_knowledge("ERR_TIMEOUT") 返回:支付网关响应超时(>5秒)
3. 日志确认:14:32:15 请求耗时 5.3s,超过 5s 阈值
**建议方案**:
- 临时方案:重试该笔订单
- 长期方案:优化支付网关超时配置,建议提升至 8s
**引用来源**:
- [来源: interfaces/_errors.md]
```
@@ -0,0 +1,88 @@
你是 Planner Agent,同时承担 Replanner 角色,负责:
1. 读取当前输入任务 {input} 以及 Executor 的最近反馈 {executor_feedback}。
2. 分析 Prometheus 告警、日志、内部文档等信息,制定可执行的下一步步骤。
3. 在执行阶段,输出 JSON,包含 decision (PLAN|EXECUTE|FINISH)、step 描述、预期要调用的工具、以及必要的上下文。
4. 调用任何腾讯云日志/主题相关工具时,region 参数必须使用连字符格式(如 ap-guangzhou),若不确定请省略以使用默认值。
5. 严格禁止编造数据,只能引用工具返回的真实内容;如果连续 3 次调用同一工具仍失败或返回空结果,需停止该方向并在最终报告的结论部分说明"无法完成"的原因。
## 最终报告输出要求(CRITICAL)
当 decision=FINISH 时,你必须:
1. **不要输出 JSON 格式**
2. **直接输出完整的 Markdown 格式报告文本**
3. **报告必须严格遵循以下模板**:
```
# 告警分析报告
---
## 📋 活跃告警清单
| 告警名称 | 级别 | 目标服务 | 首次触发时间 | 最新触发时间 | 状态 |
|---------|------|----------|-------------|-------------|------|
| [告警1名称] | [级别] | [服务名] | [时间] | [时间] | 活跃 |
| [告警2名称] | [级别] | [服务名] | [时间] | [时间] | 活跃 |
---
## 🔍 告警根因分析1 - [告警名称]
### 告警详情
- **告警级别**: [级别]
- **受影响服务**: [服务名]
- **持续时间**: [X分钟]
### 症状描述
[根据监控指标描述症状]
### 日志证据
[引用查询到的关键日志]
### 根因结论
[基于证据得出的根本原因]
---
## 🛠️ 处理方案执行1 - [告警名称]
### 已执行的排查步骤
1. [步骤1]
2. [步骤2]
### 处理建议
[给出具体的处理建议]
### 预期效果
[说明预期的效果]
---
## 🔍 告警根因分析2 - [告警名称]
[如果有第2个告警,重复上述格式]
---
## 📊 结论
### 整体评估
[总结所有告警的整体情况]
### 关键发现
- [发现1]
- [发现2]
### 后续建议
1. [建议1]
2. [建议2]
### 风险评估
[评估当前风险等级和影响范围]
```
**重要提醒**:
- 最终输出必须是纯 Markdown 文本,不要包含 JSON 结构
- 不要使用 "finalReport": "..." 这样的格式
- 直接从 "# 告警分析报告" 开始输出
- 所有内容必须基于工具查询的真实数据,严禁编造
- 如果某个步骤失败,在结论中如实说明,不要跳过
@@ -0,0 +1,10 @@
你是 AI Ops Supervisor,负责调度 planner_agent 与 executor_agent:
1. 当需要拆解任务或重新制定策略时,调用 planner_agent。
2. 当 planner_agent 输出 decision=EXECUTE 时,调用 executor_agent 执行第一步。
3. 根据 executor_agent 的反馈,评估是否需要再次调用 planner_agent,直到 decision=FINISH。
4. FINISH 后,确保向最终用户输出完整的《告警分析报告》,格式必须严格为:
告警分析报告\n---\n# 告警处理详情\n## 活跃告警清单\n## 告警根因分析N\n## 处理方案执行N\n## 结论。
5. 若步骤涉及腾讯云日志/主题工具,请确保使用连字符区域 ID(ap-guangzhou 等),或省略 region 以采用默认值。
6. 如果发现 Planner/Executor 在同一方向连续 3 次调用工具仍失败或没有数据,必须终止流程,直接输出"任务无法完成"的报告,明确告知失败原因,严禁凭空编造结果。
只允许在 planner_agent、executor_agent 与 FINISH 之间做出选择。