Spaces:
Runtime error
Runtime error
| """ | |
| Risk service – integrates ARF Bayesian risk engine, policy engine, and decision engine. | |
| Deterministic, no random fallbacks, explicit error handling. Tenant‑aware. | |
| Version: 2026-07-06 – added evaluate_intent_full with GovernanceLoop integration, | |
| skill context injection, and full HealingIntent serialisation. | |
| """ | |
| import json | |
| import logging | |
| import os | |
| import time | |
| from typing import Optional, List, Dict, Any | |
| from agentic_reliability_framework.core.governance.risk_engine import RiskEngine | |
| from agentic_reliability_framework.core.governance.intents import InfrastructureIntent | |
| from agentic_reliability_framework.core.models.event import ReliabilityEvent, HealingAction | |
| from agentic_reliability_framework.core.governance.policy_engine import PolicyEngine | |
| from agentic_reliability_framework.core.decision.decision_engine import DecisionEngine | |
| from agentic_reliability_framework.runtime.memory.rag_graph import RAGGraphMemory | |
| from agentic_reliability_framework.core.research.eclipse_probe import compute_epistemic_risk | |
| # ── Governance loop integration ────────────────────────────── | |
| from agentic_reliability_framework.core.governance.governance_loop import GovernanceLoop | |
| from agentic_reliability_framework.core.governance.cost_estimator import CostEstimator | |
| from agentic_reliability_framework.core.governance.policies import PolicyEvaluator, allow_all | |
| from agentic_reliability_framework.core.governance.stability_controller import LyapunovStabilityController | |
| from agentic_reliability_framework.core.temporal_reliability import TemporalReliabilityMonitor | |
| from agentic_reliability_framework.core.governance.healing_intent import HealingIntent | |
| # ── optional tracing ───────────────────────────────────────── | |
| try: | |
| from opentelemetry import trace | |
| _tracer = trace.get_tracer(__name__) | |
| OTEL_AVAILABLE = True | |
| except ImportError: | |
| OTEL_AVAILABLE = False | |
| _tracer = None | |
| # ── Prometheus metrics (always registered; no‑op if not scraped) ─ | |
| from prometheus_client import Counter, Histogram | |
| _EVAL_COUNTER = Counter( | |
| "arf_evaluations_total", | |
| "Total evaluation calls (intent + healing), partitioned by engine and status.", | |
| ["engine", "status"], | |
| ) | |
| _EVAL_DURATION = Histogram( | |
| "arf_evaluation_duration_seconds", | |
| "End‑to‑end latency of evaluation calls.", | |
| ["engine"], | |
| buckets=(0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0), | |
| ) | |
| _RUST_AGREEMENT = Counter( | |
| "arf_rust_agreement_total", | |
| "Agreement between Rust enforcer and Python policy evaluation.", | |
| ["result"], # "agreed" or "diverged" | |
| ) | |
| # ── optional Rust enforcer (shadow mode) ────────────────────── | |
| _RUST_ENFORCER_AVAILABLE = False | |
| _rust_evaluator = None # singleton per process | |
| _rust_policy_json: Optional[str] = None | |
| if os.getenv("ARF_USE_RUST_ENFORCER", "false").lower() == "true": | |
| try: | |
| import arf_enforcer | |
| _RUST_ENFORCER_AVAILABLE = True | |
| except ImportError: | |
| pass | |
| # Default OSS policy tree – mirrors the hard‑coded rules in the Python PolicyEvaluator | |
| _OSS_POLICY_TREE_JSON = json.dumps({ | |
| "And": [ | |
| {"Atomic": {"RegionAllowed": {"allowed_regions": ["eastus"]}}}, | |
| {"Atomic": {"ResourceTypeRestricted": { | |
| "forbidden_types": ["DATABASE_DROP", "FULL_ROLLOUT", "SYSTEM_SHUTDOWN", "SECRET_ROTATION"] | |
| }}}, | |
| {"Atomic": {"MaxPermissionLevel": {"max_level": "admin"}}} | |
| ] | |
| }) | |
| def _ensure_rust_evaluator() -> bool: | |
| """Lazy initialise the Rust policy evaluator. Returns True on success.""" | |
| global _rust_evaluator, _rust_policy_json | |
| if _rust_evaluator is not None: | |
| return True | |
| if not _RUST_ENFORCER_AVAILABLE: | |
| return False | |
| try: | |
| _rust_policy_json = _OSS_POLICY_TREE_JSON | |
| _rust_evaluator = arf_enforcer.PyPolicyEvaluator(_rust_policy_json) | |
| return True | |
| except Exception: | |
| _rust_evaluator = None | |
| return False | |
| logger = logging.getLogger(__name__) | |
| def evaluate_intent( | |
| engine: RiskEngine, | |
| intent: InfrastructureIntent, | |
| cost_estimate: Optional[float], | |
| policy_violations: List[str], | |
| tenant_id: Optional[str] = None, | |
| ) -> dict: | |
| """ | |
| Evaluate an infrastructure intent using the Bayesian risk engine. | |
| The risk score is computed using a weighted fusion of conjugate online | |
| model, optional hyperpriors, and offline HMC. The tenant_id is passed | |
| to the risk engine to select the correct per‑tenant Beta store. | |
| Parameters | |
| ---------- | |
| engine : RiskEngine | |
| Initialised ARF Bayesian risk engine (must be tenant‑aware). | |
| intent : InfrastructureIntent | |
| The infrastructure request to evaluate. | |
| cost_estimate : float or None | |
| Estimated monthly cost (used by cost‑threshold policies). | |
| policy_violations : list[str] | |
| Pre‑computed policy violation strings (from the Python evaluator). | |
| tenant_id : str, optional | |
| Tenant UUID. If provided, the risk engine will use tenant‑specific | |
| conjugate state. Required for multi‑tenant deployments. | |
| Returns | |
| ------- | |
| dict | |
| Keys: risk_score, explanation, contributions. | |
| """ | |
| t0 = time.monotonic() | |
| span = None | |
| if OTEL_AVAILABLE and _tracer: | |
| span = _tracer.start_span("risk_service.evaluate_intent") | |
| span.set_attribute("intent_type", type(intent).__name__) | |
| if tenant_id: | |
| span.set_attribute("tenant_id", tenant_id) | |
| # ── Shadow Rust enforcer (best‑effort, non‑blocking) ────── | |
| if _RUST_ENFORCER_AVAILABLE and _ensure_rust_evaluator(): | |
| try: | |
| rust_intent = { | |
| "action": getattr(intent, "intent_type", "unknown"), | |
| "component": getattr(intent, "service_name", "unknown"), | |
| "region": getattr(intent, "region", None), | |
| "resource_type": getattr(intent, "resource_type", None), | |
| "permission_level": getattr(intent, "permission_level", None), | |
| "tenant_id": tenant_id, | |
| "extra": {} | |
| } | |
| rust_raw = _rust_evaluator.evaluate( | |
| json.dumps(rust_intent), cost_estimate | |
| ) | |
| rust_violations = json.loads(rust_raw) | |
| agreed = set(rust_violations) == set(policy_violations) | |
| _RUST_AGREEMENT.labels(result="agreed" if agreed else "diverged").inc() | |
| if not agreed: | |
| msg = ( | |
| f"Rust enforcer divergence for tenant {tenant_id}: " | |
| f"Rust={sorted(rust_violations)} Python={sorted(policy_violations)}" | |
| ) | |
| logger.warning(msg) | |
| if span: | |
| span.add_event("rust_enforcer_divergence", { | |
| "rust_violations": rust_violations, | |
| "python_violations": policy_violations | |
| }) | |
| except Exception as exc: | |
| logger.debug("Rust enforcer shadow evaluation failed: %s", exc) | |
| # ── Core risk evaluation ────────────────────────────────── | |
| try: | |
| if hasattr(engine, "set_tenant"): | |
| engine.set_tenant(tenant_id) | |
| elif tenant_id: | |
| logger.warning( | |
| "RiskEngine does not yet support tenant_id; evaluations will be shared across tenants." | |
| ) | |
| score, explanation, contributions = engine.calculate_risk( | |
| intent=intent, | |
| cost_estimate=cost_estimate, | |
| policy_violations=policy_violations | |
| ) | |
| engine_label = "python" | |
| status = "success" | |
| except Exception: | |
| _EVAL_COUNTER.labels(engine="python", status="error").inc() | |
| _EVAL_DURATION.labels(engine="python").observe(time.monotonic() - t0) | |
| raise | |
| _EVAL_COUNTER.labels(engine=engine_label, status=status).inc() | |
| _EVAL_DURATION.labels(engine=engine_label).observe(time.monotonic() - t0) | |
| if span: | |
| span.set_attribute("risk_score", score) | |
| if _RUST_ENFORCER_AVAILABLE: | |
| span.set_attribute("rust_enforcer_available", True) | |
| span.end() | |
| return { | |
| "risk_score": score, | |
| "explanation": explanation, | |
| "contributions": contributions | |
| } | |
| def evaluate_intent_full( | |
| intent: InfrastructureIntent, | |
| *, | |
| risk_engine: RiskEngine, | |
| cost_estimator: Optional[CostEstimator] = None, | |
| policy_evaluator: Optional[PolicyEvaluator] = None, | |
| memory: Optional[RAGGraphMemory] = None, | |
| enable_epistemic: bool = False, | |
| hallucination_probe: Optional[Any] = None, | |
| predictive_engine: Optional[Any] = None, | |
| business_calculator: Optional[Any] = None, | |
| use_rust_enforcer: bool = False, | |
| stability_controller: Optional[LyapunovStabilityController] = None, | |
| temporal_monitor: Optional[TemporalReliabilityMonitor] = None, | |
| tenant_id: Optional[str] = None, | |
| skill_id: Optional[str] = None, | |
| skill_registry: Optional[Any] = None, | |
| context_extra: Optional[Dict[str, Any]] = None, | |
| ) -> Dict[str, Any]: | |
| """ | |
| Run the full governance loop and return a structured response containing | |
| the serialised HealingIntent with Bayesian skill posterior parameters. | |
| If stability_controller or temporal_monitor are None (the default), | |
| the governance loop will simply skip those checks. Pass stateful | |
| instances from the app state to accumulate cross‑request state. | |
| Parameters | |
| ---------- | |
| intent : InfrastructureIntent | |
| The original infrastructure request. | |
| risk_engine : RiskEngine | |
| Bayesian risk engine (tenant‑aware). | |
| cost_estimator : CostEstimator, optional | |
| Monthly cost estimator; a default instance is created if None. | |
| policy_evaluator : PolicyEvaluator, optional | |
| Policy tree evaluator; defaults to `allow_all` if None. | |
| memory : RAGGraphMemory, optional | |
| Semantic memory for similar‑incident retrieval. | |
| enable_epistemic : bool | |
| Whether to run the ECLIPSE hallucination probe and CUDL attribution. | |
| hallucination_probe : HallucinationRisk, optional | |
| Pre‑configured probe instance. | |
| predictive_engine : SimplePredictiveEngine, optional | |
| Time‑series forecasting engine. | |
| business_calculator : BusinessImpactCalculator, optional | |
| Revenue impact estimator. | |
| use_rust_enforcer : bool | |
| Whether to run the Rust policy evaluator in shadow mode. | |
| stability_controller : LyapunovStabilityController, optional | |
| Passive stability monitor; if None, stability checks are skipped. | |
| temporal_monitor : TemporalReliabilityMonitor, optional | |
| Drift detector; if None, drift detection is skipped. | |
| tenant_id : str, optional | |
| Tenant UUID for multi‑tenant state. | |
| skill_id : str, optional | |
| Skill identifier; if provided, the skill's current posterior | |
| parameters are injected into the governance loop. | |
| skill_registry : SkillRegistry, optional | |
| Instance of the skill registry (required if skill_id is given). | |
| context_extra : dict, optional | |
| Additional key‑value pairs to merge into the loop context. | |
| Returns | |
| ------- | |
| dict | |
| Keys: | |
| - risk_score : float | |
| - explanation : str | |
| - contributions : dict (empty; full trace is in healing_intent) | |
| - healing_intent : dict (serialised HealingIntent) | |
| - recommended_action : str | |
| - deterministic_id : str | |
| """ | |
| t0 = time.monotonic() | |
| span = None | |
| if OTEL_AVAILABLE and _tracer: | |
| span = _tracer.start_span("risk_service.evaluate_intent_full") | |
| span.set_attribute("intent_type", type(intent).__name__) | |
| if tenant_id: | |
| span.set_attribute("tenant_id", tenant_id) | |
| # Default components if not provided | |
| if policy_evaluator is None: | |
| policy_evaluator = PolicyEvaluator(allow_all()) | |
| if cost_estimator is None: | |
| cost_estimator = CostEstimator() | |
| # stability_controller and temporal_monitor are NOT defaulted here; | |
| # they remain None unless explicitly passed. The GovernanceLoop will skip | |
| # those checks gracefully. | |
| loop = GovernanceLoop( | |
| policy_evaluator=policy_evaluator, | |
| cost_estimator=cost_estimator, | |
| risk_engine=risk_engine, | |
| memory=memory, | |
| enable_epistemic=enable_epistemic, | |
| hallucination_probe=hallucination_probe, | |
| predictive_engine=predictive_engine, | |
| business_calculator=business_calculator, | |
| use_rust_enforcer=use_rust_enforcer, | |
| stability_controller=stability_controller, | |
| temporal_monitor=temporal_monitor, | |
| ) | |
| # ── Build context with skill posterior parameters ───────── | |
| context: Dict[str, Any] = dict(context_extra) if context_extra else {} | |
| if skill_id and skill_registry is not None: | |
| try: | |
| # Fetch the latest version for the skill | |
| versions = skill_registry.list_skill_versions(skill_id) | |
| version = versions[-1] if versions else 1 | |
| # Use public get_model() instead of direct _models access | |
| model = skill_registry.get_model(skill_id, version) | |
| if model is not None: | |
| alpha = model.alpha | |
| beta = model.beta | |
| reliability = model.mean() | |
| else: | |
| # Use default prior if no model exists yet | |
| alpha = skill_registry.default_prior_alpha | |
| beta = skill_registry.default_prior_beta | |
| reliability = alpha / (alpha + beta) | |
| context.update({ | |
| "skill_id": skill_id, | |
| "skill_version": version, | |
| "skill_ate": skill_registry.get_ate(skill_id, version), | |
| "skill_reliability_score": reliability, | |
| "skill_alpha": alpha, | |
| "skill_beta": beta, | |
| }) | |
| except Exception as e: | |
| logger.warning("Failed to inject skill context for '%s': %s", skill_id, e) | |
| # ── Execute governance loop ─────────────────────────────── | |
| healing_intent: HealingIntent = loop.run(intent, context=context) | |
| healing_dict = healing_intent.to_dict(include_advisory_context=True) | |
| risk_score = healing_intent.risk_score or 0.0 | |
| explanation = healing_intent.justification or "" | |
| # ── Metrics & span finalisation ─────────────────────────── | |
| _EVAL_COUNTER.labels(engine="governance_loop", status="success").inc() | |
| _EVAL_DURATION.labels(engine="governance_loop").observe(time.monotonic() - t0) | |
| if span: | |
| span.set_attribute("risk_score", risk_score) | |
| span.set_attribute("recommended_action", healing_dict.get("recommended_action")) | |
| span.end() | |
| return { | |
| "risk_score": risk_score, | |
| "explanation": explanation, | |
| "contributions": {}, # full trace is in healing_intent | |
| "healing_intent": healing_dict, | |
| "recommended_action": healing_dict.get("recommended_action"), | |
| "deterministic_id": healing_intent.deterministic_id, | |
| } | |
| def evaluate_healing_decision( | |
| event: ReliabilityEvent, | |
| policy_engine: PolicyEngine, | |
| decision_engine: Optional[DecisionEngine] = None, | |
| rag_graph: Optional[RAGGraphMemory] = None, | |
| model=None, | |
| tokenizer=None, | |
| tenant_id: Optional[str] = None, | |
| ) -> Dict[str, Any]: | |
| """(unchanged from original)""" | |
| # Full body kept identical to the one you provided | |
| t0 = time.monotonic() | |
| span = None | |
| if OTEL_AVAILABLE and _tracer: | |
| span = _tracer.start_span("risk_service.evaluate_healing") | |
| span.set_attribute("component", event.component) | |
| if tenant_id: | |
| span.set_attribute("tenant_id", tenant_id) | |
| if decision_engine is None and hasattr(policy_engine, 'decision_engine'): | |
| decision_engine = policy_engine.decision_engine | |
| if decision_engine is None: | |
| logger.debug("No DecisionEngine provided; creating default instance") | |
| decision_engine = DecisionEngine(rag_graph=rag_graph) | |
| orig_use = policy_engine.use_decision_engine | |
| try: | |
| policy_engine.use_decision_engine = False | |
| raw_actions = policy_engine.evaluate_policies(event) | |
| finally: | |
| policy_engine.use_decision_engine = orig_use | |
| if not raw_actions or raw_actions == [HealingAction.NO_ACTION]: | |
| if span: | |
| span.set_attribute("selected_action", HealingAction.NO_ACTION.value) | |
| span.end() | |
| _EVAL_COUNTER.labels(engine="python", status="success").inc() | |
| _EVAL_DURATION.labels(engine="python").observe(time.monotonic() - t0) | |
| return { | |
| "risk_score": 0.0, | |
| "selected_action": HealingAction.NO_ACTION.value, | |
| "expected_utility": 0.0, | |
| "alternatives": [], | |
| "explanation": "No candidate actions triggered.", | |
| "epistemic_signals": None, | |
| } | |
| reasoning_parts = [] | |
| for policy in policy_engine.policies: | |
| if any(a in policy.actions for a in raw_actions): | |
| conditions_str = ", ".join( | |
| f"{c.metric} {c.operator} {c.threshold}" for c in policy.conditions | |
| ) | |
| reasoning_parts.append( | |
| f"Policy {policy.name} triggered by {conditions_str} → actions {[a.value for a in policy.actions]}" | |
| ) | |
| reasoning_text = " ".join(reasoning_parts) | |
| evidence_text = ( | |
| f"Component: {event.component}, " | |
| f"latency_p99: {event.latency_p99}, " | |
| f"error_rate: {event.error_rate}, " | |
| f"cpu_util: {event.cpu_util}, " | |
| f"memory_util: {event.memory_util}" | |
| ) | |
| epistemic_signals = None | |
| if model is not None and tokenizer is not None: | |
| try: | |
| epistemic_signals = compute_epistemic_risk( | |
| reasoning_text, evidence_text, model, tokenizer | |
| ) | |
| except Exception as e: | |
| logger.error(f"Failed to compute epistemic risk: {e}") | |
| epistemic_signals = { | |
| "entropy": 0.0, | |
| "contradiction": 0.0, | |
| "evidence_lift": 0.0, | |
| "hallucination_risk": 0.0, | |
| } | |
| else: | |
| logger.debug("Epistemic model/tokenizer not provided; using zero signals") | |
| epistemic_signals = { | |
| "entropy": 0.0, | |
| "contradiction": 0.0, | |
| "evidence_lift": 0.0, | |
| "hallucination_risk": 0.0, | |
| } | |
| decision = decision_engine.select_optimal_action( | |
| raw_actions, event, component=event.component, | |
| epistemic_signals=epistemic_signals | |
| ) | |
| risk_score = None | |
| for alt in decision.alternatives: | |
| if alt.action == decision.best_action: | |
| risk_score = alt.risk | |
| break | |
| if risk_score is None: | |
| risk_score = decision_engine.compute_risk( | |
| decision.best_action, event, event.component) | |
| alt_list = [] | |
| for alt in decision.alternatives[:3]: | |
| alt_list.append({ | |
| "action": alt.action.value, | |
| "expected_utility": alt.utility, | |
| "risk": alt.risk, | |
| }) | |
| _EVAL_COUNTER.labels(engine="python", status="success").inc() | |
| _EVAL_DURATION.labels(engine="python").observe(time.monotonic() - t0) | |
| if span: | |
| span.set_attribute("risk_score", risk_score) | |
| span.set_attribute("selected_action", decision.best_action.value) | |
| span.set_attribute("expected_utility", decision.expected_utility) | |
| span.end() | |
| return { | |
| "risk_score": risk_score, | |
| "selected_action": decision.best_action.value, | |
| "expected_utility": decision.expected_utility, | |
| "alternatives": alt_list, | |
| "explanation": decision.explanation, | |
| "raw_decision": decision.raw_data, | |
| "epistemic_signals": epistemic_signals, | |
| } | |
| def get_system_risk() -> float: | |
| """ | |
| Return an aggregated risk score across all monitored components. | |
| This endpoint is deprecated. Use component‑level risk evaluation instead. | |
| """ | |
| raise NotImplementedError( | |
| "get_system_risk is deprecated. Use component‑level risk evaluation instead." | |
| ) | |