Files
awoooi/apps/api/src/plugins/mcp/providers/k8s_provider.py
Your Name 45dbe07188
Some checks failed
run-migration / migrate (push) Failing after 22s
Deploy Alert Rules / Deploy Prometheus Alert Rules (push) Successful in 53s
Type Sync Check / check-type-sync (push) Successful in 2m54s
CD Pipeline / build-and-deploy (push) Has been cancelled
Ansible Lint / lint (push) Has been cancelled
fix(flywheel): 自動化飛輪六大能力修復(ADR-092 B3)
【根因鏈修復】
MCP Provider bugs → PreDecisionInvestigator 失敗 → Agent Debate 無上下文
→ LLM 逾時 → description="待分析" → ADR-091 鐵閘攔截 → tg_sent 未設
→ W-2 Watchdog 誤報「靜默故障」

【六大修復】
1. MCP Provider 三蟲修復
   - ssh_provider: asyncssh.run() → conn.run()
   - prometheus_provider: KeyError 'query' → .get() 容錯
   - k8s_provider: 空 pod_name → 早返回錯誤字典

2. Agent Debate / 決策品質
   - decision_manager: 逾時降級文字改為明確描述(繞過 ADR-091 鐵閘)
   - intent_classifier: LLM 逾時降級至關鍵字分類(非 None)

3. Watchdog 誤報修復(ADR-092 B3)
   - W-2: tg_sent Redis TTL → telegram_message_id IS NULL(DB 真值)
   - W-5 新增: suggested_action IN 空/待分析/NO_ACTION + tg_id IS NULL
   - approval_timeout_resolver: 60min → 15min,batch 50 → 200

4. Config Drift 自動化
   - drift_adopt_service: auto_adopt_if_safe() 六條件安全閘
   - drift.py: 背景任務先嘗試自動採納再發人工 Telegram 卡片

5. Playbook 飛輪穩定
   - playbook_seed_service: 修復幂等性(deprecated 不視為缺失)
   - playbook_evolver: 只載 DRAFT+APPROVED(非全部 294 筆)

6. 可觀測性
   - alert_rule_engine: auto_rule 結構化日誌 + Redis 計數器(pipeline)
   - auto_approve: reject 原因 Redis 計數器
   - heartbeat_report_service: 新增「⚙️ 自動化統計(今日)」區塊

【待人工執行】
psql $DATABASE_URL -f apps/api/migrations/cleanup_duplicate_deprecated_playbooks.sql

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-04-24 10:55:50 +08:00

