121 lines
3.9 KiB
Python
121 lines
3.9 KiB
Python
|
|
"""
|
||
|
|
Structured Observability — per-decision feature attribution and metrics.
|
||
|
|
|
||
|
|
Tracks:
|
||
|
|
- Which features drove each decision
|
||
|
|
- Feature importance over time
|
||
|
|
- Decision quality metrics
|
||
|
|
- Regime-specific performance
|
||
|
|
"""
|
||
|
|
from __future__ import annotations
|
||
|
|
|
||
|
|
import json
|
||
|
|
import time
|
||
|
|
from collections import defaultdict
|
||
|
|
from dataclasses import dataclass, field
|
||
|
|
from typing import Any, Dict, List, Mapping, Optional, Tuple
|
||
|
|
|
||
|
|
from malkhut.state import MarketWorldState
|
||
|
|
from malkhut.actions import FulfilmentAction, PlannedPolicy, RiskDecision
|
||
|
|
from malkhut.features import DefaultFeatureExtractor, FeatureExtractor
|
||
|
|
|
||
|
|
|
||
|
|
@dataclass(frozen=True, slots=True)
|
||
|
|
class DecisionMetrics:
|
||
|
|
"""Per-decision metrics."""
|
||
|
|
ts_ns: int
|
||
|
|
symbol: str
|
||
|
|
action_kind: str
|
||
|
|
approved: bool
|
||
|
|
plan_latency_ns: int
|
||
|
|
entropy: float
|
||
|
|
sims: int
|
||
|
|
feature_attribution: Mapping[str, float]
|
||
|
|
regime: str
|
||
|
|
|
||
|
|
|
||
|
|
class StructuredObservability:
|
||
|
|
"""
|
||
|
|
Structured observability with per-decision feature attribution.
|
||
|
|
|
||
|
|
Tracks which features drive decisions and computes aggregate metrics.
|
||
|
|
"""
|
||
|
|
|
||
|
|
def __init__(self, feature_extractor: Optional[FeatureExtractor] = None) -> None:
|
||
|
|
self._extractor = feature_extractor or DefaultFeatureExtractor()
|
||
|
|
self._decisions: list[DecisionMetrics] = []
|
||
|
|
self._feature_importance: Dict[str, List[float]] = defaultdict(list)
|
||
|
|
self._regime_performance: Dict[str, List[float]] = defaultdict(list)
|
||
|
|
self._total_decisions = 0
|
||
|
|
|
||
|
|
def record_decision(
|
||
|
|
self,
|
||
|
|
state: MarketWorldState,
|
||
|
|
planned: PlannedPolicy,
|
||
|
|
decision: RiskDecision,
|
||
|
|
plan_ns: int,
|
||
|
|
regime: str = "unknown",
|
||
|
|
) -> None:
|
||
|
|
"""Record a decision with full feature attribution."""
|
||
|
|
fv = self._extractor.extract(state).values
|
||
|
|
|
||
|
|
# Compute feature attribution (which features are "active")
|
||
|
|
attribution = {}
|
||
|
|
for name, value in fv.items():
|
||
|
|
if abs(value) > 1e-6:
|
||
|
|
attribution[name] = value
|
||
|
|
|
||
|
|
metrics = DecisionMetrics(
|
||
|
|
ts_ns=state.ts_ns,
|
||
|
|
symbol=state.venue.symbol,
|
||
|
|
action_kind=planned.selected_action.kind.value,
|
||
|
|
approved=decision.approved,
|
||
|
|
plan_latency_ns=plan_ns,
|
||
|
|
entropy=planned.diagnostics.get("entropy", 0.0),
|
||
|
|
sims=planned.diagnostics.get("sims", 0),
|
||
|
|
feature_attribution=attribution,
|
||
|
|
regime=regime,
|
||
|
|
)
|
||
|
|
|
||
|
|
self._decisions.append(metrics)
|
||
|
|
self._total_decisions += 1
|
||
|
|
|
||
|
|
# Track feature importance
|
||
|
|
for name, value in attribution.items():
|
||
|
|
self._feature_importance[name].append(value)
|
||
|
|
|
||
|
|
# Track regime performance
|
||
|
|
self._regime_performance[regime].append(1.0 if decision.approved else 0.0)
|
||
|
|
|
||
|
|
def get_feature_importance(self, top_n: int = 10) -> List[Tuple[str, float]]:
|
||
|
|
"""Get top N features by average absolute value."""
|
||
|
|
scores = []
|
||
|
|
for name, values in self._feature_importance.items():
|
||
|
|
avg = sum(abs(v) for v in values) / len(values)
|
||
|
|
scores.append((name, avg))
|
||
|
|
scores.sort(key=lambda x: x[1], reverse=True)
|
||
|
|
return scores[:top_n]
|
||
|
|
|
||
|
|
def get_regime_approval_rate(self, regime: str) -> float:
|
||
|
|
"""Get approval rate for a specific regime."""
|
||
|
|
approvals = self._regime_performance.get(regime, [])
|
||
|
|
if not approvals:
|
||
|
|
return 0.0
|
||
|
|
return sum(approvals) / len(approvals)
|
||
|
|
|
||
|
|
@property
|
||
|
|
def total_decisions(self) -> int:
|
||
|
|
return self._total_decisions
|
||
|
|
|
||
|
|
@property
|
||
|
|
def avg_latency_ns(self) -> float:
|
||
|
|
if not self._decisions:
|
||
|
|
return 0.0
|
||
|
|
return sum(d.plan_latency_ns for d in self._decisions) / len(self._decisions)
|
||
|
|
|
||
|
|
@property
|
||
|
|
def avg_entropy(self) -> float:
|
||
|
|
if not self._decisions:
|
||
|
|
return 0.0
|
||
|
|
return sum(d.entropy for d in self._decisions) / len(self._decisions)
|