feat(adr-082): Phase 2 多 Agent 協作 — 5 角色辯證系統骨架上線

新增 5 個 Agent + Orchestrator + DecisionManager 接線:
- protocol.py: DiagnosisReport / ActionPlan / ReviewVerdict / CriticReport / DecisionPackage 型別系統
- DiagnosticianAgent: RCA 根因分析,confidence < 0.4 → ABSTAIN
- SolverAgent: 修復方案軍師,blast_radius 評分 + 降級 rule-based mock
- ReviewerAgent: 安全審查,HARD_RULES 靜態 pattern + blast_radius 閾值 (>50 revision, >80 reject)
- CriticAgent: 刻意唱反調,強制 3 問批判性思維,critical challenge → REJECT
- CoordinatorAgent: 純規則聚合,6 級決策閘,REQUEST_REVISION → 強制人工
- AgentOrchestrator: 30s 全局超時,Reviewer ‖ Critic 並行,DB Immutable Event Sourcing + Redis Streams
- DecisionManager: AIOPS_P2_ENABLED gate + _package_to_proposal_data 橋接既有 proposal_data 格式
- AgentSession DB table + 4 個複合 index
- ADR-082 決策記錄

Gate 2 修復(7 項):
- CRITICAL: DELETE FROM regex lookahead 位置錯誤(移至 FROM 後)
- CRITICAL: REQUEST_REVISION 可抵達 auto-execute 路徑(改回 requires_human_approval=True)
- IMPORTANT: _extract_json flat regex 不支援巢狀 JSON(改 find/rfind 邊界提取)
- IMPORTANT: all_degraded 遺漏 verdict.degraded(補全 4 個 Agent)
- IMPORTANT: Solver ABSTAIN guard 放行降級假設(改為無論 hypotheses 有無均跳過)
- IMPORTANT: dataclasses.asdict() Enum 未序列化導致 DB 寫入靜默失敗(加 json.dumps default handler)
- IMPORTANT: P2 gate 直讀屬性繞過父 Phase 守衛(改用 is_phase_enabled(2))

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
OG T
2026-04-15 13:33:13 +08:00
parent d51705b4ec
commit 5ddba6d6e0
11 changed files with 2171 additions and 4 deletions

View File

