diff --git a/src/intelligence/ecdt/__init__.py b/src/intelligence/ecdt/__init__.py index a588bd87..fd1edb53 100644 --- a/src/intelligence/ecdt/__init__.py +++ b/src/intelligence/ecdt/__init__.py @@ -1,11 +1,15 @@ """AEON MATRIX Enterprise Cognitive Digital Twin.""" +from .scenario_engine import ECDTScenarioEngine, ScenarioEvaluation + from .runtime import ( ECDTExecutionMode, ECDTRuntime, ) __all__ = [ + "ECDTScenarioEngine", + "ScenarioEvaluation", "ECDTExecutionMode", "ECDTRuntime", ] diff --git a/src/intelligence/ecdt/scenario_engine.py b/src/intelligence/ecdt/scenario_engine.py new file mode 100644 index 00000000..7e092f04 --- /dev/null +++ b/src/intelligence/ecdt/scenario_engine.py @@ -0,0 +1,148 @@ +"""Deterministic, side-effect-free scenario evaluation for AEON MATRIX ECDT.""" + +from __future__ import annotations + +from dataclasses import dataclass +from typing import Any, Dict, Iterable, Mapping, Sequence, Tuple + + +@dataclass(frozen=True) +class ScenarioEvaluation: + name: str + score: float + rank: int + status: str + metrics: Mapping[str, float] + trace: Tuple[str, ...] + + def to_dict(self) -> Dict[str, Any]: + return { + "name": self.name, + "score": self.score, + "rank": self.rank, + "status": self.status, + "metrics": dict(self.metrics), + "trace": list(self.trace), + } + + +class ECDTScenarioEngine: + """Evaluate and rank candidate scenarios without executing any action.""" + + def evaluate( + self, + *, + observed_state: Mapping[str, Any], + scenarios: Sequence[Mapping[str, Any]], + policy: Mapping[str, Any] | None = None, + ) -> Dict[str, Any]: + if observed_state is None: + raise ValueError("observed_state is required") + + if not scenarios: + raise ValueError("at least one scenario is required") + + allowed = True if policy is None else bool(policy.get("allowed", True)) + + evaluations = [ + self._evaluate_one( + observed_state=observed_state, + scenario=scenario, + allowed=allowed, + ) + for scenario in scenarios + ] + + ranked = sorted( + evaluations, + key=lambda item: (-item["score"], item["name"]), + ) + + results = [] + for index, item in enumerate(ranked, start=1): + result = ScenarioEvaluation( + name=item["name"], + score=item["score"], + rank=index, + status=item["status"], + metrics=item["metrics"], + trace=item["trace"], + ) + results.append(result.to_dict()) + + recommended = next( + ( + item + for item in results + if item["status"] == "ELIGIBLE" + ), + None, + ) + + return { + "status": "EVALUATED", + "executed": False, + "scenario_count": len(results), + "recommended": recommended, + "results": results, + "governance": { + "policy_allowed": allowed, + "execution_authorized": False, + }, + } + + def _evaluate_one( + self, + *, + observed_state: Mapping[str, Any], + scenario: Mapping[str, Any], + allowed: bool, + ) -> Dict[str, Any]: + name = str(scenario.get("name", "")).strip() + if not name: + raise ValueError("scenario name is required") + + impact = self._as_float( + scenario.get("impact_score", 0.0), + ) + risk = self._as_float( + scenario.get("risk_score", 0.0), + ) + confidence = self._as_float( + scenario.get("confidence", 1.0), + ) + + score = round( + (impact * confidence) - risk, + 6, + ) + + status = "ELIGIBLE" if allowed else "POLICY_BLOCKED" + + trace = ( + f"scenario={name}", + f"impact_score={impact}", + f"risk_score={risk}", + f"confidence={confidence}", + f"policy_allowed={allowed}", + f"observed_keys={','.join(sorted(map(str, observed_state.keys())))}", + ) + + return { + "name": name, + "score": score, + "status": status, + "metrics": { + "impact_score": impact, + "risk_score": risk, + "confidence": confidence, + }, + "trace": trace, + } + + @staticmethod + def _as_float(value: Any) -> float: + try: + return float(value) + except (TypeError, ValueError) as exc: + raise ValueError("scenario metric must be numeric") from exc diff --git a/tests/ecdt/test_scenario_engine.py b/tests/ecdt/test_scenario_engine.py new file mode 100644 index 00000000..4d34d948 --- /dev/null +++ b/tests/ecdt/test_scenario_engine.py @@ -0,0 +1,119 @@ +import pytest + +from src.intelligence.ecdt.scenario_engine import ECDTScenarioEngine + + +def observed(): + return { + "capacity": 0.95, + "demand": 1.10, + "inventory_risk": 0.30, + } + + +def scenarios(): + return [ + { + "name": "baseline", + "impact_score": 0.40, + "risk_score": 0.10, + "confidence": 0.90, + }, + { + "name": "dynamic_labor_scaling", + "impact_score": 0.80, + "risk_score": 0.20, + "confidence": 0.95, + }, + ] + + +def test_scenario_engine_is_deterministic(): + engine = ECDTScenarioEngine() + + first = engine.evaluate( + observed_state=observed(), + scenarios=scenarios(), + ) + second = engine.evaluate( + observed_state=observed(), + scenarios=scenarios(), + ) + + assert first == second + + +def test_best_scenario_is_recommended(): + result = ECDTScenarioEngine().evaluate( + observed_state=observed(), + scenarios=scenarios(), + ) + + assert result["status"] == "EVALUATED" + assert result["executed"] is False + assert result["recommended"]["name"] == "dynamic_labor_scaling" + assert result["recommended"]["rank"] == 1 + + +def test_policy_can_block_all_scenarios(): + result = ECDTScenarioEngine().evaluate( + observed_state=observed(), + scenarios=scenarios(), + policy={"allowed": False}, + ) + + assert result["recommended"] is None + assert all( + item["status"] == "POLICY_BLOCKED" + for item in result["results"] + ) + + +def test_trace_is_present_for_audit(): + result = ECDTScenarioEngine().evaluate( + observed_state=observed(), + scenarios=scenarios(), + ) + + assert result["results"][0]["trace"] + assert result["governance"]["execution_authorized"] is False + + +def test_missing_observed_state_raises(): + with pytest.raises(ValueError): + ECDTScenarioEngine().evaluate( + observed_state=None, + scenarios=scenarios(), + ) + + +def test_missing_scenarios_raises(): + with pytest.raises(ValueError): + ECDTScenarioEngine().evaluate( + observed_state=observed(), + scenarios=[], + ) + + +def test_invalid_metric_raises(): + data = scenarios() + data[0]["risk_score"] = "not-a-number" + + with pytest.raises(ValueError): + ECDTScenarioEngine().evaluate( + observed_state=observed(), + scenarios=data, + ) + + +def test_engine_has_no_execution_interface(): + engine = ECDTScenarioEngine() + + for name in ( + "execute", + "executor", + "apply", + "deploy", + "promote", + ): + assert not hasattr(engine, name)