Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/intelligence/ecdt/__init__.py
Original file line number Diff line number Diff line change
@@ -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",
]
148 changes: 148 additions & 0 deletions src/intelligence/ecdt/scenario_engine.py
Original file line number Diff line number Diff line change
@@ -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
119 changes: 119 additions & 0 deletions tests/ecdt/test_scenario_engine.py
Original file line number Diff line number Diff line change
@@ -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)
Loading