| 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
|
| }
|
|
|