Team Ai
Apppublic

Tribh/devops-copilot

sourceHugging Faceupdated 8mo agoView on Hugging Face
0likes
devops_tools.py179 linesDownload Raw Back to tools
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