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.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; 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.hook.PlannerSkillMetadataHook; 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.Map; import java.util.Optional; import java.util.UUID; /** * AI Ops 智能运维服务。 * 负责构建 Planner、Executor、Supervisor 多 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; @Autowired private AiOpsRuleEvaluationService aiOpsRuleEvaluationService; @Autowired private SelfEvaluationMergeService selfEvaluationMergeService; @Autowired(required = false) private SkillRegistry skillRegistry; /** * 执行 AI Ops 告警分析流程。 * * @param chatModel 大模型实例 * @param toolCallbacks Spring AI 工具回调 * @return 多 Agent 编排后的最终状态 * @throws GraphRunnerException Agent 图执行失败时抛出 */ public Optional executeAiOpsAnalysis(ChatModel chatModel, ToolCallback[] toolCallbacks) throws GraphRunnerException { return executeAiOpsAnalysis(chatModel, toolCallbacks, null, resolveSessionId(null)); } public Optional executeAiOpsAnalysis(ChatModel chatModel, ToolCallback[] toolCallbacks, AIOpsRequest request, String sessionId) throws GraphRunnerException { logger.info("Starting AI Ops multi-agent analysis"); String resolvedSessionId = isBlank(sessionId) ? resolveSessionId(request) : sessionId.trim(); long startTime = System.currentTimeMillis(); DiagnosisSession session = startDiagnosisSession(resolvedSessionId, request); diagnosisSessionRepository.save(session); // 让工具调用、Hook 和知识库检索能够拿到当前诊断会话 ID。 SessionContextHolder.setSessionId(resolvedSessionId); try { ReactAgent plannerAgent = buildPlannerAgent(chatModel, toolCallbacks); ReactAgent executorAgent = buildExecutorAgent(chatModel, toolCallbacks); SupervisorAgent supervisorAgent = SupervisorAgent.builder() .name("ai_ops_supervisor") .description("Coordinates Planner and Executor agents") .model(chatModel) .systemPrompt(promptProperties.getSupervisor()) .subAgents(List.of(plannerAgent, executorAgent)) .build(); String taskPrompt = buildTaskPrompt(request); logger.info("Invoking AI Ops supervisor agent"); Optional 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(); } } /** * 从多 Agent 执行状态中提取最终报告文本。 * * @param state Agent 图执行状态 * @return Planner 输出中的最终报告 */ public Optional extractFinalReport(OverAllState state) { logger.info("Extracting final AI Ops report"); Optional plannerFinalOutput = state.value("planner_plan") .filter(AssistantMessage.class::isInstance) .map(AssistantMessage.class::cast); if (plannerFinalOutput.isPresent()) { String reportText = plannerFinalOutput.get().getText(); logger.info("Extracted Planner final report, length: {}", reportText.length()); return Optional.of(reportText); } else { logger.warn("Unable to extract Planner final report"); 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) { persistFinalReport(sessionId, finalReport, null); } public void persistFinalReport(String sessionId, String finalReport, AIOpsRequest request) { if (isBlank(sessionId) || isBlank(finalReport)) { return; } diagnosisSessionRepository.findBySessionId(sessionId.trim()).ifPresent(session -> { session.setAnswer(finalReport); List invocations = toolInvocationRepository.findBySessionIdOrderByIdAsc(session.getSessionId()); Map evaluation = aiOpsRuleEvaluationService.evaluate(request, finalReport, invocations); session.setSelfEvaluation(selfEvaluationMergeService.mergeAiOpsRuleEvaluation( session.getSelfEvaluation(), evaluation)); diagnosisSessionRepository.save(session); }); } String buildQuerySummary(AIOpsRequest request) { if (request == null) { return "AI Ops alert analysis"; } 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(); } 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 buildKnowledgeRetrievalQuery(AIOpsRequest request) { if (request == null || !hasAlertPayload(request)) { return ""; } StringBuilder query = new StringBuilder(); appendQueryTerm(query, request.getAlertName()); appendQueryTerm(query, request.getService()); appendQueryTerm(query, request.getSeverity()); appendQueryTerm(query, request.getDescription()); appendQueryTerm(query, request.getTimeRange()); appendQueryTerm(query, request.getUserRequest()); return query.toString(); } String buildTaskPrompt(AIOpsRequest request) { StringBuilder prompt = new StringBuilder(); 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); 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("- Recommended lookup_knowledge query: ").append(knowledgeQuery).append("\n"); prompt.append("- If knowledge-base evidence is needed, call lookup_knowledge with the recommended query or a narrower query that preserves alertName and service.\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("Plans alert diagnosis steps") .model(chatModel) .systemPrompt(promptProperties.getPlanner()) .hooks(buildHooks("planner")) .outputKey("planner_plan") .build(); } /** * 构建 Executor Agent。 */ private ReactAgent buildExecutorAgent(ChatModel chatModel, ToolCallback[] toolCallbacks) { return ReactAgent.builder() .name("executor_agent") .description("Executes the current Planner step and reports feedback") .model(chatModel) .systemPrompt(promptProperties.getExecutor()) .methodTools(buildMethodToolsArray()) .tools(toolCallbacks) .hooks(buildHooks("executor")) .outputKey("executor_feedback") .build(); } /** * 根据运行模式构建方法工具数组。 * Mock 模式注入本地 QueryLogsTools;真实模式下日志查询由外部 MCP 工具提供。 */ private Object[] buildMethodToolsArray() { if (queryLogsTools != null) { return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools, queryLogsTools}; } return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools}; } 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 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("Failed to backfill AI Ops session metrics, 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 void appendQueryTerm(StringBuilder builder, String value) { if (!isBlank(value)) { if (!builder.isEmpty()) { builder.append(' '); } builder.append(value.trim()); } } private boolean isBlank(String value) { return value == null || value.trim().isEmpty(); } }