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
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@

from datetime import date, datetime, timezone
from pathlib import Path
from typing import Sequence
from typing import Any, Mapping, Sequence

import pandas as pd

Expand Down Expand Up @@ -167,3 +167,50 @@ def run_monitor(
snapshots.append(snapshot)

return snapshots


def _serialize_decision(decision: Any) -> dict[str, Any]:
positions = getattr(decision, "positions", ()) or ()
serialized_positions: list[dict[str, Any]] = []
for position in positions:
serialized_positions.append(
{
"symbol": str(getattr(position, "symbol", "") or ""),
"target_weight": float(getattr(position, "target_weight", 0.0) or 0.0),
"role": str(getattr(position, "role", "") or ""),
}
)
return {
"positions": serialized_positions,
"risk_flags": list(getattr(decision, "risk_flags", ()) or ()),
"diagnostics": dict(getattr(decision, "diagnostics", {}) or {}),
}


class PerformanceMonitor:
"""Per-run performance / decision recorder for live strategy entrypoints."""

def __init__(self, store: PerformanceStore | None = None) -> None:
self._store = store or PerformanceStore.from_env()

def record(
self,
profile_id: str,
decision: Any,
execution_result: Mapping[str, Any] | None = None,
*,
domain: str = "",
) -> dict[str, Any]:
profile = str(profile_id or "").strip()
if not profile:
return {"ok": False, "skipped": "empty_profile"}

payload = {
"strategy_profile": profile,
"domain": str(domain or "").strip(),
"recorded_at": _now_iso(),
"decision": _serialize_decision(decision),
"execution_result": dict(execution_result or {}),
}
self._store.save_live_run_record(profile, str(domain or "").strip(), payload)
return {"ok": True, "profile": profile, "domain": str(domain or "").strip()}
18 changes: 18 additions & 0 deletions src/quant_platform_kit/strategy_lifecycle/performance_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,24 @@ def load_audit_entries(self, strategy_profile: str, limit: int = 20) -> tuple[Up
entries.append(entry)
return tuple(entries)

# ── live runs (per-evaluate / per-execution records) ─────────

def _live_run_key(self, domain: str, strategy_profile: str, recorded_at: str) -> str:
safe_time = recorded_at.replace(":", "-")
return f"live_runs/{_clean_key(domain)}/{_clean_key(strategy_profile)}/{safe_time}.json"

def save_live_run_record(
self,
strategy_profile: str,
domain: str,
payload: Mapping[str, Any],
) -> None:
recorded_at = str(payload.get("recorded_at") or _now_iso())
self._write(
self._live_run_key(domain, strategy_profile, recorded_at),
{**dict(payload), "schema_version": SCHEMA_VERSION},
)

# ── dashboard ────────────────────────────────────────────────

def _dashboard_key(self) -> str:
Expand Down
36 changes: 36 additions & 0 deletions tests/test_lifecycle_performance_monitor.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
from __future__ import annotations

import tempfile
import unittest
from pathlib import Path

from quant_platform_kit.common.strategy_contracts import PositionTarget, StrategyDecision
from quant_platform_kit.strategy_lifecycle.performance_monitor import PerformanceMonitor
from quant_platform_kit.strategy_lifecycle.performance_store import PerformanceStore


class PerformanceMonitorTests(unittest.TestCase):
def test_record_persists_live_run_payload(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
store = PerformanceStore(local_root=Path(tmp))
monitor = PerformanceMonitor(store=store)
decision = StrategyDecision(
positions=(PositionTarget(symbol="BTCUSDT", target_weight=0.5, role="target"),),
risk_flags=("risk_gate:passed",),
)
result = monitor.record(
"crypto_live_pool_rotation",
decision,
{"filled_orders": 1},
domain="crypto",
)
self.assertTrue(result["ok"])
files = list(Path(tmp).rglob("live_runs/crypto/crypto_live_pool_rotation/*.json"))
self.assertEqual(len(files), 1)
payload = files[0].read_text(encoding="utf-8")
self.assertIn("BTCUSDT", payload)
self.assertIn("filled_orders", payload)


if __name__ == "__main__":
unittest.main()
Loading