Files
SuperBizAgent-java/src/main/java/com/superbiz/agent/service/AiOpsService.java
T

323 lines
14 KiB
Java
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package com.superbiz.agent.service;
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.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.AgentStep;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import com.superbiz.agent.dto.AIOpsRequest;
import com.superbiz.agent.hook.AgentLoggingHook;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.repository.ToolInvocationRepository;
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 智能运维服务
* 负责多 Agent 协作的告警分析流程
*/
@Service
public class AiOpsService {
private static final Logger logger = LoggerFactory.getLogger(AiOpsService.class);
@Autowired
private DateTimeTools dateTimeTools;
@Autowired
private InternalDocsTools internalDocsTools;
@Autowired
private QueryMetricsTools queryMetricsTools;
@Autowired(required = false) // Mock 模式下才注册
private QueryLogsTools queryLogsTools;
@Autowired
private LookupKnowledgeTool lookupKnowledgeTool;
@Autowired
private AiOpsPromptProperties promptProperties;
@Autowired
private DiagnosisSessionRepository diagnosisSessionRepository;
@Autowired
private AgentStepRepository agentStepRepository;
@Autowired
private ToolInvocationRepository toolInvocationRepository;
/**
* 执行 AI Ops 告警分析流程
*
* @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));
}
public Optional<OverAllState> executeAiOpsAnalysis(ChatModel chatModel, ToolCallback[] toolCallbacks,
AIOpsRequest request, String sessionId) throws GraphRunnerException {
logger.info("开始执行 AI Ops 多 Agent 协作流程");
String resolvedSessionId = isBlank(sessionId) ? resolveSessionId(request) : sessionId.trim();
long startTime = System.currentTimeMillis();
// 创建或更新诊断会话
DiagnosisSession session = startDiagnosisSession(resolvedSessionId, request);
diagnosisSessionRepository.save(session);
// 设置 ThreadLocal 上下文(LookupKnowledgeTool 通过此获取 sessionId)
SessionContextHolder.setSessionId(resolvedSessionId);
try {
// 构建 Planner 和 Executor Agent(每个 Agent 各自带 Hook)
ReactAgent plannerAgent = buildPlannerAgent(chatModel, toolCallbacks);
ReactAgent executorAgent = buildExecutorAgent(chatModel, toolCallbacks);
// 构建 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();
String taskPrompt = buildTaskPrompt(request);
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();
}
}
/**
* 从执行结果中提取最终报告文本
*
* @param state 执行状态
* @return 报告文本(如果存在)
*/
public Optional<String> extractFinalReport(OverAllState state) {
logger.info("开始提取最终报告...");
// 提取 Planner 最终输出(包含完整的告警分析报告)
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());
return Optional.of(reportText);
} else {
logger.warn("未能提取到 Planner 最终报告");
return Optional.empty();
}
}
public String resolveSessionId(AIOpsRequest request) {
if (request != null && !isBlank(request.getSessionId())) {
return request.getSessionId().trim();
}
return UUID.randomUUID().toString();
}
public void persistFinalReport(String sessionId, String finalReport) {
if (isBlank(sessionId) || isBlank(finalReport)) {
return;
}
diagnosisSessionRepository.findBySessionId(sessionId.trim()).ifPresent(session -> {
session.setAnswer(finalReport);
diagnosisSessionRepository.save(session);
});
}
String buildQuerySummary(AIOpsRequest request) {
if (request == null) {
return "AI Ops 告警分析";
}
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());
return summary.toString();
}
boolean hasAlertPayload(AIOpsRequest request) {
if (request == null) {
return false;
}
return !isBlank(request.getAlertName())
|| !isBlank(request.getService())
|| !isBlank(request.getSeverity())
|| !isBlank(request.getDescription())
|| !isBlank(request.getTimeRange());
}
String buildTaskPrompt(AIOpsRequest request) {
StringBuilder prompt = new StringBuilder();
prompt.append("你是企业级 SRE,接到了自动化告警排查任务。请结合工具调用,执行**规划→执行→再规划**的闭环,并最终按照固定模板输出《告警分析报告》。禁止编造虚假数据,如连续多次查询失败需诚实反馈无法完成的原因。");
prompt.append("\n\n本次告警输入:\n");
prompt.append(buildQuerySummary(request));
if (hasAlertPayload(request)) {
prompt.append("\n\nAIOps scope mode: PAYLOAD_TARGETED\n");
prompt.append("- The request includes an alert payload. Treat the supplied alert payload as the primary and only main diagnosis target.\n");
prompt.append("- The final report must focus on the supplied alert fields such as alertName, service, severity, description, and timeRange.\n");
prompt.append("- You may call queryPrometheusAlerts only to verify whether the supplied alert is still active or to identify related risk/context.\n");
prompt.append("- If queryPrometheusAlerts returns unrelated active alerts, do not create full root-cause or remediation sections for them.\n");
prompt.append("- Mention unrelated active alerts only briefly in a Related Risk section when they help explain the supplied alert.\n");
} else {
prompt.append("\n\nAIOps scope mode: AUTO_DISCOVERY\n");
prompt.append("- The request does not include alert payload fields. First call queryPrometheusAlerts to discover current active/firing alerts.\n");
prompt.append("- Prefer P0/P1 alerts or the longest-running firing alerts, then diagnose one or more alerts based on severity and evidence.\n");
prompt.append("- Use metrics, logs, and knowledge-base evidence before producing the final alert analysis report.\n");
}
return prompt.toString();
}
private DiagnosisSession startDiagnosisSession(String sessionId, AIOpsRequest request) {
DiagnosisSession session = diagnosisSessionRepository.findBySessionId(sessionId)
.orElseGet(() -> DiagnosisSession.builder()
.sessionId(sessionId)
.agentFlow("AI_OPS")
.build());
session.setQuery(buildQuerySummary(request));
session.setStatus("RUNNING");
session.setAgentFlow("AI_OPS");
session.setAnswer(null);
session.setTotalDurationMs(null);
session.setTotalTokenCount(null);
session.setStepCount(null);
session.setToolCallCount(null);
return session;
}
/**
* 构建 Planner Agent
*/
private ReactAgent buildPlannerAgent(ChatModel chatModel, ToolCallback[] toolCallbacks) {
return ReactAgent.builder()
.name("planner_agent")
.description("负责拆解告警、规划与再规划步骤")
.model(chatModel)
.systemPrompt(promptProperties.getPlanner())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "planner"))
.outputKey("planner_plan")
.build();
}
/**
* 构建 Executor Agent
*/
private ReactAgent buildExecutorAgent(ChatModel chatModel, ToolCallback[] toolCallbacks) {
return ReactAgent.builder()
.name("executor_agent")
.description("负责执行 Planner 的首个步骤并及时反馈")
.model(chatModel)
.systemPrompt(promptProperties.getExecutor())
.methodTools(buildMethodToolsArray())
.tools(toolCallbacks)
.hooks(new AgentLoggingHook(agentStepRepository, "executor"))
.outputKey("executor_feedback")
.build();
}
/**
* 动态构建方法工具数组
* 根据 cls.mock-enabled 决定是否包含 QueryLogsTools
* 工具顺序:知识库查询优先,日志查询次之,弃用工具最后
*/
private Object[] buildMethodToolsArray() {
if (queryLogsTools != null) {
// Mock 模式:包含 QueryLogsTools
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools, queryLogsTools};
} else {
// 真实模式:不包含 QueryLogsTools(由 MCP 提供日志查询功能)
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools};
}
}
/** 从 agent_step 和 tool_invocation 汇总指标回填 diagnosis_session */
private void backfillSessionMetrics(DiagnosisSession session) {
try {
List<AgentStep> steps = agentStepRepository.findBySessionIdOrderByStepIndex(session.getSessionId());
int totalTokens = 0;
int stepCount = 0;
for (AgentStep s : steps) {
stepCount++;
if (s.getTokenCount() != null) totalTokens += s.getTokenCount();
}
long toolCallCount = toolInvocationRepository.countBySessionId(session.getSessionId());
session.setTotalTokenCount(totalTokens);
session.setStepCount(stepCount);
session.setToolCallCount(Math.toIntExact(toolCallCount));
} catch (Exception e) {
logger.warn("回填会话指标失败: sessionId={}", session.getSessionId(), e);
}
}
private void appendField(StringBuilder builder, String label, String value) {
if (!isBlank(value)) {
builder.append("\n- ").append(label).append(": ").append(value.trim());
}
}
private boolean isBlank(String value) {
return value == null || value.trim().isEmpty();
}
}