Add diagnosis playbook skills

This commit is contained in:
aruo
2026-07-06 08:35:54 +08:00
parent 88e0a6c944
commit 6ccfd33ec5
23 changed files with 1002 additions and 75 deletions
@@ -0,0 +1,83 @@
package com.superbiz.agent.config;
import com.alibaba.cloud.ai.graph.skills.SkillMetadata;
import com.alibaba.cloud.ai.graph.skills.registry.SkillRegistry;
import com.alibaba.cloud.ai.graph.skills.registry.classpath.ClasspathSkillRegistry;
import org.springframework.ai.chat.prompt.SystemPromptTemplate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.io.IOException;
import java.util.List;
import java.util.Optional;
@Configuration
public class SkillConfig {
private static final String ACTIVE_SKILL_NAME = "diagnose-mysql-connection-pool";
@Bean
public SkillRegistry skillRegistry() {
SkillRegistry classpathRegistry = ClasspathSkillRegistry.builder()
.classpathPath("skills")
.basePath("target/skills-cache")
.build();
return new SingleSkillRegistry(classpathRegistry, ACTIVE_SKILL_NAME);
}
private record SingleSkillRegistry(SkillRegistry delegate, String activeSkillName) implements SkillRegistry {
@Override
public List<SkillMetadata> listAll() {
return delegate.listAll().stream()
.filter(skill -> activeSkillName.equals(skill.getName()))
.toList();
}
@Override
public String getRegistryType() {
return delegate.getRegistryType();
}
@Override
public String readSkillContent(String skillName) throws IOException {
if (!activeSkillName.equals(skillName)) {
throw new IOException("Skill not found: " + skillName);
}
return delegate.readSkillContent(skillName);
}
@Override
public String getSkillLoadInstructions() {
return delegate.getSkillLoadInstructions();
}
@Override
public SystemPromptTemplate getSystemPromptTemplate() {
return delegate.getSystemPromptTemplate();
}
@Override
public Optional<SkillMetadata> get(String skillName) {
if (!activeSkillName.equals(skillName)) {
return Optional.empty();
}
return delegate.get(skillName);
}
@Override
public int size() {
return listAll().size();
}
@Override
public boolean contains(String skillName) {
return activeSkillName.equals(skillName) && delegate.contains(skillName);
}
@Override
public void reload() {
delegate.reload();
}
}
}
@@ -0,0 +1,101 @@
package com.superbiz.agent.hook;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import com.alibaba.cloud.ai.graph.agent.Prioritized;
import com.alibaba.cloud.ai.graph.agent.hook.HookPosition;
import com.alibaba.cloud.ai.graph.agent.hook.HookPositions;
import com.alibaba.cloud.ai.graph.agent.hook.messages.AgentCommand;
import com.alibaba.cloud.ai.graph.agent.hook.messages.MessagesModelHook;
import com.alibaba.cloud.ai.graph.skills.SkillMetadata;
import com.alibaba.cloud.ai.graph.skills.registry.SkillRegistry;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.messages.Message;
import org.springframework.ai.chat.messages.SystemMessage;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
/**
* Injects planner-visible skill metadata without exposing the full skill loader tool.
*/
@Slf4j
@HookPositions(HookPosition.BEFORE_MODEL)
public class PlannerSkillMetadataHook extends MessagesModelHook {
private static final String CATALOG_MARKER = "\"skill_catalog\"";
private final SkillRegistry skillRegistry;
private final ObjectMapper objectMapper = new ObjectMapper();
public PlannerSkillMetadataHook(SkillRegistry skillRegistry) {
this.skillRegistry = skillRegistry;
}
@Override
public String getName() {
return "planner_skill_metadata_hook";
}
@Override
public int getOrder() {
return Prioritized.HIGHEST_PRECEDENCE;
}
@Override
public AgentCommand beforeModel(List<Message> previousMessages, RunnableConfig config) {
if (skillRegistry == null || skillRegistry.size() == 0 || hasCatalog(previousMessages)) {
return new AgentCommand(previousMessages);
}
try {
List<Map<String, String>> skills = skillRegistry.listAll().stream()
.map(this::toSkillSummary)
.toList();
if (skills.isEmpty()) {
return new AgentCommand(previousMessages);
}
Map<String, Object> catalog = new LinkedHashMap<>();
catalog.put("purpose", "Planner-visible diagnosis skill metadata only.");
catalog.put("rules", List.of(
"Choose at most one primary skill.",
"Do not load full skill instructions in Planner.",
"Executor reads the selected skill before evidence collection.",
"If no skill matches, set selected_skill to null."
));
catalog.put("skills", skills);
catalog.put("required_planner_output", Map.of(
"selected_skill", "skill name or null",
"selection_reason", "short reason",
"plan", "ordered execution step list"
));
Map<String, Object> payload = Map.of("skill_catalog", catalog);
String content = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(payload);
List<Message> updatedMessages = new ArrayList<>(previousMessages.size() + 1);
updatedMessages.add(new SystemMessage(content));
updatedMessages.addAll(previousMessages);
return new AgentCommand(updatedMessages);
} catch (Exception e) {
log.warn("Failed to inject planner skill metadata, fallback to original messages", e);
return new AgentCommand(previousMessages);
}
}
private Map<String, String> toSkillSummary(SkillMetadata skill) {
Map<String, String> summary = new LinkedHashMap<>();
summary.put("name", skill.getName());
summary.put("description", skill.getDescription());
return summary;
}
private boolean hasCatalog(List<Message> messages) {
return messages.stream()
.map(Message::getText)
.anyMatch(text -> text != null && text.contains(CATALOG_MARKER));
}
}
@@ -4,7 +4,10 @@ import org.springframework.ai.chat.model.ChatModel;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.agent.ReactAgent;
import com.alibaba.cloud.ai.graph.agent.flow.agent.SupervisorAgent;
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.superbiz.agent.agent.tool.DateTimeTools;
import com.superbiz.agent.agent.tool.InternalDocsTools;
import com.superbiz.agent.agent.tool.QueryLogsTools;
@@ -13,6 +16,7 @@ import com.superbiz.agent.domain.entity.AgentStep;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import com.superbiz.agent.dto.AIOpsRequest;
import com.superbiz.agent.hook.AgentLoggingHook;
import com.superbiz.agent.hook.PlannerSkillMetadataHook;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.repository.ToolInvocationRepository;
@@ -32,8 +36,8 @@ import java.util.Optional;
import java.util.UUID;
/**
* AI Ops 智能运维服务
* 负责多 Agent 协作的告警分析流程
* AI Ops 闂備礁鎼幊妯肩磽濮樿泛绀傛俊顖滅帛娴溿倝鏌熼柇锕€鏋熸俊顖氾躬閺岋繝宕煎┑鍩裤垹鈹?
* 闂佽崵濮甸崝妤呭窗閺囥垺鍎楁俊銈勭缁?Agent 闂備礁鎲¢〃鍛崲鐎n剛绀婇柡鍐ㄧ墛閸庡秹鏌涢弴銊ヤ簼闁哥喓鍋ら幃褰掑焵椤掑嫭鏅濋柍褜鍓熷畷瑙勬償閵娿儱鍤戦梺褰掑亰閸橀箖濡堕敂鍓х<?
*/
@Service
public class AiOpsService {
@@ -49,7 +53,7 @@ public class AiOpsService {
@Autowired
private QueryMetricsTools queryMetricsTools;
@Autowired(required = false) // Mock 模式下才注册
@Autowired(required = false) // Mock 婵犵妲呴崹顏堝焵椤掆偓绾绢厾娑甸埀顒佺箾閹寸偞灏い鎴濇閺呭爼鎮╁ù瀣亙闂侀潧顭堥崕閬嶅棘閳?
private QueryLogsTools queryLogsTools;
@Autowired
@@ -73,13 +77,16 @@ public class AiOpsService {
@Autowired
private SelfEvaluationMergeService selfEvaluationMergeService;
@Autowired(required = false)
private SkillRegistry skillRegistry;
/**
* 执行 AI Ops 告警分析流程
* 闂備礁婀遍悷鎶藉幢閳哄倹鏉?AI Ops 闂備礁鎲$粙鎴︽晝閵娾晩鏁嗛柣鏃傚帶缁€鍡涙煕閳╁喚鐒介柍褜鍓濆Λ鍕箒婵炶揪缍€閵嗏偓闁?
*
* @param chatModel 大模型实例
* @param toolCallbacks 工具回调数组
* @return 分析结果状态
* @throws GraphRunnerException 如果 Agent 执行失败
* @param chatModel 濠电姰鍨归悥銏ゅ礋閳ь剚绗熼埀顒€鐣烽崷顓涘亾閿濆簼绨介柡澶庢閵?
* @param toolCallbacks 闁诲氦顫夐幃鍫曞磿闁秴鐭楅柛褎顨呴悙濠囨煟閹邦剙顣虫繛鍫濈埣閺屸剝鎷呴悷閭︽缂?
* @return 闂備礁鎲$敮鎺懳涘▎鎾村€甸柣锝呯灱绾惧ジ鏌熼幆褜鍤熷ù婊庡灦閺岋絽螣閸喚鍘梺?
* @throws GraphRunnerException 濠电姷顣介埀顒€鍟块埀顒€缍婇幃?Agent 闂備礁婀遍悷鎶藉幢閳哄倹鏉稿┑鐘灪閸庤偐鍒掗崜褎鍠?
*/
public Optional<OverAllState> executeAiOpsAnalysis(ChatModel chatModel, ToolCallback[] toolCallbacks) throws GraphRunnerException {
return executeAiOpsAnalysis(chatModel, toolCallbacks, null, resolveSessionId(null));
@@ -87,27 +94,27 @@ public class AiOpsService {
public Optional<OverAllState> executeAiOpsAnalysis(ChatModel chatModel, ToolCallback[] toolCallbacks,
AIOpsRequest request, String sessionId) throws GraphRunnerException {
logger.info("开始执行 AI Ops 多 Agent 协作流程");
logger.info("Starting AI Ops multi-agent analysis");
String resolvedSessionId = isBlank(sessionId) ? resolveSessionId(request) : sessionId.trim();
long startTime = System.currentTimeMillis();
// 创建或更新诊断会话
// 闂備礁鎲$敮妤冪矙閹寸姷纾介柟鎹愵嚙缁狅綁鏌″鍐ㄥ缂佺虎鍨堕弻锟犲磼濞戞﹩鈧粓鏌i敂鐣屽⒌鐎殿噮鍓熼、妯衡攽閸垻宕堕梺?
DiagnosisSession session = startDiagnosisSession(resolvedSessionId, request);
diagnosisSessionRepository.save(session);
// 设置 ThreadLocal 上下文(LookupKnowledgeTool 通过此获取 sessionId)
// 闂佽崵濮崇粈浣规櫠娴犲鍋?ThreadLocal 濠电偞鍨堕幐鎼佹晝閿濆洨绠旈柛娑欐綑濡﹢鏌涢妷鈺婃缂佲偓閸戯箰okupKnowledgeTool 闂傚倷绶¢崑鍛┍閾忚宕查柛鎰电厛濞间即鏌曢崼婵堝缂佺媭鍨堕弻?sessionId闂?
SessionContextHolder.setSessionId(resolvedSessionId);
try {
// 构建 Planner 和 Executor Agent(每个 Agent 各自带 Hook)
// 闂備礁鎼鍛偓姘煎墰缁?Planner 闂?Executor Agent闂備焦瀵х粙鎴︽偋婵犲洦鍎婇柍鈺佸暟閳?Agent 闂備礁鎲¢懝鍓р偓姘煎灦瀹曢潧顭ㄩ崨顔芥?Hook闂?
ReactAgent plannerAgent = buildPlannerAgent(chatModel, toolCallbacks);
ReactAgent executorAgent = buildExecutorAgent(chatModel, toolCallbacks);
// 构建 Supervisor Agent(不加 Hook)
// 闂備礁鎼鍛偓姘煎墰缁?Supervisor Agent闂備焦瀵х粙鎴︽偋閸涱垳绠斿鑸靛姇缁€?Hook闂?
SupervisorAgent supervisorAgent = SupervisorAgent.builder()
.name("ai_ops_supervisor")
.description("负责调度 Planner 与 Executor 的多 Agent 控制器")
.description("Coordinates Planner and Executor agents")
.model(chatModel)
.systemPrompt(promptProperties.getSupervisor())
.subAgents(List.of(plannerAgent, executorAgent))
@@ -115,19 +122,19 @@ public class AiOpsService {
String taskPrompt = buildTaskPrompt(request);
logger.info("调用 Supervisor Agent 开始编排...");
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());
@@ -146,25 +153,25 @@ public class AiOpsService {
}
/**
* 从执行结果中提取最终报告文本
* 濠电偛顕慨瀵糕偓娑掓櫆閺呭爼鎮剧仦鎯т粧閻庡厜鍋撻柍褜鍓涢崚鎺楀Ω閳轰礁鍤戝┑鐘才堥崑鎾绘煠閸偄鐏存鐐存崌楠炲洭顢楅埀顒傚緤閸ф鐓涢柛顐h壘娴滃墽绱撻崒娆戭槮闁绘锕ラ幈銊╁Χ婢跺﹤绐涙繝鐢靛Т閸燁垶鎮楅鈧弻?
*
* @param state 执行状态
* @return 报告文本(如果存在)
* @param state 闂備礁婀遍悷鎶藉幢閳哄倹鏉搁梻浣虹帛椤牓宕洪弽顓炵劦?
* @return 闂備胶顢婄紙浼村磿闁秴绠熼柨鐔哄Т濡﹢鏌涢妷锝呭闁圭兘浜堕弻銊モ槈濡厧顣哄銈傛暘閸パ冨殤濠电姴锕ら崯浼村箺閻樼粯鐓曢柨鏃囧吹閸樻粎绱?
*/
public Optional<String> extractFinalReport(OverAllState state) {
logger.info("开始提取最终报告...");
logger.info("闁诲孩顔栭崰鎺楀磻閹炬枼鏀芥い鏃傗拡閸庢劗鎲告0浣虹獢鐎规洩缍佸浠嬪Ω閿旇法甯涚紓鍌氬€风粈渚€鎮ф繝鍐╁弿闁靛牆顦?..");
// 提取 Planner 最终输出(包含完整的告警分析报告)
// 闂備礁婀辩划顖炲礉閺嚶颁汗?Planner 闂備礁鎼悧鍐磻閹惧墎纾藉ù锝呮憸婢э絿绱掓0婵嗗籍鐎规洘鐟╅幃顔锯偓闈涙憸椤︹晠姊洪崨濠勫ⅹ闁瑰啿閰i獮鍡涘醇閳垛晛浜鹃柣鐔哄濠€浼存煛閸☆厾绉柟顖氬暣瀹曠喖顢曢敐鍛畼闂佽崵濮崑鎾绘煥閺囨浜鹃梺鎼炲妼闁帮絽顕i幖浣哥疀妞ゆ挾鍊幘缁樼厱婵炴垶锕╅悡顓犵磼?
Optional<AssistantMessage> plannerFinalOutput = state.value("planner_plan")
.filter(AssistantMessage.class::isInstance)
.map(AssistantMessage.class::cast);
if (plannerFinalOutput.isPresent()) {
String reportText = plannerFinalOutput.get().getText();
logger.info("成功提取到 Planner 最终报告,长度: {}", reportText.length());
logger.info("闂備胶鎳撻悺銊╁礉閺囩喐鍙忔繛鎴欏灩缁犵敻鏌熼柇锕€澧紒鎻掓健閺?Planner 闂備礁鎼悧鍐磻閹惧墎纾藉ù锝呮憸婢ф稑鈹戦鍝勨偓婵嗙暦閵婏妇绡€闊洦娲滈ˇ顕€姊婚崒妤€浜鹃梺鍓茬厛閸犳牠顢? {}", reportText.length());
return Optional.of(reportText);
} else {
logger.warn("未能提取到 Planner 最终报告");
logger.warn("Unable to extract Planner final report");
return Optional.empty();
}
}
@@ -197,16 +204,16 @@ public class AiOpsService {
String buildQuerySummary(AIOpsRequest request) {
if (request == null) {
return "AI Ops 告警分析";
return "AI Ops alert analysis";
}
StringBuilder summary = new StringBuilder("AI Ops 告警分析");
appendField(summary, "告警", request.getAlertName());
appendField(summary, "服务", request.getService());
appendField(summary, "等级", request.getSeverity());
appendField(summary, "时间范围", request.getTimeRange());
appendField(summary, "描述", request.getDescription());
appendField(summary, "请求", request.getUserRequest());
StringBuilder summary = new StringBuilder("AI Ops alert analysis");
appendField(summary, "alert", request.getAlertName());
appendField(summary, "service", request.getService());
appendField(summary, "severity", request.getSeverity());
appendField(summary, "timeRange", request.getTimeRange());
appendField(summary, "description", request.getDescription());
appendField(summary, "request", request.getUserRequest());
return summary.toString();
}
@@ -238,8 +245,8 @@ public class AiOpsService {
String buildTaskPrompt(AIOpsRequest request) {
StringBuilder prompt = new StringBuilder();
prompt.append("你是企业级 SRE,接到了自动化告警排查任务。请结合工具调用,执行**规划→执行→再规划**的闭环,并最终按照固定模板输出《告警分析报告》。禁止编造虚假数据,如连续多次查询失败需诚实反馈无法完成的原因。");
prompt.append("\n\n本次告警输入:\n");
prompt.append("You are an enterprise SRE handling an automated alert diagnosis task. Combine tool evidence, run a plan-execute-replan loop, and output the final alert analysis report. Do not fabricate data; if repeated queries fail, clearly state why the task cannot be completed.");
prompt.append("\n\nAlert input:\n");
prompt.append(buildQuerySummary(request));
if (hasAlertPayload(request)) {
String knowledgeQuery = buildKnowledgeRetrievalQuery(request);
@@ -278,53 +285,69 @@ public class AiOpsService {
}
/**
* 构建 Planner Agent
* 闂備礁鎼鍛偓姘煎墰缁?Planner Agent
*/
private ReactAgent buildPlannerAgent(ChatModel chatModel, ToolCallback[] toolCallbacks) {
return ReactAgent.builder()
.name("planner_agent")
.description("负责拆解告警、规划与再规划步骤")
.description("Plans alert diagnosis steps")
.model(chatModel)
.systemPrompt(promptProperties.getPlanner())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "planner"))
.hooks(buildHooks("planner"))
.outputKey("planner_plan")
.build();
}
/**
* 构建 Executor Agent
* 闂備礁鎼鍛偓姘煎墰缁?Executor Agent
*/
private ReactAgent buildExecutorAgent(ChatModel chatModel, ToolCallback[] toolCallbacks) {
return ReactAgent.builder()
.name("executor_agent")
.description("负责执行 Planner 的首个步骤并及时反馈")
.description("Executes the current Planner step and reports feedback")
.model(chatModel)
.systemPrompt(promptProperties.getExecutor())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "executor"))
.hooks(buildHooks("executor"))
.outputKey("executor_feedback")
.build();
}
/**
* 动态构建方法工具数组
* 根据 cls.mock-enabled 决定是否包含 QueryLogsTools
* 工具顺序:知识库查询优先,日志查询次之,弃用工具最后
* 闂備礁鎲¢弻锝夊礉瀹ュ鐒垫い鎴f硶閸斿秹鏌f惔顔肩仩妞ゆ洘鐟╅幃婊兾熼懡銈呭箥婵犵數鍋涢ˇ鏉棵洪弽銊ヮ嚤闁圭増婢樼粈鍌炴⒑閸噮鍎愭繛鍫濆缁?
* 闂備礁鎼粔鐑斤綖婢跺﹦鏆?cls.mock-enabled 闂備礁鎲¢崝鏇㈠疮閸ф鍋╁Δ锝呭暙閸欏﹥銇勯弽銊ь暡闁稿骸锕弻娑㈠冀瑜庨崳褰掓煙?QueryLogsTools
* 闁诲氦顫夐幃鍫曞磿闁秴鐭楅柟绋跨昂娴滄粓鏌涘┑鍡楊伀缁炬澘绉归弻銊モ槈濞嗘劗娈ら梺缁樻惈缁辨洟骞忛悩璇插耿婵°倕鍟惃鎴︽⒑閸濆嫯顫﹂柛搴㈡尦椤㈡艾螖娴e壊鍤ゅ┑鈽嗗灠閹碱偆鏁妷鈺傜叆婵炴垶蓱濠€鐗堜繆椤愮喐娅堢紒鐘崇☉铻栧ù锝呮惈瀵劑鏌i悩鍙夊偍闁搞劍妞介、鏇㈠礂閼测斁鏋欓柣搴到婢у海绮堟径灞稿亾濞堝灝鏋涢柛鐔跺嵆瀵偊濡堕崪浣告櫊闂侀潧顦崕鍗烆嚗閺冨牊鐓涢柛顐h壘娴滈箖姊?
*/
private Object[] buildMethodToolsArray() {
if (queryLogsTools != null) {
// Mock 模式:包含 QueryLogsTools
// Mock 婵犵妲呴崹顏堝焵椤掆偓绾绢厾娑甸埀顒勬⒑閹稿海鈯曢柤鐟板⒔閳ь剙鐏氶敃銏犵暦?QueryLogsTools
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools, queryLogsTools};
} else {
// 真实模式:不包含 QueryLogsTools(由 MCP 提供日志查询功能)
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools};
}
// Real mode excludes local QueryLogsTools because logs are provided by MCP.
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools};
}
/** 从 agent_step 和 tool_invocation 汇总指标回填 diagnosis_session */
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()
};
}
/** 濠?agent_step 闂?tool_invocation 婵犳鍠氶幊鎾趁洪敃鍌氱劦妞ゆ帒鍊荤敮娑㈡倵閸偄鍝虹€殿喕绮欏畷鎯邦槼缂佲偓閳ь剚绻?diagnosis_session */
private void backfillSessionMetrics(DiagnosisSession session) {
try {
List<AgentStep> steps = agentStepRepository.findBySessionIdOrderByStepIndex(session.getSessionId());
@@ -340,7 +363,7 @@ public class AiOpsService {
session.setStepCount(stepCount);
session.setToolCallCount(Math.toIntExact(toolCallCount));
} catch (Exception e) {
logger.warn("回填会话指标失败: sessionId={}", session.getSessionId(), e);
logger.warn("闂備焦鎮堕崕鎶藉磻濞戙垺鏅查柣鎰綑椤曡鲸鎱ㄥΟ铏癸紞婵☆垰鐗撻弻鐔虹矙閹稿骸顦╅梺缁樼壄缁叉儳顕ラ崟顒佺秶妞ゆ劑鍎? sessionId={}", session.getSessionId(), e);
}
}
@@ -4,7 +4,10 @@ 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;
@@ -13,6 +16,7 @@ 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;
@@ -99,6 +103,9 @@ public class ChatService {
@Autowired
private KnowledgeDomainService knowledgeDomainService;
@Autowired(required = false)
private SkillRegistry skillRegistry;
@Autowired
private ToolTraceSummaryService toolTraceSummaryService;
@@ -160,6 +167,7 @@ public class ChatService {
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");
@@ -271,7 +279,7 @@ public class ChatService {
.systemPrompt(systemPrompt)
.methodTools(buildMethodToolsArray())
.tools(getToolCallbacks())
.hooks(new AgentLoggingHook(agentStepRepository, "intelligent_assistant"))
.hooks(buildHooks("intelligent_assistant"))
.build();
}
@@ -516,7 +524,7 @@ public class ChatService {
.description("负责拆解问题、规划步骤")
.model(chatModel)
.systemPrompt(prompt.toString())
.hooks(new AgentLoggingHook(agentStepRepository, "planner"))
.hooks(buildHooks("planner"))
.outputKey("planner_plan")
.build();
}
@@ -553,11 +561,30 @@ public class ChatService {
.systemPrompt(prompt.toString())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "executor"))
.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;
@@ -0,0 +1,38 @@
---
name: diagnose-aiops-alert
description: Diagnose AIOps alert payloads, active Prometheus alerts, alert scope control, HighCPUUsage, HighMemoryUsage, SlowResponse, ServiceUnavailable, and alert-driven incident reports. Use in AIOps flows or when the user asks to diagnose current alerts.
---
# AIOps Alert Diagnosis
## Workflow
1. Determine scope mode.
- Payload present: treat the supplied alert as the primary diagnosis target.
- No payload: call `queryPrometheusAlerts` first and choose P0/P1 or the longest-running firing alert.
2. For payload mode, preserve alert name, service, severity, description, and time range in the `lookup_knowledge` query.
3. Confirm active alert state with `queryPrometheusAlerts` when useful, but do not diagnose unrelated alerts as the main target.
4. Query metrics/logs that match the alert type and service.
5. Produce a report that distinguishes confirmed evidence, related risks, and missing evidence.
## Required Evidence
- Alert state from payload or `queryPrometheusAlerts`.
- `lookup_knowledge` when playbook or runbook guidance is needed.
- Logs and metrics aligned to the alert type.
## Stop Conditions
- If payload mode returns unrelated active alerts, mention them only as related risk.
- If three calls in the same direction fail or return no data, stop that direction and report the failure.
- Do not invent metric values, log lines, or remediation execution results.
## Report Rules
- Use the existing alert analysis report structure.
- Keep the supplied alert as the main diagnosis target in payload mode.
- Include confidence and evidence gaps.
## Eval Anchor
RAG cases: `aiops-payment-latency-alert`, `aiops-prometheus-alert-scope`.
@@ -0,0 +1,34 @@
---
name: diagnose-jvm-memory-risk
description: Diagnose JVM memory risk, high heap usage, OOM risk, OutOfMemoryError, frequent Full GC, memory leak, pod OOMKilled, or order-service memory alerts. Use when memory, JVM, heap, GC, OOM, or OOMKilled appears.
---
# JVM Memory Risk Diagnosis
## Workflow
1. Extract affected service, memory threshold, heap size, GC symptoms, pod/container events, and time window.
2. Call `query_metrics` or alert tools for heap usage, memory usage, GC count/time, and active memory alerts.
3. Call `query_logs` for Full GC warnings, OutOfMemoryError, OOMKilled, restart events, or allocation-heavy stack traces.
4. Call `lookup_knowledge` when JVM memory troubleshooting or remediation guidance is needed.
5. Decide whether the supported risk is high memory pressure, confirmed OOM, suspected leak, or insufficient evidence.
## Required Evidence
- `query_metrics` for resource pressure claims.
- `query_logs` for OOM, GC, or restart evidence.
## Stop Conditions
- High memory usage alone is not proof of memory leak.
- OOM risk is stronger when high memory metrics align with Full GC, OOMKilled, or OutOfMemoryError logs.
- If evidence is incomplete, return LOW_CONFID wording and list the missing metrics/logs.
## Report Rules
- Include immediate mitigation, heap/GC investigation, leak investigation, and monitoring recommendations.
- Do not say the issue can be ignored while memory remains above threshold.
## Eval Anchor
Fixed diagnosis case: `jvm-memory-risk`.
@@ -0,0 +1,35 @@
---
name: diagnose-mysql-connection-pool
description: Diagnose MySQL, HikariCP, database connection pool exhaustion, connection acquisition timeout, slow SQL, connection leak, or database saturation issues. Use when the user mentions MySQL pool, HikariCP, connection pool, database timeout, order-service timeout, or connection exhaustion.
---
# MySQL Connection Pool Diagnosis
## Workflow
1. Extract service, database, timeout symptom, and time window.
2. Call `lookup_knowledge` with MySQL, HikariCP, connection pool, and the affected service.
3. Call `query_logs` for connection acquisition timeout, active/max pool counts, waiting threads, leak warnings, slow query, or lock waits.
4. Call `query_metrics` when metrics are available for active connections, idle connections, wait time, DB latency, and error rate.
5. Decide whether the evidence supports pool exhaustion, slow SQL causing saturation, connection leak, or insufficient evidence.
## Required Evidence
- `lookup_knowledge` for pool configuration and diagnosis guidance.
- `query_logs` for concrete pool or SQL symptoms.
- `query_metrics` when making saturation or capacity claims.
## Stop Conditions
- Confirmed pool exhaustion requires log or metric evidence such as active equals max, waiting threads, acquisition timeout, or leak warnings.
- If only request timeout is present without pool evidence, state that the pool hypothesis is unconfirmed.
- If logs show slow SQL but not pool saturation, report slow SQL as the stronger supported cause.
## Report Rules
- Include current evidence, likely root cause, missing evidence, short-term mitigation, and long-term fix.
- Avoid saying "fully confirmed" unless at least two evidence sources align.
## Eval Anchor
Fixed diagnosis case: `mysql-pool-exhausted`.
@@ -0,0 +1,37 @@
---
name: diagnose-payment-timeout
description: Diagnose payment API, payment gateway, ERR_TIMEOUT, gateway timeout, payment-service latency, or payment request timeout issues. Use when the user mentions payment timeout, ERR_TIMEOUT, ERR_GATEWAY_TIMEOUT, slow payment, or payment-service latency.
---
# Payment Timeout Diagnosis
## Workflow
1. Identify the affected payment service, error code, endpoint, and time window from the user request.
2. Call `read_skill` only once for this playbook, then follow the evidence order below.
3. Call `lookup_knowledge` with a narrow query containing payment, timeout, the error code if present, and the affected service.
4. Call `query_logs` for payment-service timeout, downstream dependency timeout, gateway timeout, or request duration above threshold.
5. Call `query_metrics` or alert tools for latency, error rate, saturation, and active alerts when metrics are available.
6. Compare knowledge guidance with logs and metrics before stating a root cause.
## Required Evidence
- `lookup_knowledge` for error-code or payment timeout guidance.
- `query_logs` for concrete timeout or dependency evidence.
- `query_metrics` when the question asks for impact, latency, or current alert state.
## Stop Conditions
- If only knowledge is available and logs/metrics are missing, return LOW_CONFID language.
- If tools fail or return no evidence, state which evidence is missing and do not claim a confirmed root cause.
- Do not repeatedly call `lookup_knowledge` with synonym-only queries after a relevant result.
## Report Rules
- Separate immediate mitigation from long-term remediation.
- Cite the evidence source type for each key conclusion.
- Do not claim payment provider failure unless logs or metrics support an upstream dependency issue.
## Eval Anchor
Fixed diagnosis case: `payment-timeout`.
@@ -0,0 +1,34 @@
---
name: diagnose-redis-timeout
description: Diagnose Redis timeout, Redis connection timeout, cache dependency timeout, Redis cluster unavailable, hot key, network latency, or payment-service Redis dependency failures. Use when Redis or cache timeout appears in the user request, logs, or alert payload.
---
# Redis Timeout Diagnosis
## Workflow
1. Extract affected service, Redis operation, host/cluster, timeout value, and time window.
2. Call `lookup_knowledge` for Redis timeout or cache troubleshooting guidance when knowledge evidence is needed.
3. Call `query_logs` for Redis connection timeout, retry count, host, command latency, hot key, or dependency errors.
4. Call `query_metrics` when available for Redis latency, connection count, CPU, memory, error rate, or network saturation.
5. Distinguish client timeout, Redis saturation, network issue, and missing evidence.
## Required Evidence
- `query_logs` is mandatory for a concrete Redis timeout claim.
- `lookup_knowledge` is recommended for remediation and configuration guidance.
- `query_metrics` is required before claiming Redis resource saturation.
## Stop Conditions
- If only one Redis timeout log exists and no metrics are available, return LOW_CONFID wording.
- If Redis is only mentioned as a possible downstream dependency, do not make it the root cause without supporting logs.
## Report Rules
- State whether the supported issue is client-side timeout, Redis cluster issue, network issue, or unconfirmed.
- Include retry/backoff, timeout tuning, connection pool, and monitoring recommendations only when relevant.
## Eval Anchor
Fixed diagnosis case: `redis-timeout`.
@@ -0,0 +1,33 @@
---
name: diagnose-slow-response
description: Diagnose slow response, high P95/P99 latency, API latency regression, slow request, downstream latency, or user-service response time alerts. Use when the user mentions P99, P95, response time, slow endpoint, latency, or SlowResponse alerts.
---
# Slow Response Diagnosis
## Workflow
1. Extract service, endpoint, latency percentile, threshold, and time window.
2. Call `query_metrics` or alert tools to confirm latency and impact.
3. Call `query_logs` for slow request records, endpoint duration, downstream timing, cache misses, or database query timeout.
4. Call `lookup_knowledge` when process guidance, service-specific runbook, or known failure mode evidence is needed.
5. Classify the supported cause: database slow query, downstream dependency, cache miss, resource saturation, or insufficient evidence.
## Required Evidence
- `query_metrics` for latency or alert confirmation.
- `query_logs` for endpoint-level or dependency-level evidence.
## Stop Conditions
- If metrics show latency but logs do not identify a cause, say impact is confirmed but root cause is not.
- If logs identify slow SQL or dependency latency, use that as a candidate cause and mark confidence based on metric alignment.
## Report Rules
- Include impacted endpoints, observed latency, suspected bottleneck, evidence gaps, and next checks.
- Do not say there is no risk when P95/P99 remains above threshold.
## Eval Anchor
Fixed diagnosis case: `slow-response`.
@@ -56,17 +56,17 @@ class AiOpsServiceTest {
request.setSeverity("P1");
request.setTimeRange("last_15m");
request.setDescription("P95 latency is high");
request.setUserRequest("结合日志和指标排查支付超时");
request.setUserRequest("check logs and metrics for payment timeout");
String summary = service.buildQuerySummary(request);
assertTrue(summary.contains("AI Ops 告警分析"));
assertTrue(summary.contains("告警: payment-service-latency-high"));
assertTrue(summary.contains("服务: payment-service"));
assertTrue(summary.contains("等级: P1"));
assertTrue(summary.contains("时间范围: last_15m"));
assertTrue(summary.contains("描述: P95 latency is high"));
assertTrue(summary.contains("请求: 结合日志和指标排查支付超时"));
assertTrue(summary.contains("AI Ops alert analysis"));
assertTrue(summary.contains("alert: payment-service-latency-high"));
assertTrue(summary.contains("service: payment-service"));
assertTrue(summary.contains("severity: P1"));
assertTrue(summary.contains("timeRange: last_15m"));
assertTrue(summary.contains("description: P95 latency is high"));
assertTrue(summary.contains("request: "));
}
@Test
@@ -99,8 +99,8 @@ class AiOpsServiceTest {
assertTrue(prompt.contains("Related Risk"));
assertTrue(prompt.contains("Recommended lookup_knowledge query: HighCPUUsage payment-service P1 CPU usage is above 80% last_15m"));
assertTrue(prompt.contains("preserves alertName and service"));
assertTrue(prompt.contains("告警: HighCPUUsage"));
assertTrue(prompt.contains("服务: payment-service"));
assertTrue(prompt.contains("alert: HighCPUUsage"));
assertTrue(prompt.contains("service: payment-service"));
assertFalse(prompt.contains("AIOps scope mode: AUTO_DISCOVERY"));
}
@@ -112,11 +112,11 @@ class AiOpsServiceTest {
request.setSeverity(" ");
request.setDescription("P95 latency above threshold");
request.setTimeRange("last_10m");
request.setUserRequest("结合日志和指标排查");
request.setUserRequest("check logs and metrics");
String query = service.buildKnowledgeRetrievalQuery(request);
assertEquals("HighLatency payment-service P95 latency above threshold last_10m 结合日志和指标排查", query);
assertEquals("HighLatency payment-service P95 latency above threshold last_10m check logs and metrics", query);
}
@Test
@@ -142,16 +142,16 @@ class AiOpsServiceTest {
void persistFinalReportUpdatesDiagnosisSessionAnswer() {
DiagnosisSession session = DiagnosisSession.builder()
.sessionId("aiops-session-001")
.query("AI Ops 告警分析")
.query("AI Ops alert analysis")
.status("SUCCESS")
.agentFlow("AI_OPS")
.build();
when(diagnosisSessionRepository.findBySessionId("aiops-session-001")).thenReturn(Optional.of(session));
when(toolInvocationRepository.findBySessionIdOrderByIdAsc("aiops-session-001")).thenReturn(List.of());
service.persistFinalReport("aiops-session-001", "# 告警分析报告\nHighCPUUsage payment-service analysis with evidence summary.");
service.persistFinalReport("aiops-session-001", "# 闁告稑锕ㄩ鐔煎礆閸℃鈧粙骞庨妷銉﹀暈\nHighCPUUsage payment-service analysis with evidence summary.");
assertEquals("# 告警分析报告\nHighCPUUsage payment-service analysis with evidence summary.", session.getAnswer());
assertEquals("# 闁告稑锕ㄩ鐔煎礆閸℃鈧粙骞庨妷銉﹀暈\nHighCPUUsage payment-service analysis with evidence summary.", session.getAnswer());
assertTrue(session.getSelfEvaluation().contains("aiops_rule_evaluation"));
verify(diagnosisSessionRepository).save(session);
}
@@ -1,5 +1,8 @@
package com.superbiz.agent.service;
import com.alibaba.cloud.ai.graph.agent.ReactAgent;
import com.alibaba.cloud.ai.graph.skills.registry.SkillRegistry;
import com.alibaba.cloud.ai.graph.skills.registry.classpath.ClasspathSkillRegistry;
import com.superbiz.agent.agent.tool.DateTimeTools;
import com.superbiz.agent.agent.tool.QueryLogsTools;
import com.superbiz.agent.agent.tool.QueryMetricsTools;
@@ -194,6 +197,56 @@ class ChatServiceSequentialAgentTest {
assertSame(queryMetricsTools, methodTools[3]);
}
@Test
void createReactAgentInjectsSkillCatalogThroughAlibabaHook() throws Exception {
ChatService chatService = createChatService();
ScriptedChatModel chatModel = new ScriptedChatModel();
SkillRegistry skillRegistry = ClasspathSkillRegistry.builder()
.classpathPath("skills")
.basePath("target/test-skills-cache")
.build();
ReflectionTestUtils.setField(chatService, "skillRegistry", skillRegistry);
ReactAgent agent = chatService.createReactAgent(chatModel, "BASE_TEST_PROMPT");
agent.call("diagnose mysql connection pool exhaustion");
assertTrue(chatModel.promptText.contains("BASE_TEST_PROMPT"));
assertTrue(chatModel.promptText.contains("## Skills System"));
assertTrue(chatModel.promptText.contains("diagnose-mysql-connection-pool"));
assertTrue(chatModel.promptText.contains("read_skill"));
}
@Test
void plannerGetsSkillMetadataAndExecutorGetsReadSkillTool() throws Exception {
ChatService chatService = createChatService();
ScriptedChatModel chatModel = new ScriptedChatModel();
SkillRegistry skillRegistry = ClasspathSkillRegistry.builder()
.classpathPath("skills")
.basePath("target/test-skills-cache")
.build();
ReflectionTestUtils.setField(chatService, "skillRegistry", skillRegistry);
chatService.executeChatComplex(
chatModel,
new ToolCallback[0],
"diagnose mysql connection pool exhaustion",
List.of(),
"planner-skill-metadata-session"
);
assertTrue(chatModel.plannerPromptText.contains("\"skill_catalog\""));
assertTrue(chatModel.plannerPromptText.contains("diagnose-mysql-connection-pool"));
assertTrue(chatModel.plannerPromptText.contains("\"selected_skill\""));
assertFalse(chatModel.plannerPromptText.contains("## Skills System"));
assertFalse(chatModel.plannerPromptText.contains("read_skill"));
assertTrue(chatModel.executorPromptText.contains("## Skills System"));
assertTrue(chatModel.executorPromptText.contains("diagnose-mysql-connection-pool"));
assertTrue(chatModel.executorPromptText.contains("read_skill"));
assertFalse(chatModel.verifierPromptText.contains("diagnose-mysql-connection-pool"));
assertFalse(chatModel.verifierPromptText.contains("read_skill"));
}
private ChatService createChatService() {
ChatService chatService = new ChatService();
@@ -245,6 +298,9 @@ class ChatServiceSequentialAgentTest {
private static final class ScriptedChatModel implements ChatModel {
private final java.util.ArrayList<String> agentCalls = new java.util.ArrayList<>();
private String promptText = "";
private String plannerPromptText = "";
private String executorPromptText = "";
private String verifierPromptText = "";
private boolean sawVerifierPrompt;
private final java.util.List<String> verifierOutputs;
private int verifierOutputIndex;
@@ -283,12 +339,15 @@ class ChatServiceSequentialAgentTest {
String text;
if (promptText.contains("PLANNER_TEST_PROMPT")) {
agentCalls.add("chat_planner");
plannerPromptText = promptText;
text = "PLANNER_PLAN";
} else if (promptText.contains("EXECUTOR_TEST_PROMPT")) {
agentCalls.add("chat_executor");
executorPromptText = promptText;
text = "EXECUTOR_FINAL_ANSWER";
} else if (promptText.contains("VERIFIER_TEST_PROMPT")) {
agentCalls.add("chat_verifier");
verifierPromptText = promptText;
sawVerifierPrompt = true;
int index = Math.min(verifierOutputIndex, verifierOutputs.size() - 1);
text = verifierOutputs.get(index);
@@ -0,0 +1,44 @@
package com.superbiz.agent.service;
import com.alibaba.cloud.ai.graph.agent.hook.skills.ReadSkillTool;
import com.alibaba.cloud.ai.graph.skills.registry.SkillRegistry;
import com.superbiz.agent.config.SkillConfig;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class SkillCatalogServiceTest {
@Test
void loadsDiagnosisSkillsFromClasspathRegistry() {
SkillRegistry registry = newRegistry();
assertEquals(1, registry.size());
assertTrue(registry.contains("diagnose-mysql-connection-pool"));
}
@Test
void readSkillReturnsFullInstructionsFromOfficialTool() {
ReadSkillTool tool = new ReadSkillTool(newRegistry());
String skill = tool.apply(new ReadSkillTool.ReadSkillRequest("diagnose-mysql-connection-pool"), null);
assertTrue(skill.contains("## Workflow"));
assertTrue(skill.contains("query_logs"));
assertTrue(skill.contains("Fixed diagnosis case: `mysql-pool-exhausted`"));
}
@Test
void readSkillToolReturnsUnknownSkillError() {
ReadSkillTool tool = new ReadSkillTool(newRegistry());
String result = tool.apply(new ReadSkillTool.ReadSkillRequest("missing-skill"), null);
assertTrue(result.contains("Skill not found: missing-skill"));
}
private SkillRegistry newRegistry() {
return new SkillConfig().skillRegistry();
}
}