Tribh/devops-copilot
0
1"""2devops_tools.py — DevOps tool registry using configurable thresholds.3"""4from devops_copilot.tools.registry import registry5from devops_copilot.core.log_storage import LogStorage6from devops_copilot.core.config import thresholds7from devops_copilot.core.observability import (8 track_tool_metrics, ACTIVE_INCIDENTS,9 REMEDIATION_SUCCESS, REMEDIATION_FAILURE,10 ANOMALY_DETECTION_TIME11)12from devops_copilot.utils.logger import logger13from typing import Optional, Dict, Any, List14import time15import httpx16from prometheus_client.parser import text_string_to_metric_families17import numpy as np18from collections import deque19 20# Shared async log store instance21log_store = LogStorage()22 23 24@registry.register(name="search_logs", description="Search production logs for a service.")25@track_tool_metrics("search_logs")26def search_logs(service: str, level: Optional[str] = None, minutes_ago: int = 5) -> str:27 window = thresholds.window_seconds(service) // 60 # use per-service window28 start_time = time.time() - (max(minutes_ago, window) * 60)29 logs = log_store.query_logs_sync(service=service, level=level, start_time=start_time)30 if not logs:31 return f"No logs found for {service} in the last {minutes_ago} minutes."32 formatted = "\n".join(f"[{l['level']}] {l['message']}" for l in logs)33 return f"Latest logs for {service}:\n{formatted}"34 35 36@registry.register(name="get_metrics", description="Get error rates and health metrics for a service.")37@track_tool_metrics("get_metrics")38def get_metrics(service: str) -> Dict[str, Any]:39 window = thresholds.window_seconds(service)40 error_rate = log_store.get_error_rate_sync(service, window_seconds=window)41 threshold = thresholds.error_rate_threshold(service)42 43 # Require minimum log volume before triggering anomaly (avoids cold-start noise)44 total_logs = log_store.query_logs_sync(45 service=service,46 start_time=time.time() - window,47 )48 enough_data = len(total_logs) >= thresholds.MIN_LOG_VOLUME49 50 is_anomaly = enough_data and (error_rate > threshold)51 status = "CRITICAL" if is_anomaly else ("HEALTHY" if enough_data else "INSUFFICIENT_DATA")52 53 if is_anomaly:54 ACTIVE_INCIDENTS.inc()55 spike_start = log_store.get_spike_start_sync(service)56 if spike_start:57 mttd = min(time.time() - spike_start, thresholds.MTTD_CEILING_SECONDS)58 ANOMALY_DETECTION_TIME.observe(mttd)59 logger.info(f"MTTD for {service}: {mttd:.2f}s (threshold={threshold*100:.0f}%, window={window}s)")60 61 return {62 "service": service,63 "error_rate": f"{error_rate*100:.2f}%",64 "threshold": f"{threshold*100:.0f}%",65 "window_seconds": window,66 "status": status,67 "anomaly_detected": is_anomaly,68 "log_count": len(total_logs),69 "timestamp": time.time()70 }71 72 73@registry.register(name="restart_service", description="Restart a failing service. REQUIRES APPROVAL.")74@track_tool_metrics("restart_service")75def restart_service(service: str, reason: str) -> str:76 logger.warning(f"RESTARTING SERVICE: {service} | reason: {reason}")77 try:78 REMEDIATION_SUCCESS.labels(service=service).inc()79 ACTIVE_INCIDENTS.dec()80 log_store.clear_spike_sync(service)81 return f"Service {service} successfully restarted. Reason: {reason}"82 except Exception as e:83 REMEDIATION_FAILURE.labels(service=service).inc()84 raise e85 86 87@registry.register(name="slack_notify", description="Send a message to the DevOps Slack channel.")88@track_tool_metrics("slack_notify")89def slack_notify(channel: str, message: str) -> str:90 logger.info(f"SLACK [{channel}]: {message}")91 return f"Notification sent to #{channel}"92 93 94# Metric history for anomaly detection95metric_history: Dict[str, deque] = {}96 97@registry.register(name="scrape_metrics", description="Fetch live Prometheus metrics from an external service and check for anomalies.")98@track_tool_metrics("scrape_metrics")99def scrape_metrics(service: str) -> str:100 url = thresholds.get_service_url(service)101 if not url:102 return f"Error: No metrics URL configured for service '{service}'."103 104 try:105 logger.info(f"Scraping metrics for {service} from {url}...")106 response = httpx.get(url, timeout=5.0)107 response.raise_for_status()108 109 metrics_text = response.text110 parsed_metrics = []111 anomaly_flag = False112 113 # Look for DocuGenie specific metrics or general health114 for family in text_string_to_metric_families(metrics_text):115 if family.name.startswith("docugenie_") or family.name.startswith("fastapi_"):116 for sample in family.samples:117 val = sample.value118 parsed_metrics.append(f"{sample.name}: {val}")119 120 # Statistical Anomaly Detection121 history_key = f"{service}:{sample.name}"122 if history_key not in metric_history:123 metric_history[history_key] = deque(maxlen=20)124 125 hist = metric_history[history_key]126 if len(hist) > 5:127 mean = np.mean(hist)128 std = np.std(hist)129 if std > 0:130 z_score = (val - mean) / std131 if abs(z_score) > 3:132 anomaly_flag = True133 logger.warning(f"🚨 ANOMALY detected in {sample.name}: Z-Score={z_score:.2f}")134 135 hist.append(val)136 137 if not parsed_metrics:138 return f"Service {service} is reachable, but no specific application metrics were found."139 140 status_prefix = "⚠️ ANOMALY DETECTED! " if anomaly_flag else ""141 summary = "\n".join(parsed_metrics[:10]) 142 return f"{status_prefix}Live metrics for {service}:\n{summary}"143 144 except Exception as e:145 logger.error(f"Failed to scrape metrics for {service}: {e}")146 return f"Error: Could not reach metrics endpoint for {service}. {str(e)}"147 148 149@registry.register(name="trigger_chaos_simulation", description="Load a failure scenario into the system for demonstration.")150@track_tool_metrics("trigger_chaos_simulation")151def trigger_chaos_simulation(scenario: str) -> str:152 """153 Scenarios: 'latency_spike', 'rate_limit', 'memory_leak'154 """155 scenarios = {156 "latency_spike": {157 "service": "docugenie",158 "metrics": ["docugenie_retrieval_latency_seconds_sum: 45.2", "docugenie_queries_total: 100"],159 "logs": "WARNING: Vector search query timed out after 2.0s\nERROR: Upstream latency spike detected in FAISS."160 },161 "rate_limit": {162 "service": "docugenie",163 "metrics": ["fastapi_responses_total{status='429'}: 15"],164 "logs": "ERROR: Azure OpenAI Rate Limit reached (code 429). Retrying in 60s..."165 }166 }167 168 if scenario not in scenarios:169 return f"Unknown scenario '{scenario}'. Available: {list(scenarios.keys())}"170 171 data = scenarios[scenario]172 logger.warning(f"CHAOS SIMULATION TRIGGERED: {scenario}")173 174 # Ingest logs into the local log store for the agent to find175 for line in data["logs"].split("\n"):176 log_store.ingest_log(data["service"], "ERROR" if "ERROR" in line else "WARNING", line)177 178 return f"Chaos scenario '{scenario}' successfully loaded into {data['service']}. AI SRE has been notified."179 