| import time |
| from datetime import datetime |
| from agent_brain import NagaMLOpsAgent |
| from pipeline_engine import MLPipelineEngine |
| import utils |
|
|
| class SelfHealingOrchestrator: |
| def __init__(self): |
| self.agent = NagaMLOpsAgent() |
| self.engine = MLPipelineEngine() |
| self.total_runs = 0 |
| self.faults_detected = 0 |
| self.faults_healed = 0 |
| self.total_heal_time_sec = 0.0 |
|
|
| def get_kpi_stats(self) -> dict: |
| """ |
| Calculates high level system health KPIs for dashboard counters. |
| """ |
| success_rate = 100.0 |
| if self.total_runs > 0: |
| failed_unhealed = self.faults_detected - self.faults_healed |
| success_rate = round(max(0.0, ((self.total_runs - failed_unhealed) / self.total_runs) * 100.0), 1) |
|
|
| avg_mtth = 0.0 |
| if self.faults_healed > 0: |
| avg_mtth = round(self.total_heal_time_sec / self.faults_healed, 2) |
|
|
| return { |
| "health_score_pct": success_rate, |
| "total_incidents": self.faults_detected, |
| "auto_healed_count": self.faults_healed, |
| "avg_mtth_sec": avg_mtth |
| } |
|
|
| def run_pipeline_check(self, custom_code: str = None) -> dict: |
| """ |
| Executes standard pipeline monitoring run without fault injection. |
| """ |
| self.total_runs += 1 |
| exec_result = self.engine.execute_pipeline(code=custom_code) |
| return { |
| "status": exec_result["telemetry"]["status"], |
| "telemetry": exec_result["telemetry"], |
| "logs": exec_result["logs"], |
| "code": exec_result["code"], |
| "kpis": self.get_kpi_stats() |
| } |
|
|
| def trigger_and_heal_fault(self, fault_name: str, custom_code: str = None) -> dict: |
| """ |
| Full autonomous closed-loop: |
| 1. Inject fault / load scenario |
| 2. Run broken pipeline & detect failure |
| 3. Invoke Naga Agentic LLM for Root Cause Analysis & Python code patch synthesis |
| 4. Apply patch code to engine |
| 5. Verify re-execution & commit incident log |
| """ |
| start_heal_timer = time.time() |
| self.total_runs += 1 |
| self.faults_detected += 1 |
|
|
| |
| if custom_code: |
| broken_code = self.engine.set_custom_code(custom_code) |
| elif fault_name: |
| broken_code = self.engine.load_fault_scenario(fault_name) |
| else: |
| broken_code = self.engine.current_code |
|
|
| |
| pre_heal_run = self.engine.execute_pipeline(broken_code) |
| pre_telemetry = pre_heal_run["telemetry"] |
| pre_logs = pre_heal_run["logs"] |
|
|
| agent_timeline = [] |
| agent_timeline.append(f"⏱ [{time.strftime('%H:%M:%S')}] 🚨 ANOMALY DETECTED! Status: {pre_telemetry['status']}. Initiating Naga AI Agentic Diagnosis...") |
|
|
| |
| diagnosis = self.agent.diagnose_and_heal( |
| script_code=broken_code, |
| execution_logs=pre_logs, |
| telemetry=pre_telemetry, |
| fault_name=fault_name |
| ) |
|
|
| fault_cat = diagnosis.get("fault_category", "UNKNOWN_FAULT") |
| rca = diagnosis.get("root_cause_analysis", "No detailed RCA provided.") |
| explanation = diagnosis.get("explanation_for_engineers", "Patch synthesized.") |
| patched_code = diagnosis.get("patch_code", broken_code) |
|
|
| agent_timeline.append(f"⏱ [{time.strftime('%H:%M:%S')}] 🧠 Root Cause Analysis ({fault_cat}): {rca[:180]}...") |
| agent_timeline.append(f"⏱ [{time.strftime('%H:%M:%S')}] 🛠 Synthesized Python Patch Code. Applying patch to execution environment...") |
|
|
| |
| post_heal_run = self.engine.execute_pipeline(patched_code) |
| post_telemetry = post_heal_run["telemetry"] |
| post_logs = post_heal_run["logs"] |
|
|
| heal_duration = round(time.time() - start_heal_timer, 2) |
| verified_success = (post_telemetry["status"] in ["HEALTHY", "NORMAL"]) and (post_telemetry["accuracy"] >= 0.70 or "syntax" not in post_logs.lower()) |
|
|
| if verified_success: |
| self.faults_healed += 1 |
| self.total_heal_time_sec += heal_duration |
| agent_timeline.append(f"⏱ [{time.strftime('%H:%M:%S')}] ✅ VERIFICATION SUCCESSFUL! Pipeline restored to HEALTHY (Accuracy: {post_telemetry['accuracy']*100:.1f}%, Time to Heal: {heal_duration}s).") |
| final_status = "RESOLVED_AND_HEALTHY" |
| else: |
| agent_timeline.append(f"⏱ [{time.strftime('%H:%M:%S')}] ⚠️ Verification Partial. Pipeline output logged for engineer review.") |
| final_status = "HEAL_ATTEMPTED_NEEDS_REVIEW" |
|
|
| |
| diff_html = utils.generate_code_diff_html(broken_code, patched_code) |
|
|
| |
| incident_id = f"INC-{datetime.now().strftime('%Y%m%d-%H%M%S')}" |
| incident_data = { |
| "id": incident_id, |
| "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), |
| "fault_name": fault_name or "Custom Script Anomaly", |
| "fault_category": fault_cat, |
| "final_status": final_status, |
| "heal_duration_sec": heal_duration, |
| "pre_telemetry": pre_telemetry, |
| "post_telemetry": post_telemetry, |
| "root_cause_analysis": rca, |
| "explanation": explanation, |
| "verification_checklist": diagnosis.get("verification_checklist", []), |
| "original_code": broken_code, |
| "patched_code": patched_code, |
| "pre_logs": pre_logs, |
| "post_logs": post_logs |
| } |
| |
| report_file = utils.save_incident_report(incident_data) |
| agent_timeline.append(f"⏱ [{time.strftime('%H:%M:%S')}] 📝 Saved Incident Audit Report: {report_file}") |
|
|
| return { |
| "incident_id": incident_id, |
| "status": final_status, |
| "fault_category": fault_cat, |
| "rca": rca, |
| "explanation": explanation, |
| "heal_duration": heal_duration, |
| "pre_telemetry": pre_telemetry, |
| "post_telemetry": post_telemetry, |
| "timeline": "\n".join(agent_timeline), |
| "diff_html": diff_html, |
| "broken_code": broken_code, |
| "patched_code": patched_code, |
| "pre_logs": pre_logs, |
| "post_logs": post_logs, |
| "kpis": self.get_kpi_stats(), |
| "telemetry_history": self.engine.execution_history |
| } |
|
|