@@ -0,0 +1,366 @@
"""
AWOOOI AIOps Phase 2 — Agent Orchestrator排程指揮官
======================================================
職責:序列化 5 個 Agent 的執行,管理 Redis Streams 審計軌跡,記錄 DB
執行順序:
Diagnostician(snapshot)
→ Solver(diagnosis)
→ [Reviewer(plan) ‖ Critic(diagnosis, plan)] # 並行
→ Coordinator(diagnosis, plan, verdict, critic)
設計原則:
1. 每個 Agent 已有內建 5s 熔斷(在各自的 run() 中)— Orchestrator 不重複
2. 全流程 30s 全局超時GLOBAL_TIMEOUT_SEC
3. Redis Streams XADDbest-effort失敗不阻塞主流程
4. DB 記錄:每個 Agent turn 一行 agent_sessionsImmutable Event Sourcing
5. AIOPS_P2_ENABLED=False 時禁止使用(呼叫方負責 gate
Redis Stream key: aiops:p2:events
DB table: agent_sessions
ADR-082: Phase 2 多 Agent 協作
2026-04-15 ogt + Claude Sonnet 4.6(亞太): Phase 2 初始建立
"""
from __future__ import annotations
import asyncio
import dataclasses
import json
import time
import uuid
from datetime import UTC, datetime
from typing import TYPE_CHECKING
import structlog
from sqlalchemy import insert
from src.agents.coordinator_agent import get_coordinator_agent
from src.agents.critic_agent import compute_input_hash as critic_hash
from src.agents.critic_agent import get_critic_agent
from src.agents.diagnostician_agent import compute_input_hash as diag_hash
from src.agents.diagnostician_agent import get_diagnostician_agent
from src.agents.protocol import (
ActionPlan,
AgentRole,
AgentSessionStatus,
AgentVote,
CriticReport,
DecisionPackage,
DiagnosisReport,
ReviewVerdict,
)
from src.agents.reviewer_agent import compute_input_hash as reviewer_hash
from src.agents.reviewer_agent import get_reviewer_agent
from src.agents.solver_agent import compute_input_hash as solver_hash
from src.agents.solver_agent import get_solver_agent
from src.db.base import get_db_context
from src.db.models import AgentSession
if TYPE_CHECKING:
from src.services.evidence_snapshot import EvidenceSnapshot
logger = structlog.get_logger(__name__)
# 全局超時(所有 Agent 加起來)
GLOBAL_TIMEOUT_SEC = 30.0
# Redis Stream key
STREAM_KEY = "aiops:p2:events"
STREAM_MAXLEN = 10_000 # 保留最近 10k 條事件
# ─────────────────────────────────────────────────────────────────────────────
# Public API
# ─────────────────────────────────────────────────────────────────────────────
async def run_agent_debate(
snapshot: "EvidenceSnapshot",
incident_id: str,
) -> DecisionPackage:
"""
執行完整的 5 Agent 辯證流程。
Args:
snapshot: Phase 1 感官快照EvidenceSnapshot
incident_id: 關聯 incident ID
Returns:
DecisionPackage包含 recommended_action + confidence + requires_human_approval
Raises:
不拋出 — 所有錯誤內部吸收,最壞情況返回 requires_human_approval=True 的 Package
"""
session_id = str(uuid.uuid4())
start_ms = int(time.monotonic() * 1000)
logger.info(
"agent_debate_start",
session_id=session_id,
incident_id=incident_id,
snapshot_id=snapshot.snapshot_id,
)
try:
package = await asyncio.wait_for(
_debate(session_id, incident_id, snapshot),
timeout=GLOBAL_TIMEOUT_SEC,
)
elapsed = int(time.monotonic() * 1000) - start_ms
logger.info(
"agent_debate_done",
session_id=session_id,
incident_id=incident_id,
requires_human=package.requires_human_approval,
confidence=package.confidence,
status=package.session_status.value,
elapsed_ms=elapsed,
)
return package
except asyncio.TimeoutError:
elapsed = int(time.monotonic() * 1000) - start_ms
logger.error(
"agent_debate_global_timeout",
session_id=session_id,
elapsed_ms=elapsed,
timeout_sec=GLOBAL_TIMEOUT_SEC,
)
return _timeout_package(elapsed)
except Exception:
elapsed = int(time.monotonic() * 1000) - start_ms
logger.exception("agent_debate_fatal", session_id=session_id)
return _error_package(elapsed)
# ─────────────────────────────────────────────────────────────────────────────
# Internal — debate pipeline
# ─────────────────────────────────────────────────────────────────────────────
async def _debate(
session_id: str,
incident_id: str,
snapshot: "EvidenceSnapshot",
) -> DecisionPackage:
"""實際的 Agent 辯證流程(在全局超時保護內執行)。"""
# ── Step 1: Diagnostician ──────────────────────────────────────────────
diagnostician = get_diagnostician_agent()
diagnosis = await diagnostician.run(snapshot)
await _record_turn(
session_id=session_id,
incident_id=incident_id,
role=AgentRole.DIAGNOSTICIAN,
input_hash=diag_hash(snapshot),
output=diagnosis,
latency_ms=diagnosis.latency_ms,
vote=diagnosis.vote,
degraded=diagnosis.degraded,
)
# ── Step 2: Solver ─────────────────────────────────────────────────────
solver = get_solver_agent()
plan = await solver.run(diagnosis)
await _record_turn(
session_id=session_id,
incident_id=incident_id,
role=AgentRole.SOLVER,
input_hash=solver_hash(diagnosis),
output=plan,
latency_ms=plan.latency_ms,
vote=plan.vote,
degraded=plan.degraded,
)
# ── Step 3: Reviewer ‖ Critic並行──────────────────────────────────
reviewer = get_reviewer_agent()
critic = get_critic_agent()
verdict, critic_report = await asyncio.gather(
reviewer.run(plan),
critic.run(diagnosis, plan),
)
await asyncio.gather(
_record_turn(
session_id=session_id,
incident_id=incident_id,
role=AgentRole.REVIEWER,
input_hash=reviewer_hash(plan),
output=verdict,
latency_ms=verdict.latency_ms,
vote=verdict.vote,
degraded=verdict.degraded,
),
_record_turn(
session_id=session_id,
incident_id=incident_id,
role=AgentRole.CRITIC,
input_hash=critic_hash(diagnosis, plan),
output=critic_report,
latency_ms=critic_report.latency_ms,
vote=critic_report.vote,
degraded=critic_report.degraded,
),
return_exceptions=True, # MINOR-1: best-effort — 一個審計失敗不取消另一個
)
# ── Step 4: Coordinator ────────────────────────────────────────────────
coordinator = get_coordinator_agent()
package = await coordinator.run(diagnosis, plan, verdict, critic_report)
await _record_turn(
session_id=session_id,
incident_id=incident_id,
role=AgentRole.COORDINATOR,
input_hash=_coordinator_input_hash(diagnosis, plan, verdict, critic_report),
output=package,
latency_ms=package.latency_ms,
vote=AgentVote.APPROVE if not package.requires_human_approval else AgentVote.ABSTAIN,
degraded=package.all_agents_degraded,
)
return package
# ─────────────────────────────────────────────────────────────────────────────
# Internal — DB + Redis recording
# ─────────────────────────────────────────────────────────────────────────────
async def _record_turn(
session_id: str,
incident_id: str,
role: AgentRole,
input_hash: str,
output: object,
latency_ms: int,
vote: AgentVote,
degraded: bool,
) -> None:
"""
寫入 agent_sessions DB 行 + Redis Stream XADD兩者皆 best-effort
Immutable Event Sourcing — 只 INSERT永不 UPDATE / DELETE。
"""
output_json = _serialize_output(output)
# ── DB 寫入 ──────────────────────────────────────────────────────────
try:
async with get_db_context() as db:
await db.execute(
insert(AgentSession).values(
id=str(uuid.uuid4()),
session_id=session_id,
incident_id=incident_id,
agent_role=role.value,
input_hash=input_hash,
output_json=output_json,
latency_ms=latency_ms,
vote=vote.value,
degraded=degraded,
created_at=datetime.now(UTC),
)
)
except Exception:
logger.exception(
"agent_session_db_write_failed",
session_id=session_id,
role=role.value,
)
# ── Redis Stream XADDbest-effort失敗靜默────────────────────────
try:
from src.core.redis_client import get_redis
redis = get_redis()
await redis.xadd(
STREAM_KEY,
{
"session_id": session_id,
"incident_id": incident_id,
"role": role.value,
"vote": vote.value,
"latency_ms": str(latency_ms),
"degraded": "1" if degraded else "0",
},
maxlen=STREAM_MAXLEN,
approximate=True,
)
except Exception:
# Redis 失敗不阻塞主流程(觀測性非關鍵路徑)
logger.warning(
"agent_stream_xadd_failed",
session_id=session_id,
role=role.value,
)
def _serialize_output(output: object) -> dict:
"""
將 Agent 輸出 dataclass 轉為可 JSON 序列化的 dict。
Gate 2: dataclasses.asdict() 保留 Enum 實例json.dumps 會拋 TypeError。
用 json.dumps(default=...) 先把 Enum 轉 .value 再 loads 回 dict。
"""
import json
from enum import Enum
try:
if dataclasses.is_dataclass(output) and not isinstance(output, type):
raw = dataclasses.asdict(output) # type: ignore[arg-type]
return json.loads(
json.dumps(raw, default=lambda o: o.value if isinstance(o, Enum) else str(o))
)
return {"raw": str(output)[:2000]}
except Exception:
return {"error": "serialization_failed"}
def _coordinator_input_hash(
diagnosis: DiagnosisReport,
plan: ActionPlan,
verdict: ReviewVerdict,
critic: CriticReport,
) -> str:
"""Coordinator 輸入 fingerprint4 個 Agent 輸出的組合)。"""
import hashlib
key = (
diagnosis.evidence_snapshot_id
+ (diagnosis.top_hypothesis.description if diagnosis.top_hypothesis else "")
+ (plan.top_candidate.action if plan.top_candidate else "")
+ verdict.vote.value
+ str(critic.challenge_count)
)
return hashlib.sha256(key.encode()).hexdigest()[:16]
# ─────────────────────────────────────────────────────────────────────────────
# Degraded fallback packages
# ─────────────────────────────────────────────────────────────────────────────
def _timeout_package(elapsed_ms: int) -> DecisionPackage:
"""全局超時 → 強制人工審核。"""
return DecisionPackage(
recommended_action=None,
confidence=0.0,
requires_human_approval=True,
debate_summary=f"[超時] 辯證流程超過 {GLOBAL_TIMEOUT_SEC}s強制升級人工審核",
session_status=AgentSessionStatus.TIMEOUT,
latency_ms=elapsed_ms,
blocked_reason=f"全局超時 > {GLOBAL_TIMEOUT_SEC}s",
all_agents_degraded=True,
)
def _error_package(elapsed_ms: int) -> DecisionPackage:
"""未預期錯誤 → 強制人工審核。"""
return DecisionPackage(
recommended_action=None,
confidence=0.0,
requires_human_approval=True,
debate_summary="[錯誤] 辯證流程發生未預期錯誤,強制升級人工審核",
session_status=AgentSessionStatus.FAILED,
latency_ms=elapsed_ms,
blocked_reason="Orchestrator 未預期錯誤",
all_agents_degraded=True,
)

