"""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()