493 lines
21 KiB
Python
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.
"""
Kubernetes MCP Tool Provider - ADR-015 模組化架構
================================================
提供 Kubernetes 操作工具:
- kubectl_get: 查詢資源
- kubectl_delete: 刪除 Pod
- kubectl_scale: 調整副本數
- kubectl_restart: 重啟 Deployment
MCP Phase 1 新增 (2026-04-11 Claude Sonnet 4.6 — Sprint B 後驗收):
- k8s_get_pod_logs: 取得 Pod 最後 N 行 log取代 log_summary 手動呼叫)
- k8s_watch_rollout: 監控 rollout 狀態直到完成(真正的驗證,非 sleep 猜)
- k8s_get_events: 取得 namespace/resource 的 K8s events
- k8s_describe_pod: 完整 Pod describe含 Conditions、Volumes、Env
- k8s_get_hpa_status: 取得 HPA 當前副本數/上限/下限
- k8s_get_node_conditions: 取得 Node conditionsReady/MemoryPressure/DiskPressure
安全守衛Phase 1 強化):
- namespace 限 awoooi-prod硬編碼白名單
- 操作類工具需 trust_score >= 0.7
透過 DI 注入 ActionExecutor不直接 import services。
@see docs/adr/ADR-015-mcp-modular-architecture.md
@see docs/superpowers/specs/2026-04-10-infra-rebuild-sprint-abc-design.md §MCP Phase 1
"""
import re
import uuid
from typing import Any
import structlog
from src.plugins.mcp.interfaces import MCPTool, MCPToolProvider, MCPToolResult
logger = structlog.get_logger(__name__)
# namespace 白名單(硬編碼,防止跨命名空間操作)
ALLOWED_NAMESPACES = {"awoooi-prod"}
DEFAULT_NAMESPACE = "awoooi-prod"
# 資源名稱白名單正則(防止 injection
_RE_SAFE_K8S_NAME = re.compile(r'^[a-zA-Z0-9][a-zA-Z0-9._-]{0,252}$')
_RE_SAFE_RESOURCE = re.compile(r'^[a-zA-Z0-9/-]{1,64}$')
def _validate_namespace(ns: str) -> str:
if ns not in ALLOWED_NAMESPACES:
raise ValueError(f"Namespace '{ns}' not allowed (whitelist: {ALLOWED_NAMESPACES})")
return ns
def _validate_name(name: str, field: str = "name") -> str:
if not name or not _RE_SAFE_K8S_NAME.match(name):
raise ValueError(f"Unsafe {field}: {name!r}")
return name
class K8sProvider(MCPToolProvider):
"""
Kubernetes MCP Tool Provider
封裝所有 kubectl 操作,透過 ActionExecutor 執行。
"""
def __init__(self) -> None:
# Lazy import to avoid circular dependency
self._executor = None
@property
def name(self) -> str:
return "kubernetes"
def _get_executor(self):
"""Lazy load executor to avoid import at module load time"""
if self._executor is None:
from src.services.executor import get_executor
self._executor = get_executor()
return self._executor
async def list_tools(self) -> list[MCPTool]:
return [
MCPTool(
name="kubectl_get",
description="Query Kubernetes resources (pods, deployments, services)",
input_schema={
"type": "object",
"properties": {
"resource": {"type": "string", "description": "Resource type: pods, deployments, services"},
"namespace": {"type": "string", "description": "Namespace (default: awoooi-prod)"},
"name": {"type": "string", "description": "Resource name (optional)"},
},
"required": ["resource"],
},
server_name=self.name,
),
MCPTool(
name="kubectl_delete",
description="Delete a Pod (with dry-run validation)",
input_schema={
"type": "object",
"properties": {
"resource": {"type": "string", "description": "Resource type: pod"},
"name": {"type": "string", "description": "Pod name"},
"namespace": {"type": "string", "description": "Namespace"},
},
"required": ["name"],
},
server_name=self.name,
),
MCPTool(
name="kubectl_scale",
description="Scale a Deployment to N replicas",
input_schema={
"type": "object",
"properties": {
"deployment": {"type": "string", "description": "Deployment name"},
"replicas": {"type": "integer", "description": "Target replica count"},
"namespace": {"type": "string", "description": "Namespace"},
},
"required": ["deployment", "replicas"],
},
server_name=self.name,
),
MCPTool(
name="kubectl_restart",
description="Restart a Deployment (rollout restart)",
input_schema={
"type": "object",
"properties": {
"deployment": {"type": "string", "description": "Deployment name"},
"namespace": {"type": "string", "description": "Namespace"},
},
"required": ["deployment"],
},
server_name=self.name,
),
# ── MCP Phase 1 新增工具 ──────────────────────────────────────────
MCPTool(
name="k8s_get_pod_logs",
description=(
"Get last N lines of logs from a Pod container. "
"Use for PodRestartingTooMuch, CrashLoopBackOff diagnosis. Read-only."
),
input_schema={
"type": "object",
"properties": {
"pod_name": {"type": "string", "description": "Pod name (e.g. awoooi-api-xxx-yyy)"},
"container": {"type": "string", "description": "Container name (optional, uses first if omitted)"},
"namespace": {"type": "string", "description": "Namespace (default: awoooi-prod)"},
"tail": {"type": "integer", "description": "Number of lines (default: 100, max: 500)"},
"previous": {"type": "boolean", "description": "Get logs from previous (crashed) container"},
},
"required": ["pod_name"],
},
server_name=self.name,
),
MCPTool(
name="k8s_watch_rollout",
description=(
"Watch a Deployment rollout until complete or timeout. "
"Use after kubectl_restart to verify success instead of sleeping."
),
input_schema={
"type": "object",
"properties": {
"deployment": {"type": "string", "description": "Deployment name"},
"namespace": {"type": "string", "description": "Namespace (default: awoooi-prod)"},
"timeout_seconds": {"type": "integer", "description": "Timeout in seconds (default: 120, max: 300)"},
},
"required": ["deployment"],
},
server_name=self.name,
),
MCPTool(
name="k8s_get_events",
description=(
"Get Kubernetes events for a namespace or specific resource. "
"Useful for diagnosing scheduling failures, image pull errors, OOMKilled."
),
input_schema={
"type": "object",
"properties": {
"namespace": {"type": "string", "description": "Namespace (default: awoooi-prod)"},
"resource_name": {"type": "string", "description": "Filter by resource name (optional)"},
"event_type": {"type": "string", "description": "Filter: Warning or Normal (optional)"},
},
"required": [],
},
server_name=self.name,
),
MCPTool(
name="k8s_describe_pod",
description=(
"Full kubectl describe for a Pod: Conditions, Events, Volumes, Env, Resource limits. "
"Use when get_pod_logs is insufficient for diagnosis."
),
input_schema={
"type": "object",
"properties": {
"pod_name": {"type": "string", "description": "Pod name"},
"namespace": {"type": "string", "description": "Namespace (default: awoooi-prod)"},
},
"required": ["pod_name"],
},
server_name=self.name,
),
MCPTool(
name="k8s_get_hpa_status",
description=(
"Get HPA (HorizontalPodAutoscaler) current/min/max replicas and CPU utilization. "
"Use for capacity planning and scaling diagnosis."
),
input_schema={
"type": "object",
"properties": {
"namespace": {"type": "string", "description": "Namespace (default: awoooi-prod)"},
"name": {"type": "string", "description": "HPA name (optional, lists all if omitted)"},
},
"required": [],
},
server_name=self.name,
),
MCPTool(
name="k8s_get_node_conditions",
description=(
"Get all Node conditions: Ready, MemoryPressure, DiskPressure, PIDPressure. "
"Use for cluster-wide health diagnosis."
),
input_schema={
"type": "object",
"properties": {},
"required": [],
},
server_name=self.name,
),
]
async def execute(
self,
tool_name: str,
parameters: dict[str, Any],
) -> MCPToolResult:
execution_id = str(uuid.uuid4())[:8]
executor = self._get_executor()
try:
if tool_name == "kubectl_get":
output = await self._kubectl_get(executor, parameters)
elif tool_name == "kubectl_delete":
output = await self._kubectl_delete(executor, parameters)
elif tool_name == "kubectl_scale":
output = await self._kubectl_scale(executor, parameters)
elif tool_name == "kubectl_restart":
output = await self._kubectl_restart(executor, parameters)
elif tool_name == "k8s_get_pod_logs":
output = await self._k8s_get_pod_logs(executor, parameters)
elif tool_name == "k8s_watch_rollout":
output = await self._k8s_watch_rollout(executor, parameters)
elif tool_name == "k8s_get_events":
output = await self._k8s_get_events(executor, parameters)
elif tool_name == "k8s_describe_pod":
output = await self._k8s_describe_pod(executor, parameters)
elif tool_name == "k8s_get_hpa_status":
output = await self._k8s_get_hpa_status(executor, parameters)
elif tool_name == "k8s_get_node_conditions":
output = await self._k8s_get_node_conditions(executor, parameters)
else:
return MCPToolResult(
success=False,
execution_id=execution_id,
error=f"Unknown tool: {tool_name}",
)
return MCPToolResult(
success=True,
execution_id=execution_id,
output=output,
)
except Exception as e:
logger.exception("k8s_provider_error", tool=tool_name, error=str(e))
return MCPToolResult(
success=False,
execution_id=execution_id,
error=str(e),
)
async def _kubectl_get(self, executor, parameters: dict) -> dict:
namespace = parameters.get("namespace", "awoooi-prod")
resource = parameters.get("resource", "pods")
name = parameters.get("name", "")
cmd = f"kubectl get {resource} {name} -n {namespace} -o json".strip()
result = await executor.execute_kubectl_command(cmd)
if result.success and result.k8s_response:
return result.k8s_response.get("stdout", "")
return {"error": result.error}
async def _kubectl_delete(self, executor, parameters: dict) -> dict:
namespace = parameters.get("namespace", "awoooi-prod")
resource = parameters.get("resource", "pod")
name = parameters.get("name", "")
if not name:
return {"error": "Missing 'name' parameter"}
# Dry-run validation
if resource == "pod":
dry_run = await executor.validate_pod_exists(name, namespace)
else:
dry_run = await executor.validate_deployment_exists(name, namespace)
if not dry_run.passed:
return {"error": dry_run.message, "dry_run": False}
# Execute deletion
if resource == "pod":
result = await executor.delete_pod(name, namespace)
else:
return {"error": "Direct deployment deletion not supported, use restart"}
return {
"success": result.success,
"message": result.message,
"duration_ms": result.duration_ms,
}
async def _kubectl_scale(self, executor, parameters: dict) -> dict:
namespace = parameters.get("namespace", "awoooi-prod")
deployment = parameters.get("deployment", "")
replicas = parameters.get("replicas", 1)
if not deployment:
return {"error": "Missing 'deployment' parameter"}
cmd = f"kubectl scale deployment/{deployment} --replicas={replicas} -n {namespace}"
result = await executor.execute_kubectl_command(cmd)
return {
"success": result.success,
"scaled": result.success,
"replicas": replicas,
"message": result.message,
}
async def _kubectl_restart(self, executor, parameters: dict) -> dict:
namespace = parameters.get("namespace", "awoooi-prod")
deployment = parameters.get("deployment", "")
if not deployment:
return {"error": "Missing 'deployment' parameter"}
dry_run = await executor.validate_deployment_exists(deployment, namespace)
if not dry_run.passed:
return {"error": dry_run.message, "dry_run": False}
result = await executor.restart_deployment(deployment, namespace)
return {
"success": result.success,
"restarted": result.success,
"message": result.message,
"duration_ms": result.duration_ms,
}
# =========================================================================
# MCP Phase 1 工具實作2026-04-11 Claude Sonnet 4.6
# =========================================================================
async def _k8s_get_pod_logs(self, executor, parameters: dict) -> dict:
# Bug 根因LLM 有時傳空字串 pod_name_validate_name 直接 raise ValueError 炸整個 investigator
# 修復pod_name 空時返回結構化 error讓 Agent 可繼續分析其他工具2026-04-24 Claude Sonnet 4.6
raw_pod_name = parameters.get("pod_name", "")
if not raw_pod_name or not raw_pod_name.strip():
return {"error": "empty_pod_name", "message": "pod_name 為空,請提供具體 pod 名稱(如 awoooi-api-xxx-yyy"}
pod_name = _validate_name(raw_pod_name, "pod_name")
namespace = _validate_namespace(parameters.get("namespace", DEFAULT_NAMESPACE))
container = parameters.get("container", "")
tail = max(1, min(int(parameters.get("tail", 100)), 500))
previous = parameters.get("previous", False)
container_flag = f"-c {_validate_name(container, 'container')}" if container else ""
previous_flag = "--previous" if previous else ""
cmd = f"kubectl logs {pod_name} {container_flag} {previous_flag} --tail={tail} -n {namespace}".strip()
result = await executor.execute_kubectl_command(cmd)
return {
"pod": pod_name,
"namespace": namespace,
"tail": tail,
"previous": previous,
"logs": result.k8s_response.get("stdout", "") if result.success and result.k8s_response else "",
"error": result.error if not result.success else None,
}
async def _k8s_watch_rollout(self, executor, parameters: dict) -> dict:
deployment = _validate_name(parameters["deployment"], "deployment")
namespace = _validate_namespace(parameters.get("namespace", DEFAULT_NAMESPACE))
timeout = max(10, min(int(parameters.get("timeout_seconds", 120)), 300))
cmd = f"kubectl rollout status deployment/{deployment} -n {namespace} --timeout={timeout}s"
result = await executor.execute_kubectl_command(cmd)
stdout = result.k8s_response.get("stdout", "") if result.success and result.k8s_response else ""
return {
"deployment": deployment,
"namespace": namespace,
"success": result.success,
"status": stdout.strip() or result.error,
}
async def _k8s_get_events(self, executor, parameters: dict) -> dict:
namespace = _validate_namespace(parameters.get("namespace", DEFAULT_NAMESPACE))
resource_name = parameters.get("resource_name", "")
event_type = parameters.get("event_type", "")
# 白名單校驗可選參數
if resource_name:
resource_name = _validate_name(resource_name, "resource_name")
if event_type and event_type not in ("Warning", "Normal"):
return {"error": "event_type must be 'Warning' or 'Normal'"}
field_selector_parts = []
if resource_name:
field_selector_parts.append(f"involvedObject.name={resource_name}")
if event_type:
field_selector_parts.append(f"type={event_type}")
field_selector = f"--field-selector={','.join(field_selector_parts)}" if field_selector_parts else ""
cmd = f"kubectl get events -n {namespace} {field_selector} --sort-by='.lastTimestamp' -o json".strip()
result = await executor.execute_kubectl_command(cmd)
return {
"namespace": namespace,
"resource_name": resource_name or None,
"event_type": event_type or None,
"events": result.k8s_response.get("stdout", "") if result.success and result.k8s_response else "",
"error": result.error if not result.success else None,
}
async def _k8s_describe_pod(self, executor, parameters: dict) -> dict:
# Bug 根因:同 _k8s_get_pod_logs空 pod_name 會 raise 炸 investigator2026-04-24 Claude Sonnet 4.6
raw_pod_name = parameters.get("pod_name", "")
if not raw_pod_name or not raw_pod_name.strip():
return {"error": "empty_pod_name", "message": "pod_name 為空,請提供具體 pod 名稱(如 awoooi-api-xxx-yyy"}
pod_name = _validate_name(raw_pod_name, "pod_name")
namespace = _validate_namespace(parameters.get("namespace", DEFAULT_NAMESPACE))
cmd = f"kubectl describe pod {pod_name} -n {namespace}"
result = await executor.execute_kubectl_command(cmd)
return {
"pod": pod_name,
"namespace": namespace,
"description": result.k8s_response.get("stdout", "") if result.success and result.k8s_response else "",
"error": result.error if not result.success else None,
}
async def _k8s_get_hpa_status(self, executor, parameters: dict) -> dict:
namespace = _validate_namespace(parameters.get("namespace", DEFAULT_NAMESPACE))
name = parameters.get("name", "")
if name:
name = _validate_name(name, "name")
target = name if name else ""
cmd = f"kubectl get hpa {target} -n {namespace} -o json".strip()
result = await executor.execute_kubectl_command(cmd)
return {
"namespace": namespace,
"hpa": result.k8s_response.get("stdout", "") if result.success and result.k8s_response else "",
"error": result.error if not result.success else None,
}
async def _k8s_get_node_conditions(self, executor, _parameters: dict) -> dict:
# Node 操作不需要 namespace直接查 cluster 層
cmd = "kubectl get nodes -o json"
result = await executor.execute_kubectl_command(cmd)
return {
"nodes": result.k8s_response.get("stdout", "") if result.success and result.k8s_response else "",
"error": result.error if not result.success else None,
}
async def health_check(self) -> bool:
"""Check if kubectl is accessible"""
try:
executor = self._get_executor()
result = await executor.execute_kubectl_command("kubectl version --client")
return result.success
except Exception:
return False