26S26h

QK·6일 전
#!/usr/bin/env python3
"""Robusta-style Synthetic Incident Context Generator for RCA Agent."""
import argparse
import json
import os
import random
from datetime import datetime, timedelta, timezone
from typing import Any, Dict, List, Tuple

DEFAULT_CLUSTER_NAME = os.getenv("RCA_CLUSTER_NAME", "prod-kr01-k8s")
DEFAULT_NAMESPACE = os.getenv("RCA_NAMESPACE", "storage")
DEFAULT_TARGET_NODE = os.getenv("RCA_TARGET_NODE", "kr01-aistor-dn03")
DEFAULT_TARGET_POD = os.getenv("RCA_TARGET_POD", "minio-aistor-pool-0-2")
DEFAULT_MOUNT_DRIVE = os.getenv("RCA_MOUNT_DRIVE", "/data/nvme02")
DEFAULT_SERVICE_TIER = os.getenv("RCA_SERVICE_TIER", "P1-CRITICAL")


def generate_time_series(base_val: float, spike_val: float, anomaly_idx: int,
                         steps: int = 15, noise: float = 0.05) -> List[Dict[str, Any]]:
    """Prometheus Range Query 응답 형식([timestamp, value])을 모사합니다."""
    now = datetime.now(timezone.utc)
    series: List[Dict[str, Any]] = []
    for i in range(steps):
        t = now - timedelta(minutes=(steps - i))
        ts_sec = int(t.timestamp())
        if i < anomaly_idx:
            v = base_val * (1.0 + random.uniform(-noise, noise))
        else:
            ramp = min(1.0, 0.4 + (i - anomaly_idx) * 0.2)
            v = base_val + (spike_val - base_val) * ramp * (1.0 + random.uniform(-noise, noise))
        series.append({"timestamp": ts_sec, "value": f"{v:.2f}"})
    return series


def build_scenario_slow_drive(node: str, pod: str, drive: str) -> Tuple[str, str, Dict[str, Any], List[str], List[Dict[str, Any]]]:
    """특정 NVMe I/O Hang / Slow Drive로 인한 분산 락 및 타임아웃 장애."""
    alarm_name = "MinIODriveLatencyHigh"
    symptom_summary = f"Client S3 PutObject timeout (HTTP 503 Slow Down) and I/O wait spike on drive {drive} on {node}"

    metrics = {
        "minio_node_drive_total_duration_microseconds": {
            "query": f'minio_node_drive_total_duration_microseconds{{pod="{pod}", drive="{drive}"}}',
            "baseline": "15,000 µs (15ms)",
            "values": generate_time_series(15000, 4800000, 8),
        },
        "go_goroutines": {
            "query": f'go_goroutines{{pod="{pod}"}}',
            "baseline": "2,200",
            "values": generate_time_series(2200, 18900, 8),
        },
        "minio_s3_requests_waiting_total": {
            "query": f'minio_s3_requests_waiting_total{{pod="{pod}"}}',
            "baseline": "12",
            "values": generate_time_series(12, 1450, 9),
        },
        "node_disk_io_time_seconds_total": {
            "query": f'rate(node_disk_io_time_seconds_total{{instance="{node}"}}[1m])',
            "baseline": "0.15",
            "values": generate_time_series(0.15, 0.99, 8),
        },
    }

    base_t = datetime.now(timezone.utc) - timedelta(minutes=5)
    logs = [
        f"{base_t.strftime('%H:%M:%S')} [INFO] (api) S3 API request received: PutObject bucket=corp-ai-training-data size=12400192B",
        f"{(base_t + timedelta(seconds=15)).strftime('%H:%M:%S')} [WARN] (storage) Drive {drive} on {node} responded slowly: 4.821s latency threshold exceeded",
        f"{(base_t + timedelta(seconds=22)).strftime('%H:%M:%S')} [ERROR] (dsync) dsync: timeout while acquiring write lock on peer {node}:9000 for path corp-ai-training-data/weights_v2.bin",
        f"{(base_t + timedelta(seconds=23)).strftime('%H:%M:%S')} [ERROR] (api) PutObject(bucket=corp-ai-training-data) err: context deadline exceeded while waiting on {drive}",
        f"{(base_t + timedelta(seconds=28)).strftime('%H:%M:%S')} [WARN] (http) Returned HTTP 503 Slow Down to client ai-training-worker-6789f",
        f"{(base_t + timedelta(seconds=35)).strftime('%H:%M:%S')} [WARN] (storage) Drive {drive} marked as degraded due to multiple I/O timeouts",
    ]

    events = [
        {
            "type": "Warning",
            "reason": "StorageLatencySpike",
            "message": f"Persistent volume attached to {pod} via {drive} exceeded SLO response time (>3000ms)",
            "source": f"directpv-driver-{node}",
            "count": 14,
        },
        {
            "type": "Warning",
            "reason": "Unhealthy",
            "message": f"Readiness probe failed on {pod}: HTTP probe failed with statuscode: 503",
            "source": "kubelet",
            "count": 3,
        },
    ]
    return alarm_name, symptom_summary, metrics, logs, events