View File

@@ -942,6 +942,52 @@ EXPERT_RULES: dict[str, dict[str, Any]] = {
}
def _package_to_proposal_data(package: Any) -> dict[str, Any]:
"""
將 Phase 2 DecisionPackage 轉換為 proposal_data dict。
proposal_data 格式是既有 decision_manager 決策流程的共同語言,
所有下游auto_approve / Telegram / approval_record都依此格式讀取。
ADR-082: Phase 2 多 Agent 協作
2026-04-15 ogt + Claude Sonnet 4.6(亞太)
"""
action = package.recommended_action or ""
confidence = package.confidence
# requires_human_approval → risk_level 映射
# Phase 2 Reviewer + Critic 已做安全審查,不需要 human = low/medium
if package.requires_human_approval:
risk_level = "high" if confidence > 0 else "critical"
elif confidence >= 0.8:
risk_level = "low"
elif confidence >= 0.6:
risk_level = "medium"
else:
risk_level = "high"
return {
"action": action,
"kubectl_command": action if action.startswith("kubectl") else "",
"description": package.debate_summary[:500] if package.debate_summary else "",
"reasoning": package.debate_summary[:1000] if package.debate_summary else "",
"confidence": confidence,
"risk_level": risk_level,
"source": "phase2_agent_debate",
"requires_human_review": package.requires_human_approval,
# Phase 2 診斷摘要(供 Telegram 卡片顯示)
"debate_summary": package.debate_summary or "",
"all_agents_degraded": package.all_agents_degraded,
"blocked_reason": package.blocked_reason or "",
"session_status": package.session_status.value if package.session_status else "",
# MINOR-2: 補齊 Expert System / LLM 路徑存在的 keys防止下游 .get() 靜默空缺
"is_rule_based": False,
"matched_rule": "",
"rule_id": "",
"from_cache": False,
}
def expert_analyze(incident: Incident) -> dict[str, Any]:
"""
Expert System 規則引擎分析
@@ -1670,6 +1716,31 @@ class DecisionManager:
# ADR-070: 原有 MCP 收集路徑Phase 0 保留)
mcp_context = await self._collect_mcp_context(incident)
# ADR-082 Phase 2: 5 Agent 辯證feature flag 守衛)
# AIOPS_P2_ENABLED=True → 走 AgentOrchestrator 路徑,跳過 Playbook / LLM
# 需要 EvidenceSnapshot若 P1 未開啟則自行收集
# 2026-04-15 ogt + Claude Sonnet 4.6(亞太)
if aiops_flags.is_phase_enabled(2): # Gate 2: 用 is_phase_enabled 統一父 Phase 守衛
p2_snapshot = evidence_snapshot
if p2_snapshot is None:
try:
from src.services.pre_decision_investigator import get_pre_decision_investigator
p2_snapshot = await get_pre_decision_investigator().investigate(incident)
except Exception:
logger.warning(
"p2_snapshot_collect_failed",
incident_id=incident.incident_id,
)
if p2_snapshot is not None:
from src.services.agent_orchestrator import run_agent_debate
package = await run_agent_debate(
snapshot=p2_snapshot,
incident_id=incident.incident_id,
)
return _package_to_proposal_data(package)
# snapshot 仍為 None → 降級繼續走原路徑(不阻塞)
logger.warning("p2_no_snapshot_fallback", incident_id=incident.incident_id)
# Phase 7.5: 先嘗試 Playbook 匹配
playbook_result = await self._try_playbook_match(incident)
if playbook_result: