petter2025's picture
Upload folder using huggingface_hub
b9c169a verified
Raw
History Blame
20.4 kB
"""
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."
)