def build_scenario_network_partition(node: str, pod: str, drive: str) -> Tuple[str, str, Dict[str, Any], List[str], List[Dict[str, Any]]]:
    """노드 간 네트워크 순단/패킷 드롭으로 인한 dsync 분산 락 실패."""
    alarm_name = "MinIONodePeerUnreachable"
    symptom_summary = f"Peer interconnect dropped between {node} and other cluster nodes, causing dsync quorum loss"

    metrics = {
        "minio_node_interconnect_rx_errors_total": {
            "query": f'rate(minio_node_interconnect_rx_errors_total{{instance="{node}"}}[1m])',
            "baseline": "0.0",
            "values": generate_time_series(0.0, 420.0, 7),
        },
        "minio_cluster_nodes_offline_total": {
            "query": "minio_cluster_nodes_offline_total",
            "baseline": "0",
            "values": generate_time_series(0, 1, 8),
        },
    }

    logs = [
        f"[ERROR] (peer-rpc) Failed to dial peer node {node}:9000: dial tcp {node}:9000: i/o timeout",
        "[ERROR] (dsync) Lost quorum on locker dsync for distributed pool: expected 8 votes, received 4",
        f"[WARN] (cluster) Node {node} temporarily marked as offline in peer table",
    ]

    events = [{
        "type": "Warning",
        "reason": "NodeNetworkFlapping",
        "message": f"Node {node} NIC dropped BGP routes / interface packet drop spike",
        "source": "cilium-agent",
        "count": 8,
    }]
    return alarm_name, symptom_summary, metrics, logs, events


def generate_incident_payload(scenario: str, output_file: str, cluster: str,
                              namespace: str, node: str, pod: str, drive: str) -> str:
    """Robusta-style Incident Bundle JSON을 생성합니다."""
    now = datetime.now(timezone.utc)
    incident_id = f"INC-{cluster.split('-')[0].upper()}-{now.strftime('%Y%m%d%H%M%S')}"

    if scenario == "slow_drive":
        alarm_name, symptom_summary, metrics, logs, events = build_scenario_slow_drive(node, pod, drive)
    else:
        alarm_name, symptom_summary, metrics, logs, events = build_scenario_network_partition(node, pod, drive)

    payload = {
        "schema_version": "v1.2",
        "collector_source": "robusta-headless-standalone",
        "incident_meta": {
            "incident_id": incident_id,
            "alarm_name": alarm_name,
            "severity": DEFAULT_SERVICE_TIER,
            "timestamp": now.isoformat(),
            "cluster_name": cluster,
            "target": {
                "namespace": namespace,
                "node_name": node,
                "pod_name": pod,
                "mount_path": drive,
            },
        },
        "symptom_summary": symptom_summary,
        "evidence": {
            "k8s_events": events,
            "pod_logs_tail": logs,
            "prometheus_metrics": metrics,
        },
        "system_context": {
            "component": "minio-aistor",
            "storage_driver": "directpv",
            "k8s_version": "v1.28.10",
        },
    }

    with open(output_file, "w", encoding="utf-8") as f:
        json.dump(payload, f, ensure_ascii=False, indent=2)

    return os.path.abspath(output_file)


def main() -> None:
    parser = argparse.ArgumentParser(description="RCA Agent 검증용 가상 Incident Context JSON 제너레이터")
    parser.add_argument("--scenario", choices=["slow_drive", "network_partition"], default="slow_drive", help="장애 시나리오 선택")
    parser.add_argument("--output", default="synthetic_incident_payload.json", help="저장할 JSON 파일명")
    parser.add_argument("--cluster", default=DEFAULT_CLUSTER_NAME, help="클러스터명")
    parser.add_argument("--namespace", default=DEFAULT_NAMESPACE, help="타깃 네임스페이스")
    parser.add_argument("--node", default=DEFAULT_TARGET_NODE, help="타깃 노드 호스트명")
    parser.add_argument("--pod", default=DEFAULT_TARGET_POD, help="타깃 MinIO 파드명")
    parser.add_argument("--drive", default=DEFAULT_MOUNT_DRIVE, help="타깃 드라이브 마운트 경로")
    args = parser.parse_args()

    out_path = generate_incident_payload(
        scenario=args.scenario,
        output_file=args.output,
        cluster=args.cluster,
        namespace=args.namespace,
        node=args.node,
        pod=args.pod,
        drive=args.drive,
    )

    print("=" * 65)
    print("[✓] 가상 Incident Context JSON 생성 완료!")
    print(f"[*] 파일 경로       : {out_path}")
    print(f"[*] 장애 시나리오   : {args.scenario}")
    print(f"[*] 대상 인프라     : Node={args.node}, Pod={args.pod}, Drive={args.drive}")
    print("=" * 65)


if __name__ == "__main__":
    main()
profile
engineer

0개의 댓글