From 2cd3f0c3d410cb1511d0889244d7bf627257c214 Mon Sep 17 00:00:00 2001 From: BobbyAxerol Date: Wed, 29 Jul 2026 08:57:00 +0000 Subject: [PATCH 1/4] Add native event artifact memory policy --- backends/native_event.py | 500 ++++++++++++++++-- benchmarks/phase34a_native_event_memory.json | 50 ++ benchmarks/phase34a_native_event_memory.md | 13 + .../run_phase34a_native_event_memory.py | 179 +++++++ docs/endpoint.md | 38 ++ endpoint.py | 22 + engines.py | 15 + tests/test_phase34a_native_event_artifacts.py | 171 ++++++ upgrade/implement.md | 55 ++ 9 files changed, 993 insertions(+), 50 deletions(-) create mode 100644 benchmarks/phase34a_native_event_memory.json create mode 100644 benchmarks/phase34a_native_event_memory.md create mode 100644 benchmarks/run_phase34a_native_event_memory.py create mode 100644 tests/test_phase34a_native_event_artifacts.py diff --git a/backends/native_event.py b/backends/native_event.py index b5ce992..c148dbe 100644 --- a/backends/native_event.py +++ b/backends/native_event.py @@ -6,7 +6,8 @@ from __future__ import annotations -from dataclasses import dataclass, field, replace +from dataclasses import asdict, dataclass, field, replace +from pathlib import Path from typing import Dict, List, Optional, Sequence, Union import numpy as np @@ -126,6 +127,9 @@ class NativeEventConfig: execution: ExecutionConfig = field(default_factory=ExecutionConfig) fee_rate: Union[float, Dict[str, float]] = 0.0 use_funding: bool = True + report_level: str = "audit" + audit_sink: str = "memory" + audit_sink_path: Optional[str] = None def __post_init__(self) -> None: if isinstance(self.fee_rate, dict): @@ -133,6 +137,163 @@ def __post_init__(self) -> None: raise ValueError("fee_rate must be >= 0") elif float(self.fee_rate) < 0.0: raise ValueError("fee_rate must be >= 0") + object.__setattr__(self, "report_level", _normalize_native_event_report_level(self.report_level)) + object.__setattr__(self, "audit_sink", _normalize_native_event_audit_sink(self.audit_sink)) + + +@dataclass(frozen=True) +class NativeEventArtifactPlan: + keep_equity_path: bool + keep_position_path: bool + keep_fee_path: bool + keep_funding_path: bool + keep_margin_path: bool + keep_fill_ledger: bool + keep_command_terminal_state: bool + keep_event_ledger: bool + keep_command_tape: bool + materialize_pandas: bool + materialize_python_objects: bool + materialize_active_orders: bool + + +@dataclass(frozen=True) +class CompactFillLedger: + bar: np.ndarray + command_index: np.ndarray + original_index: np.ndarray + order_id_code: np.ndarray + symbol_code: np.ndarray + side: np.ndarray + qty: np.ndarray + price: np.ndarray + fee: np.ndarray + id_values: tuple[str, ...] + symbols: tuple[str, ...] + + @property + def fill_count(self) -> int: + return int(len(self.bar)) + + +@dataclass(frozen=True) +class CompactCommandLedger: + original_index: np.ndarray + command_bar: np.ndarray + action: np.ndarray + symbol_code: np.ndarray + side: np.ndarray + order_type: np.ndarray + order_id_code: np.ndarray + target_order_id_code: np.ndarray + parent_order_id_code: np.ndarray + group_id_code: np.ndarray + oco_group_id_code: np.ndarray + status: np.ndarray + reject_code: np.ndarray + fill_bar: np.ndarray + fill_qty: np.ndarray + fill_price: np.ndarray + fill_fee: np.ndarray + active: np.ndarray + waiting_parent: np.ndarray + working_qty: np.ndarray + working_price: np.ndarray + working_trigger: np.ndarray + id_values: tuple[str, ...] + symbols: tuple[str, ...] + + +@dataclass(frozen=True) +class CompactOrderEventLedger: + bar: np.ndarray + command_index: np.ndarray + event_type: np.ndarray + status: np.ndarray + related_command_index: np.ndarray + + @property + def event_count(self) -> int: + return int(len(self.bar)) + + +def _normalize_native_event_report_level(report_level: str) -> str: + level = str(report_level or "audit").lower().strip() + aliases = {"full": "audit", "debug": "audit", "research": "standard", "optimizer": "score", "scoring": "score"} + level = aliases.get(level, level) + if level not in {"score", "minimal", "standard", "audit"}: + raise ValueError("native_event report_level must be score, minimal, standard, audit, or full") + return level + + +def _normalize_native_event_audit_sink(audit_sink: str) -> str: + sink = str(audit_sink or "memory").lower().strip() + if sink not in {"none", "memory", "jsonl", "parquet"}: + raise ValueError("native_event audit_sink must be none, memory, jsonl, or parquet") + return sink + + +def _native_event_artifact_plan(report_level: str) -> NativeEventArtifactPlan: + level = _normalize_native_event_report_level(report_level) + if level == "score": + return NativeEventArtifactPlan( + keep_equity_path=True, + keep_position_path=True, + keep_fee_path=True, + keep_funding_path=True, + keep_margin_path=True, + keep_fill_ledger=False, + keep_command_terminal_state=True, + keep_event_ledger=False, + keep_command_tape=False, + materialize_pandas=True, + materialize_python_objects=False, + materialize_active_orders=False, + ) + if level == "minimal": + return NativeEventArtifactPlan( + keep_equity_path=True, + keep_position_path=True, + keep_fee_path=True, + keep_funding_path=True, + keep_margin_path=True, + keep_fill_ledger=True, + keep_command_terminal_state=True, + keep_event_ledger=False, + keep_command_tape=False, + materialize_pandas=True, + materialize_python_objects=False, + materialize_active_orders=False, + ) + if level == "standard": + return NativeEventArtifactPlan( + keep_equity_path=True, + keep_position_path=True, + keep_fee_path=True, + keep_funding_path=True, + keep_margin_path=True, + keep_fill_ledger=True, + keep_command_terminal_state=True, + keep_event_ledger=False, + keep_command_tape=False, + materialize_pandas=True, + materialize_python_objects=True, + materialize_active_orders=False, + ) + return NativeEventArtifactPlan( + keep_equity_path=True, + keep_position_path=True, + keep_fee_path=True, + keep_funding_path=True, + keep_margin_path=True, + keep_fill_ledger=True, + keep_command_terminal_state=True, + keep_event_ledger=True, + keep_command_tape=True, + materialize_pandas=True, + materialize_python_objects=True, + materialize_active_orders=True, + ) @dataclass @@ -756,6 +917,9 @@ def run_order_commands( slot_size: Optional[Union[float, Dict[str, float]]] = None, min_qty: Optional[Union[float, Dict[str, float]]] = None, min_notional: Optional[Union[float, Dict[str, float]]] = None, + report_level: Optional[str] = None, + audit_sink: Optional[str] = None, + audit_sink_path: Optional[str] = None, ) -> BacktestResultV2: """ Execute Phase 30B lifecycle `OrderCommand` tapes through event v2. @@ -765,6 +929,11 @@ def run_order_commands( phase. """ idx = validate_datetime(datetime_index) + requested_report_level = self.config.report_level if report_level is None else report_level + level = _normalize_native_event_report_level(requested_report_level) + plan = _native_event_artifact_plan(level) + sink = self.config.audit_sink if audit_sink is None else _normalize_native_event_audit_sink(audit_sink) + sink_path = self.config.audit_sink_path if audit_sink_path is None else audit_sink_path if symbols is None: symbol_list = list(closes.keys()) else: @@ -895,13 +1064,39 @@ def run_order_commands( use_funding=bool(self.config.use_funding), ) - fills = self._build_fills( - compiled_commands.sorted_commands, - idx, - fill_bar, - fill_qty, - fill_price, - fill_fee, + fill_ledger = self._build_compact_fill_ledger( + compiled_commands=compiled_commands, + fill_bar=fill_bar, + fill_qty=fill_qty, + fill_price=fill_price, + fill_fee=fill_fee, + ) + command_ledger = self._build_compact_command_ledger( + compiled_commands=compiled_commands, + command_status=command_status, + reject_code=reject_code, + fill_bar=fill_bar, + fill_qty=fill_qty, + fill_price=fill_price, + fill_fee=fill_fee, + active=active, + waiting_parent=waiting_parent, + working_qty=working_qty, + working_price=working_price, + working_trigger=working_trigger, + ) + event_ledger = self._build_compact_order_event_ledger( + event_count=int(event_count), + event_bar=event_bar, + event_command=event_command, + event_type=event_type, + event_status=event_status, + event_related_command=event_related_command, + ) + fills = ( + self._build_fills(compiled_commands.sorted_commands, idx, fill_bar, fill_qty, fill_price, fill_fee) + if plan.materialize_python_objects + else () ) equity = pd.Series(equity_arr, index=idx, name="equity") positions = pd.DataFrame( @@ -920,36 +1115,86 @@ def run_order_commands( }, index=idx, ) - command_report = self._build_command_report( - compiled_commands, - command_status, - reject_code, - fill_bar, - fill_qty, - fill_price, - fill_fee, - active, - waiting_parent, - working_qty, - working_price, - working_trigger, - ) - order_events = self._build_order_events( - idx=idx, - compiled_commands=compiled_commands, - event_count=int(event_count), - event_bar=event_bar, - event_command=event_command, - event_type=event_type, - event_status=event_status, - event_related_command=event_related_command, - ) - if command_report.empty: + if level in {"standard", "audit"}: + command_report = self._build_command_report( + compiled_commands, + command_status, + reject_code, + fill_bar, + fill_qty, + fill_price, + fill_fee, + active, + waiting_parent, + working_qty, + working_price, + working_trigger, + ) + else: + command_report = pd.DataFrame() + if level == "audit" and sink != "none": + order_events = self._build_order_events( + idx=idx, + compiled_commands=compiled_commands, + event_count=int(event_count), + event_bar=event_bar, + event_command=event_command, + event_type=event_type, + event_status=event_status, + event_related_command=event_related_command, + ) + else: + order_events = pd.DataFrame() + if command_report.empty or not plan.materialize_active_orders: active_orders = pd.DataFrame() else: active_orders = command_report[ (command_report["active"] == True) | (command_report["waiting_parent"] == True) # noqa: E712 ].copy() + audit_artifacts = self._write_native_event_audit_sink( + sink=sink, + sink_path=sink_path, + command_report=command_report, + order_events=order_events, + fill_ledger=fill_ledger, + command_ledger=command_ledger, + event_ledger=event_ledger, + report_level=level, + ) + lifecycle_counters = { + "fill_count": int(fill_ledger.fill_count), + "event_count": int(event_count), + "rejected_count": int(np.sum(command_status == ORDER_STATUS_REJECTED)), + "canceled_count": int(np.sum(command_status == ORDER_STATUS_CANCELED)), + "filled_command_count": int(np.sum(command_status == ORDER_STATUS_FILLED)), + "pending_command_count": int(np.sum(command_status == ORDER_STATUS_PENDING)), + "expired_event_count": int(np.sum(event_ledger.event_type == ORDER_EVENT_EXPIRE)), + } + metadata = { + "backend": "native_event", + "engine": "event_v2_lifecycle", + "report_level": level, + "report_level_requested": str(requested_report_level), + "artifact_plan": asdict(plan), + "audit_sink": sink, + "audit_sink_path": sink_path, + "audit_artifacts": audit_artifacts, + "fee_rate_oneway": self._fee_rate_metadata(fee_rates, symbol_list), + "slippage_bps": self.config.execution.slippage_bps, + "order_report": command_report, + "command_report": command_report, + "order_events": order_events, + "active_orders": active_orders, + "compact_fill_ledger": fill_ledger if plan.keep_fill_ledger else None, + "compact_command_ledger": command_ledger if plan.keep_command_terminal_state else None, + "compact_order_event_ledger": event_ledger if plan.keep_event_ledger and sink == "memory" else None, + "id_values": compiled_commands.id_values, + "quantity_constraints": constraints.as_dict(), + "quantity_preflight": quantity_preflight, + "initial_buying_power": self.config.account.initial_capital * float(np.mean(leverages)), + "liquidation_reason": int(liq_reason), + "lifecycle_counters": lifecycle_counters, + } return BacktestResultV2( equity=equity, @@ -961,7 +1206,7 @@ def run_order_commands( leverage=float(np.mean(leverages)), liquidated=bool(liq_flag), liquidation_bar=int(liq_idx), - orders=self._commands_to_order_intents(compiled_commands.sorted_commands), + orders=self._commands_to_order_intents(compiled_commands.sorted_commands) if plan.materialize_python_objects else (), fills=tuple(fills), fees=pd.Series(fee_arr, index=idx, name="fees"), funding=pd.Series(funding_arr, index=idx, name="funding"), @@ -973,21 +1218,7 @@ def run_order_commands( index=idx, ), diagnostics=diagnostics, - metadata={ - "backend": "native_event", - "engine": "event_v2_lifecycle", - "fee_rate_oneway": self._fee_rate_metadata(fee_rates, symbol_list), - "slippage_bps": self.config.execution.slippage_bps, - "order_report": command_report, - "command_report": command_report, - "order_events": order_events, - "active_orders": active_orders, - "id_values": compiled_commands.id_values, - "quantity_constraints": constraints.as_dict(), - "quantity_preflight": quantity_preflight, - "initial_buying_power": self.config.account.initial_capital * float(np.mean(leverages)), - "liquidation_reason": int(liq_reason), - }, + metadata=metadata, ) def run_strategy( @@ -1012,6 +1243,9 @@ def run_strategy( min_notional: Optional[Union[float, Dict[str, float]]] = None, execution_mode: str = "fast", command_effective_phase: str = "next_bar", + report_level: Optional[str] = None, + audit_sink: Optional[str] = None, + audit_sink_path: Optional[str] = None, ) -> BacktestResultV2: """ Run a reactive strategy against native-event v2 lifecycle semantics. @@ -1028,6 +1262,9 @@ def run_strategy( execution_mode = str(execution_mode).lower().strip() if execution_mode not in {"fast", "audit"}: raise ValueError("execution_mode must be 'fast' or 'audit'") + requested_report_level = self.config.report_level if report_level is None else report_level + level = _normalize_native_event_report_level(requested_report_level) + plan = _native_event_artifact_plan(level) idx = validate_datetime(datetime_index) symbol_list = list(symbols) if symbols is not None else list(closes.keys()) @@ -1153,13 +1390,17 @@ def run_strategy( slot_size=slot_size, min_qty=min_qty, min_notional=min_notional, + report_level=level, + audit_sink=audit_sink, + audit_sink_path=audit_sink_path, ) final_result.metadata.update( { "engine": "event_v2_reactive_incremental", "reactive_execution_mode": execution_mode, "command_effective_phase": "next_bar", - "emitted_command_tape": tuple(emitted), + "emitted_command_tape": tuple(emitted) if plan.keep_command_tape else (), + "emitted_command_tape_retained": bool(plan.keep_command_tape), "emitted_command_count": len(emitted), "ignored_commands_after_end": int(ignored_commands_after_end), "strategy_callback_count": int(callback_count), @@ -1837,6 +2078,165 @@ def size_order(symbol: str, notional: float, price: float, side: OrderSide = Ord return size_order + @staticmethod + def _build_compact_fill_ledger( + *, + compiled_commands: CompiledOrderCommandArrays, + fill_bar: np.ndarray, + fill_qty: np.ndarray, + fill_price: np.ndarray, + fill_fee: np.ndarray, + ) -> CompactFillLedger: + mask = (fill_bar >= 0) & (fill_qty != 0.0) + command_index = np.nonzero(mask)[0].astype(np.int64) + return CompactFillLedger( + bar=np.ascontiguousarray(fill_bar[mask], dtype=np.int64), + command_index=np.ascontiguousarray(command_index, dtype=np.int64), + original_index=np.ascontiguousarray(compiled_commands.original_index[mask], dtype=np.int64), + order_id_code=np.ascontiguousarray(compiled_commands.command_order_id[mask], dtype=np.int64), + symbol_code=np.ascontiguousarray(compiled_commands.command_symbol[mask], dtype=np.int64), + side=np.ascontiguousarray(compiled_commands.command_side[mask], dtype=np.int64), + qty=np.ascontiguousarray(fill_qty[mask], dtype=np.float64), + price=np.ascontiguousarray(fill_price[mask], dtype=np.float64), + fee=np.ascontiguousarray(fill_fee[mask], dtype=np.float64), + id_values=tuple(compiled_commands.id_values), + symbols=tuple(compiled_commands.symbols), + ) + + @staticmethod + def _build_compact_command_ledger( + *, + compiled_commands: CompiledOrderCommandArrays, + command_status: np.ndarray, + reject_code: np.ndarray, + fill_bar: np.ndarray, + fill_qty: np.ndarray, + fill_price: np.ndarray, + fill_fee: np.ndarray, + active: np.ndarray, + waiting_parent: np.ndarray, + working_qty: np.ndarray, + working_price: np.ndarray, + working_trigger: np.ndarray, + ) -> CompactCommandLedger: + return CompactCommandLedger( + original_index=np.ascontiguousarray(compiled_commands.original_index, dtype=np.int64), + command_bar=np.ascontiguousarray(compiled_commands.command_bar, dtype=np.int64), + action=np.ascontiguousarray(compiled_commands.command_action, dtype=np.int64), + symbol_code=np.ascontiguousarray(compiled_commands.command_symbol, dtype=np.int64), + side=np.ascontiguousarray(compiled_commands.command_side, dtype=np.int64), + order_type=np.ascontiguousarray(compiled_commands.command_type, dtype=np.int64), + order_id_code=np.ascontiguousarray(compiled_commands.command_order_id, dtype=np.int64), + target_order_id_code=np.ascontiguousarray(compiled_commands.command_target_order_id, dtype=np.int64), + parent_order_id_code=np.ascontiguousarray(compiled_commands.command_parent_order_id, dtype=np.int64), + group_id_code=np.ascontiguousarray(compiled_commands.command_group_id, dtype=np.int64), + oco_group_id_code=np.ascontiguousarray(compiled_commands.command_oco_group_id, dtype=np.int64), + status=np.ascontiguousarray(command_status, dtype=np.int64), + reject_code=np.ascontiguousarray(reject_code, dtype=np.int64), + fill_bar=np.ascontiguousarray(fill_bar, dtype=np.int64), + fill_qty=np.ascontiguousarray(fill_qty, dtype=np.float64), + fill_price=np.ascontiguousarray(fill_price, dtype=np.float64), + fill_fee=np.ascontiguousarray(fill_fee, dtype=np.float64), + active=np.ascontiguousarray(active, dtype=np.int64), + waiting_parent=np.ascontiguousarray(waiting_parent, dtype=np.int64), + working_qty=np.ascontiguousarray(working_qty, dtype=np.float64), + working_price=np.ascontiguousarray(working_price, dtype=np.float64), + working_trigger=np.ascontiguousarray(working_trigger, dtype=np.float64), + id_values=tuple(compiled_commands.id_values), + symbols=tuple(compiled_commands.symbols), + ) + + @staticmethod + def _build_compact_order_event_ledger( + *, + event_count: int, + event_bar: np.ndarray, + event_command: np.ndarray, + event_type: np.ndarray, + event_status: np.ndarray, + event_related_command: np.ndarray, + ) -> CompactOrderEventLedger: + n = max(int(event_count), 0) + return CompactOrderEventLedger( + bar=np.ascontiguousarray(event_bar[:n], dtype=np.int64), + command_index=np.ascontiguousarray(event_command[:n], dtype=np.int64), + event_type=np.ascontiguousarray(event_type[:n], dtype=np.int64), + status=np.ascontiguousarray(event_status[:n], dtype=np.int64), + related_command_index=np.ascontiguousarray(event_related_command[:n], dtype=np.int64), + ) + + @staticmethod + def _write_native_event_audit_sink( + *, + sink: str, + sink_path: Optional[str], + command_report: pd.DataFrame, + order_events: pd.DataFrame, + fill_ledger: CompactFillLedger, + command_ledger: CompactCommandLedger, + event_ledger: CompactOrderEventLedger, + report_level: str, + ) -> Dict: + if sink in {"none", "memory"} or report_level != "audit": + return {} + if not sink_path: + raise ValueError("native_event audit_sink='jsonl' or 'parquet' requires audit_sink_path") + root = Path(sink_path) + root.mkdir(parents=True, exist_ok=True) + if sink == "jsonl": + command_path = root / "command_report.jsonl" + event_path = root / "order_events.jsonl" + fill_path = root / "fill_ledger.jsonl" + command_report.to_json(command_path, orient="records", lines=True, date_format="iso") + order_events.to_json(event_path, orient="records", lines=True, date_format="iso") + pd.DataFrame( + { + "bar": fill_ledger.bar, + "command_index": fill_ledger.command_index, + "original_index": fill_ledger.original_index, + "order_id_code": fill_ledger.order_id_code, + "symbol_code": fill_ledger.symbol_code, + "side": fill_ledger.side, + "qty": fill_ledger.qty, + "price": fill_ledger.price, + "fee": fill_ledger.fee, + } + ).to_json(fill_path, orient="records", lines=True, date_format="iso") + return { + "format": "jsonl", + "command_report": str(command_path), + "order_events": str(event_path), + "fill_ledger": str(fill_path), + "event_count": int(event_ledger.event_count), + "fill_count": int(fill_ledger.fill_count), + } + command_path = root / "command_report.parquet" + event_path = root / "order_events.parquet" + fill_path = root / "fill_ledger.parquet" + command_report.to_parquet(command_path, index=False) + order_events.to_parquet(event_path, index=False) + pd.DataFrame( + { + "bar": fill_ledger.bar, + "command_index": fill_ledger.command_index, + "original_index": fill_ledger.original_index, + "order_id_code": fill_ledger.order_id_code, + "symbol_code": fill_ledger.symbol_code, + "side": fill_ledger.side, + "qty": fill_ledger.qty, + "price": fill_ledger.price, + "fee": fill_ledger.fee, + } + ).to_parquet(fill_path, index=False) + return { + "format": "parquet", + "command_report": str(command_path), + "order_events": str(event_path), + "fill_ledger": str(fill_path), + "event_count": int(event_ledger.event_count), + "fill_count": int(fill_ledger.fill_count), + } + @staticmethod def _build_command_report( compiled_commands: CompiledOrderCommandArrays, diff --git a/benchmarks/phase34a_native_event_memory.json b/benchmarks/phase34a_native_event_memory.json new file mode 100644 index 0000000..f756bc9 --- /dev/null +++ b/benchmarks/phase34a_native_event_memory.json @@ -0,0 +1,50 @@ +[ + { + "audit_sink": "memory", + "command_report_rows": 0, + "commands": 1575, + "events": 3075, + "fills": 1500, + "fills_materialized": 0, + "final_equity": 100006.59999999916, + "levels": 10, + "order_event_rows": 0, + "orders_materialized": 0, + "peak_rss_mb": 333.08984375, + "report_level": "minimal", + "rows": 3000, + "seconds": 1.22510303882882 + }, + { + "audit_sink": "memory", + "command_report_rows": 1575, + "commands": 1575, + "events": 3075, + "fills": 1500, + "fills_materialized": 1500, + "final_equity": 100006.59999999916, + "levels": 10, + "order_event_rows": 0, + "orders_materialized": 1500, + "peak_rss_mb": 336.640625, + "report_level": "standard", + "rows": 3000, + "seconds": 1.045848773792386 + }, + { + "audit_sink": "memory", + "command_report_rows": 1575, + "commands": 1575, + "events": 3075, + "fills": 1500, + "fills_materialized": 1500, + "final_equity": 100006.59999999916, + "levels": 10, + "order_event_rows": 3075, + "orders_materialized": 1500, + "peak_rss_mb": 292.08984375, + "report_level": "audit", + "rows": 3000, + "seconds": 1.009142744820565 + } +] diff --git a/benchmarks/phase34a_native_event_memory.md b/benchmarks/phase34a_native_event_memory.md new file mode 100644 index 0000000..5ad7492 --- /dev/null +++ b/benchmarks/phase34a_native_event_memory.md @@ -0,0 +1,13 @@ +# Phase 34A Native Event Artifact Memory Benchmark + +| report_level | seconds | peak RSS MB | commands | fills | events | command rows | event rows | fills obj | orders obj | +|---|---:|---:|---:|---:|---:|---:|---:|---:|---:| +| minimal | 1.225103 | 333.090 | 1575 | 1500 | 3075 | 0 | 0 | 0 | 0 | +| standard | 1.045849 | 336.641 | 1575 | 1500 | 3075 | 1575 | 0 | 1500 | 1500 | +| audit | 1.009143 | 292.090 | 1575 | 1500 | 3075 | 1575 | 3075 | 1500 | 1500 | + +Notes: + +- Each row runs in a fresh subprocess. +- Peak RSS includes Python import, pandas, and Numba/cache overhead; on small workloads it is not expected to be monotonic by artifact level. +- The artifact contract is verified by command/event row counts and materialized Python object counts; larger command-heavy runs are needed for stable RSS deltas. diff --git a/benchmarks/run_phase34a_native_event_memory.py b/benchmarks/run_phase34a_native_event_memory.py new file mode 100644 index 0000000..8e3b188 --- /dev/null +++ b/benchmarks/run_phase34a_native_event_memory.py @@ -0,0 +1,179 @@ +from __future__ import annotations + +import argparse +import json +import resource +import subprocess +import sys +import time +from pathlib import Path + +import numpy as np +import pandas as pd + +from quantbt import AccountConfig, ExecutionConfig, NativeEventBackend, NativeEventConfig +from quantbt.core.orders import OrderAction, OrderCommand +from quantbt.core.schema import OrderSide, OrderType, TimeInForce + + +def _rss_mb() -> float: + value = float(resource.getrusage(resource.RUSAGE_SELF).ru_maxrss) + if sys.platform == "darwin": + return value / (1024.0 * 1024.0) + return value / 1024.0 + + +def _market(rows: int): + idx = pd.date_range("2020-01-01", periods=rows, freq="15min", tz="UTC") + x = np.arange(rows, dtype=np.float64) + close = pd.Series(100.0 + np.sin(x / 17.0) * 2.0 + x * 0.0001, index=idx) + high = close + 1.2 + low = close - 1.2 + return idx, {"BTC": close}, {"BTC": high}, {"BTC": low} + + +def _commands(idx: pd.DatetimeIndex, levels: int, cycle: int): + commands = [] + order_id = 0 + for bar in range(1, len(idx), cycle): + commands.append(OrderCommand(timestamp=idx[bar], action=OrderAction.CANCEL_ALL, symbol="BTC")) + anchor = 100.0 + np.sin(bar / 17.0) * 2.0 + bar * 0.0001 + for level in range(1, levels + 1): + commands.append( + OrderCommand( + timestamp=idx[bar], + symbol="BTC", + side=OrderSide.BUY, + order_type=OrderType.LIMIT, + qty=0.01, + price=float(anchor - 0.08 * level), + tif=TimeInForce.GTC, + order_id=f"entry-{order_id}", + tag=f"GRID-C{bar}-L{level}", + metadata={"campaign_id": f"C{bar}", "level_id": str(level)}, + ) + ) + order_id += 1 + commands.append( + OrderCommand( + timestamp=idx[bar], + symbol="BTC", + side=OrderSide.SELL, + order_type=OrderType.LIMIT, + qty=0.01, + price=float(anchor + 0.08 * level), + tif=TimeInForce.GTC, + reduce_only=True, + order_id=f"exit-{order_id}", + tag=f"GRID-C{bar}-X{level}", + metadata={"campaign_id": f"C{bar}", "level_id": str(level), "leg_role": "take_profit"}, + ) + ) + order_id += 1 + return tuple(commands) + + +def _run_child(args) -> dict: + idx, close, high, low = _market(args.rows) + commands = _commands(idx, levels=args.levels, cycle=args.cycle) + backend = NativeEventBackend( + NativeEventConfig( + account=AccountConfig(initial_capital=100_000.0, leverage=10.0), + execution=ExecutionConfig(slippage_bps=0.0), + fee_rate=0.0, + use_funding=False, + report_level=args.report_level, + audit_sink=args.audit_sink, + audit_sink_path=args.audit_sink_path, + ) + ) + start = time.perf_counter() + result = backend.run_order_commands(idx, commands, close, high, low, symbols=["BTC"]) + elapsed = time.perf_counter() - start + payload = { + "report_level": result.metadata["report_level"], + "audit_sink": result.metadata["audit_sink"], + "rows": int(args.rows), + "levels": int(args.levels), + "commands": int(len(commands)), + "fills": int(result.metadata["lifecycle_counters"]["fill_count"]), + "events": int(result.metadata["lifecycle_counters"]["event_count"]), + "seconds": float(elapsed), + "peak_rss_mb": float(_rss_mb()), + "command_report_rows": int(len(result.metadata["command_report"])), + "order_event_rows": int(len(result.metadata["order_events"])), + "fills_materialized": int(len(result.fills)), + "orders_materialized": int(len(result.orders)), + "final_equity": float(result.equity.iloc[-1]), + } + print(json.dumps(payload, sort_keys=True)) + return payload + + +def _run_parent(args) -> list[dict]: + rows = [] + for level in ("minimal", "standard", "audit"): + cmd = [ + sys.executable, + __file__, + "--child", + "--rows", + str(args.rows), + "--levels", + str(args.levels), + "--cycle", + str(args.cycle), + "--report-level", + level, + ] + completed = subprocess.run(cmd, check=True, capture_output=True, text=True) + rows.append(json.loads(completed.stdout.strip().splitlines()[-1])) + if args.json_out: + Path(args.json_out).write_text(json.dumps(rows, indent=2, sort_keys=True) + "\n") + if args.md_out: + lines = [ + "# Phase 34A Native Event Artifact Memory Benchmark", + "", + "| report_level | seconds | peak RSS MB | commands | fills | events | command rows | event rows | fills obj | orders obj |", + "|---|---:|---:|---:|---:|---:|---:|---:|---:|---:|", + ] + for row in rows: + lines.append( + "| {report_level} | {seconds:.6f} | {peak_rss_mb:.3f} | {commands} | {fills} | {events} | " + "{command_report_rows} | {order_event_rows} | {fills_materialized} | {orders_materialized} |".format(**row) + ) + lines.extend( + [ + "", + "Notes:", + "", + "- Each row runs in a fresh subprocess.", + "- Peak RSS includes Python import, pandas, and Numba/cache overhead; on small workloads it is not expected to be monotonic by artifact level.", + "- The artifact contract is verified by command/event row counts and materialized Python object counts; larger command-heavy runs are needed for stable RSS deltas.", + ] + ) + Path(args.md_out).write_text("\n".join(lines) + "\n") + print(json.dumps(rows, indent=2, sort_keys=True)) + return rows + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--child", action="store_true") + parser.add_argument("--rows", type=int, default=5_000) + parser.add_argument("--levels", type=int, default=15) + parser.add_argument("--cycle", type=int, default=50) + parser.add_argument("--report-level", default="audit") + parser.add_argument("--audit-sink", default="memory") + parser.add_argument("--audit-sink-path", default=None) + parser.add_argument("--json-out", default="benchmarks/phase34a_native_event_memory.json") + parser.add_argument("--md-out", default="benchmarks/phase34a_native_event_memory.md") + args = parser.parse_args() + if args.child: + _run_child(args) + else: + _run_parent(args) + + +if __name__ == "__main__": + main() diff --git a/docs/endpoint.md b/docs/endpoint.md index acef52b..cd4264e 100644 --- a/docs/endpoint.md +++ b/docs/endpoint.md @@ -867,6 +867,39 @@ result = bt.simulate( ) ``` +Native-event artifact policy: + +```python +bt = QuantBTEndpoint.native_event_lifecycle( + initial_capital=100_000, + leverage=5, + report_level="minimal", # minimal | standard | audit | full + audit_sink="none", # none | memory | jsonl | parquet +) +``` + +`report_level` changes only artifact retention. It must not change equity, +positions, fees, funding, margin, liquidation, or lifecycle counters. + +| Level | Intended use | Retained artifacts | +|---|---|---| +| `minimal` | WFO/service loops | accounting paths, diagnostics, compact fill/command ledgers, no Python fills/orders, no event DataFrame | +| `standard` | research | minimal artifacts plus Python fills and command terminal report | +| `audit` / `full` | certification | full command report, event report, active-order report, Python fills/orders, compact ledgers | + +For long audits, use a disk sink: + +```python +bt = QuantBTEndpoint.native_event_lifecycle( + report_level="audit", + audit_sink="jsonl", + audit_sink_path="/tmp/quantbt_native_event_audit", +) +``` + +`jsonl` and `parquet` sinks require an explicit `audit_sink_path`; QuantBT does +not silently create long-lived audit bundles in arbitrary project folders. + Execution rules: - market orders fill on the bar close with slippage; @@ -1009,6 +1042,11 @@ result.metadata["reactive_incremental_compile_replays"] # 0 result.metadata["emitted_command_tape"] # replayable OrderCommand tape ``` +For reactive strategies, `report_level="minimal"` intentionally omits +`emitted_command_tape` from metadata while preserving +`emitted_command_count`. Use `report_level="audit"` when a replayable command +tape is required for certification. + Scoped cancel-all: ```python diff --git a/endpoint.py b/endpoint.py index 9072867..a52aa6c 100644 --- a/endpoint.py +++ b/endpoint.py @@ -199,6 +199,8 @@ class EndpointConfig: nautilus_depth_config: Optional[NautilusExecutionDepthConfig] = None option_config: object = None report_level: str = "full" + audit_sink: str = "memory" + audit_sink_path: Optional[str] = None strategy_class: object = None walkforward_config: Optional[WalkForwardConfig] = None walkforward_target_mode: str = "signal_notional" @@ -1787,6 +1789,9 @@ def _intrabar_execution_kwargs(self, symbol: str) -> Dict: slot_size=self.config.slot_size, min_qty=self.config.min_qty, min_notional=self.config.min_notional, + report_level=self.config.report_level, + audit_sink=self.config.audit_sink, + audit_sink_path=self.config.audit_sink_path, ) sizing_mode = IntrabarSizingMode(str(self.config.metadata.get("intrabar_sizing_mode", IntrabarSizingMode.UNITS.value))) fixed_notional = float(self.config.metadata.get("fixed_notional", self.config.alloc_per_trade if not isinstance(self.config.alloc_per_trade, dict) else self.config.alloc_per_trade.get(symbol, 0.0))) @@ -1858,6 +1863,9 @@ def _run_single(self, data, signal, signal_col, datetime_index, symbols): slot_size=self.config.slot_size, min_qty=self.config.min_qty, min_notional=self.config.min_notional, + report_level=self.config.report_level, + audit_sink=self.config.audit_sink, + audit_sink_path=self.config.audit_sink_path, ) markers = _intrabar_marker_columns(frame) if backend == "native_vectorized" and markers: @@ -1899,6 +1907,9 @@ def _run_orders(self, data, orders, order_commands, datetime_index, symbols): slot_size=self.config.slot_size, min_qty=self.config.min_qty, min_notional=self.config.min_notional, + report_level=self.config.report_level, + audit_sink=self.config.audit_sink, + audit_sink_path=self.config.audit_sink_path, ) self._store_result(self.engine.result) return self.result @@ -1932,6 +1943,9 @@ def _run_native_event_strategy(self, data, strategy, datetime_index, symbols): slot_size=self.config.slot_size, min_qty=self.config.min_qty, min_notional=self.config.min_notional, + report_level=self.config.report_level, + audit_sink=self.config.audit_sink, + audit_sink_path=self.config.audit_sink_path, ) self._store_result(self.engine.result) return self.result @@ -1981,6 +1995,9 @@ def _run_structured_orders(self, data, datetime_index, symbols): slot_size=self.config.slot_size, min_qty=self.config.min_qty, min_notional=self.config.min_notional, + report_level=self.config.report_level, + audit_sink=self.config.audit_sink, + audit_sink_path=self.config.audit_sink_path, ) result = self.engine.result result.metadata.update( @@ -2124,6 +2141,9 @@ def _run_arbitrage(self, data, signal, signal_col, closes, highs, lows, hedge_ra execution=self.config.execution, fee_rate=self.config.v2_fee_rate, use_funding=self.config.use_funding, + report_level=self.config.report_level, + audit_sink=self.config.audit_sink, + audit_sink_path=self.config.audit_sink_path, ) ) else: @@ -2723,6 +2743,8 @@ def _endpoint_run_config_payload(config: EndpointConfig) -> Dict: "symbols": _jsonable(config.symbols), "metadata": _jsonable(config.metadata), "report_level": config.report_level, + "audit_sink": config.audit_sink, + "audit_sink_path": config.audit_sink_path, } if config.nautilus_config is not None: payload["nautilus"] = _jsonable( diff --git a/engines.py b/engines.py index 6603039..44336c8 100644 --- a/engines.py +++ b/engines.py @@ -69,6 +69,9 @@ def __init__( strategy=None, event_engine_version: str = "v1", reactive_execution_mode: str = "fast", + report_level: str = "audit", + audit_sink: str = "memory", + audit_sink_path: Optional[str] = None, datetime_index: Optional[Union[pd.DatetimeIndex, pd.Series]] = None, closes: Optional[SeriesMap] = None, highs: Optional[SeriesMap] = None, @@ -109,6 +112,9 @@ def __init__( self.strategy = strategy self.event_engine_version = str(event_engine_version).lower().strip() self.reactive_execution_mode = str(reactive_execution_mode).lower().strip() + self.report_level = str(report_level) + self.audit_sink = str(audit_sink) + self.audit_sink_path = audit_sink_path self.datetime_index = datetime_index self.closes = closes self.highs = highs @@ -205,6 +211,9 @@ def _run_native_event(self) -> BacktestResultV2: execution=self.execution, fee_rate=self.fee_rate, use_funding=self.use_funding, + report_level=self.report_level, + audit_sink=self.audit_sink, + audit_sink_path=self.audit_sink_path, ) ) @@ -235,6 +244,9 @@ def _run_native_event(self) -> BacktestResultV2: min_qty=self.min_qty, min_notional=self.min_notional, execution_mode=self.reactive_execution_mode, + report_level=self.report_level, + audit_sink=self.audit_sink, + audit_sink_path=self.audit_sink_path, ) if self.basket is not None: @@ -297,6 +309,9 @@ def _run_native_event(self) -> BacktestResultV2: slot_size=self.slot_size, min_qty=self.min_qty, min_notional=self.min_notional, + report_level=self.report_level, + audit_sink=self.audit_sink, + audit_sink_path=self.audit_sink_path, ) orders = self.orders diff --git a/tests/test_phase34a_native_event_artifacts.py b/tests/test_phase34a_native_event_artifacts.py new file mode 100644 index 0000000..5b643ca --- /dev/null +++ b/tests/test_phase34a_native_event_artifacts.py @@ -0,0 +1,171 @@ +from __future__ import annotations + +import pandas as pd + +from quantbt import AccountConfig, ExecutionConfig, NativeEventBackend, NativeEventConfig, QuantBTEndpoint +from quantbt.core.orders import OrderAction, OrderCommand +from quantbt.core.schema import OrderSide, OrderType, TimeInForce + + +def _market(n: int = 8): + idx = pd.date_range("2024-01-01", periods=n, freq="1h", tz="UTC") + close = pd.Series([100.0, 100.0, 101.0, 103.0, 98.0, 99.0, 104.0, 100.0][:n], index=idx) + high = pd.Series([101.0, 102.0, 104.0, 106.0, 100.0, 101.0, 106.0, 103.0][:n], index=idx) + low = pd.Series([99.0, 98.0, 99.0, 100.0, 94.0, 96.0, 101.0, 98.0][:n], index=idx) + frame = pd.DataFrame({"open": close, "high": high, "low": low, "close": close, "volume": 1_000.0}, index=idx) + return idx, frame, {"BTC": close}, {"BTC": high}, {"BTC": low} + + +def _commands(idx): + return [ + OrderCommand( + timestamp=idx[1], + symbol="BTC", + side=OrderSide.BUY, + order_type=OrderType.LIMIT, + qty=1.0, + price=99.0, + tif=TimeInForce.GTC, + order_id="entry", + ), + OrderCommand( + timestamp=idx[2], + symbol="BTC", + side=OrderSide.SELL, + order_type=OrderType.LIMIT, + qty=1.0, + price=105.0, + tif=TimeInForce.GTC, + reduce_only=True, + order_id="take-profit", + ), + OrderCommand( + timestamp=idx[2], + symbol="BTC", + side=OrderSide.SELL, + order_type=OrderType.LIMIT, + qty=1.0, + price=97.0, + tif=TimeInForce.GTD, + reduce_only=True, + order_id="expires", + expires_at=idx[4], + ), + OrderCommand(timestamp=idx[5], action=OrderAction.CANCEL, target_order_id="take-profit"), + ] + + +def _backend(report_level: str, audit_sink: str = "memory", audit_sink_path=None): + return NativeEventBackend( + NativeEventConfig( + account=AccountConfig(initial_capital=10_000.0, leverage=10.0), + execution=ExecutionConfig(slippage_bps=0.0), + fee_rate=0.0, + use_funding=False, + report_level=report_level, + audit_sink=audit_sink, + audit_sink_path=audit_sink_path, + ) + ) + + +def _assert_accounting_equal(left, right): + pd.testing.assert_series_equal(left.equity, right.equity) + pd.testing.assert_series_equal(left.returns, right.returns) + pd.testing.assert_frame_equal(left.positions, right.positions) + pd.testing.assert_series_equal(left.fees, right.fees) + pd.testing.assert_series_equal(left.funding, right.funding) + pd.testing.assert_frame_equal(left.margin, right.margin) + pd.testing.assert_frame_equal(left.diagnostics, right.diagnostics) + assert left.liquidated == right.liquidated + assert left.liquidation_bar == right.liquidation_bar + assert left.metadata["lifecycle_counters"] == right.metadata["lifecycle_counters"] + + +def test_native_event_report_levels_preserve_accounting_and_reduce_artifacts(): + idx, _, close, high, low = _market() + commands = _commands(idx) + + audit = _backend("audit").run_order_commands(idx, commands, close, high, low) + standard = _backend("standard").run_order_commands(idx, commands, close, high, low) + minimal = _backend("minimal").run_order_commands(idx, commands, close, high, low) + + _assert_accounting_equal(audit, standard) + _assert_accounting_equal(audit, minimal) + + assert audit.metadata["report_level"] == "audit" + assert standard.metadata["report_level"] == "standard" + assert minimal.metadata["report_level"] == "minimal" + + assert not audit.metadata["command_report"].empty + assert not audit.metadata["order_events"].empty + assert len(audit.fills) == audit.metadata["lifecycle_counters"]["fill_count"] + + assert not standard.metadata["command_report"].empty + assert standard.metadata["order_events"].empty + assert len(standard.fills) == audit.metadata["lifecycle_counters"]["fill_count"] + + assert minimal.metadata["command_report"].empty + assert minimal.metadata["order_events"].empty + assert minimal.fills == () + assert minimal.orders == () + assert minimal.metadata["compact_fill_ledger"].fill_count == audit.metadata["lifecycle_counters"]["fill_count"] + assert minimal.metadata["compact_order_event_ledger"] is None + assert minimal.metadata["compact_command_ledger"].status.tolist() == audit.metadata["compact_command_ledger"].status.tolist() + + +def test_native_event_audit_jsonl_sink_writes_trace_without_accounting_drift(tmp_path): + idx, _, close, high, low = _market() + commands = _commands(idx) + + memory = _backend("audit").run_order_commands(idx, commands, close, high, low) + disk = _backend("audit", audit_sink="jsonl", audit_sink_path=tmp_path).run_order_commands( + idx, + commands, + close, + high, + low, + ) + + _assert_accounting_equal(memory, disk) + artifacts = disk.metadata["audit_artifacts"] + assert artifacts["format"] == "jsonl" + assert artifacts["fill_count"] == memory.metadata["lifecycle_counters"]["fill_count"] + assert (tmp_path / "command_report.jsonl").exists() + assert (tmp_path / "order_events.jsonl").exists() + assert (tmp_path / "fill_ledger.jsonl").exists() + + +def test_endpoint_propagates_native_event_report_level_and_reactive_tape_policy(): + idx, frame, _, _, _ = _market() + + class Strategy: + def on_bar_close(self, context): + if context.bar_index == 0: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol="BTC", + side=OrderSide.BUY, + order_type=OrderType.MARKET, + qty=1.0, + tif=TimeInForce.IOC, + order_id="entry", + ) + ] + return [] + + endpoint = QuantBTEndpoint.native_event_strategy( + initial_capital=10_000, + leverage=10, + use_funding=False, + report_level="minimal", + ) + result = endpoint.simulate(data=frame, strategy=Strategy(), symbols=["BTC"]) + + assert result.metadata["report_level"] == "minimal" + assert result.metadata["emitted_command_count"] == 1 + assert result.metadata["emitted_command_tape"] == () + assert result.metadata["emitted_command_tape_retained"] is False + assert result.metadata["lifecycle_counters"]["fill_count"] == 1 + assert result.fills == () diff --git a/upgrade/implement.md b/upgrade/implement.md index 6a05c4f..43df844 100644 --- a/upgrade/implement.md +++ b/upgrade/implement.md @@ -6611,6 +6611,8 @@ Goal: ### Phase 34A - Native Event Artifact And Memory Contract +Status: implemented on `feat/30-native-event-lifecycle`. + Scope: - Wire `report_level` through native-event endpoints, configs, backend, kernel @@ -6646,6 +6648,59 @@ Acceptance: - Audit can retain full trace through memory or chunked disk sink. - Public `.simulate()` remains source-compatible. +Implemented: + +- Added `NativeEventConfig.report_level`, `audit_sink`, and `audit_sink_path`. +- Added `NativeEventArtifactPlan` with explicit artifact-retention flags. +- Added compact struct-of-arrays ledgers: + - `CompactFillLedger`; + - `CompactCommandLedger`; + - `CompactOrderEventLedger`. +- Wired report policy through: + - `QuantBTEndpoint`; + - `BacktestEngineV2`; + - native-event lifecycle v2 backend; + - reactive native-event strategy replay. +- `full` normalizes to `audit`; existing default behavior remains + audit-compatible. +- `minimal` keeps accounting paths and compact ledgers but omits heavy Python + fills/orders and command/event DataFrames. +- `standard` keeps command terminal report and Python fills but omits full + lifecycle event DataFrame. +- `audit` keeps full command report, order events, active-order report, Python + fills/orders, compact ledgers, and optional disk sink artifacts. +- Reactive minimal mode records `emitted_command_count` but does not retain the + full `emitted_command_tape`. +- Added `audit_sink="jsonl"` and `audit_sink="parquet"` support with explicit + `audit_sink_path`; no silent project-folder writes. +- Updated endpoint docs for native-event report levels and audit sinks. + +Validation: + +```bash +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q tests/test_phase34a_native_event_artifacts.py +# 3 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q tests/test_phase30a_native_event_lifecycle_contract.py tests/test_phase30b_native_event_lifecycle_kernel.py tests/test_phase30c_native_event_endpoint_lifecycle.py tests/test_phase30d_native_event_reactive_runner.py tests/test_phase30e_native_event_incremental_runner.py +# 33 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q tests/test_endpoint.py tests/test_phase14c_prepared_report_levels.py +# 26 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run python benchmarks/run_phase34a_native_event_memory.py --rows 3000 --levels 10 --cycle 40 +# artifact retention benchmark recorded in benchmarks/phase34a_native_event_memory.md +``` + +Benchmark interpretation: + +- `minimal` produced zero command-report rows, zero event-report rows, zero + materialized Python fills, and zero materialized Python orders for the test + workload while preserving final equity and lifecycle counters. +- The small subprocess RSS numbers include Python import, pandas, and + Numba/cache overhead, so they are not used as a strict memory delta claim. + Larger Phase 34B/34C optimization-batch benchmarks are still required before + claiming stable RSS reduction percentages. + ### Phase 34B - Prepared Native Event Score Path Scope: From 970c9bb915de3c5f98da723cae164da550ae08ac Mon Sep 17 00:00:00 2001 From: BobbyAxerol Date: Wed, 29 Jul 2026 09:17:41 +0000 Subject: [PATCH 2/4] Add prepared native event score path --- __init__.py | 15 +- backends/native_event.py | 38 ++- .../phase34b_native_event_prepared_score.json | 12 + .../phase34b_native_event_prepared_score.md | 12 + ...un_phase34b_native_event_prepared_score.py | 173 +++++++++++ core/__init__.py | 4 +- core/results.py | 103 +++++- docs/endpoint.md | 32 ++ endpoint.py | 242 ++++++++++++++- metrics/performance.py | 293 ++++++++++++++++-- optimization/__init__.py | 2 + optimization/evaluators/__init__.py | 2 + optimization/evaluators/native_event.py | 32 ++ ...st_phase34b_native_event_prepared_score.py | 143 +++++++++ upgrade/implement.md | 53 ++++ 15 files changed, 1109 insertions(+), 47 deletions(-) create mode 100644 benchmarks/phase34b_native_event_prepared_score.json create mode 100644 benchmarks/phase34b_native_event_prepared_score.md create mode 100644 benchmarks/run_phase34b_native_event_prepared_score.py create mode 100644 optimization/evaluators/native_event.py create mode 100644 tests/test_phase34b_native_event_prepared_score.py diff --git a/__init__.py b/__init__.py index 8d2ecf8..ed57c24 100644 --- a/__init__.py +++ b/__init__.py @@ -47,7 +47,14 @@ from .backtester import BacktestEngine from .portfolio import MultiSymbolPortfolio -from .endpoint import EndpointConfig, PreparedIntrabarRunner, QuantBTEndpoint, QuantBTPreparedContext, format_metrics_report +from .endpoint import ( + EndpointConfig, + PreparedIntrabarRunner, + PreparedNativeEventStrategyRunner, + QuantBTEndpoint, + QuantBTPreparedContext, + format_metrics_report, +) from .walkforward import ( DuplicatePruner, EarlyStoppingCallback, @@ -94,6 +101,7 @@ OptimizationTrialRecord, OptunaOptimizer, PreparedIntrabarEvaluator, + PreparedNativeEventStrategyEvaluator, PreparedPortfolioEvaluator, PreparedSignalEvaluator, ReportMetricObjective, @@ -136,7 +144,7 @@ ) from .adapters.nautilus import NautilusBacktestEngine from .core.types import BacktestResult -from .core.results import BacktestResultV2, OptionBacktestResult +from .core.results import BacktestResultV2, NativeAccountingArrays, NativeEventScoreResult, OptionBacktestResult from .core.execution_contract import ( EXECUTION_CONTRACT_REGISTRY, AmbiguityPolicy, @@ -441,7 +449,9 @@ "NautilusBacktestEngine", "NativeEventBackend", "NativeEventConfig", + "NativeAccountingArrays", "NativeActiveOrderSnapshot", + "NativeEventScoreResult", "NativeEventStrategyError", "NativeEventStrategyProtocol", "NativeFillEvent", @@ -464,6 +474,7 @@ "PortfolioRebalancePolicy", "PortfolioSizingMode", "QuantBTEndpoint", + "PreparedNativeEventStrategyRunner", "QuantBTPreparedContext", "format_metrics_report", "CANONICAL_OPTION_CHAIN_COLUMNS", diff --git a/backends/native_event.py b/backends/native_event.py index c148dbe..2eaa0fe 100644 --- a/backends/native_event.py +++ b/backends/native_event.py @@ -1246,6 +1246,9 @@ def run_strategy( report_level: Optional[str] = None, audit_sink: Optional[str] = None, audit_sink_path: Optional[str] = None, + market_arrays: Optional[PreparedMarketArrays] = None, + opens_arr: Optional[np.ndarray] = None, + volumes_arr: Optional[np.ndarray] = None, ) -> BacktestResultV2: """ Run a reactive strategy against native-event v2 lifecycle semantics. @@ -1268,18 +1271,29 @@ def run_strategy( idx = validate_datetime(datetime_index) symbol_list = list(symbols) if symbols is not None else list(closes.keys()) - market_arrays = self.prepare_market_arrays( - datetime_index=idx, - closes=closes, - highs=highs, - lows=lows, - funding_rate=funding_rate, - symbols=symbol_list, - ) - open_dict = align_series(opens, symbol_list, idx, fallback=align_series(closes, symbol_list, idx)) - volume_dict = align_series(volumes, symbol_list, idx, fallback={s: pd.Series(0.0, index=idx) for s in symbol_list}) - opens_arr = np.ascontiguousarray(np.column_stack([open_dict[s].to_numpy(dtype=np.float64) for s in symbol_list])) - volumes_arr = np.ascontiguousarray(np.column_stack([volume_dict[s].to_numpy(dtype=np.float64) for s in symbol_list])) + if market_arrays is None: + market_arrays = self.prepare_market_arrays( + datetime_index=idx, + closes=closes, + highs=highs, + lows=lows, + funding_rate=funding_rate, + symbols=symbol_list, + ) + elif market_arrays.signature != self._market_signature(idx, symbol_list): + raise ValueError("prepared market arrays do not match datetime_index/symbols") + if opens_arr is None: + open_dict = align_series(opens, symbol_list, idx, fallback=align_series(closes, symbol_list, idx)) + opens_arr = np.ascontiguousarray(np.column_stack([open_dict[s].to_numpy(dtype=np.float64) for s in symbol_list])) + else: + opens_arr = np.ascontiguousarray(opens_arr, dtype=np.float64) + if volumes_arr is None: + volume_dict = align_series(volumes, symbol_list, idx, fallback={s: pd.Series(0.0, index=idx) for s in symbol_list}) + volumes_arr = np.ascontiguousarray(np.column_stack([volume_dict[s].to_numpy(dtype=np.float64) for s in symbol_list])) + else: + volumes_arr = np.ascontiguousarray(volumes_arr, dtype=np.float64) + if opens_arr.shape != market_arrays.closes.shape or volumes_arr.shape != market_arrays.closes.shape: + raise ValueError("prepared opens/volumes arrays must match market array shape") contract_sizes = self._per_symbol_array(contract_size, symbol_list, default=1.0) constraints = build_quantity_constraints( diff --git a/benchmarks/phase34b_native_event_prepared_score.json b/benchmarks/phase34b_native_event_prepared_score.json new file mode 100644 index 0000000..ee6b064 --- /dev/null +++ b/benchmarks/phase34b_native_event_prepared_score.json @@ -0,0 +1,12 @@ +{ + "metric_parity": true, + "peak_rss_mb": 337.92578125, + "prepared_endpoint_result_retained": false, + "prepared_score_seconds": 0.6344219469465315, + "prepared_scores": 12, + "public_audit_seconds": 1.7633190099149942, + "public_last_report_level": "audit", + "rows": 600, + "speedup": 2.779410482884201, + "trials": 12 +} diff --git a/benchmarks/phase34b_native_event_prepared_score.md b/benchmarks/phase34b_native_event_prepared_score.md new file mode 100644 index 0000000..89b7b61 --- /dev/null +++ b/benchmarks/phase34b_native_event_prepared_score.md @@ -0,0 +1,12 @@ +# Phase 34B Native Event Prepared Score Benchmark + +- Rows: `600` +- Trials: `12` +- Public audit seconds: `1.763319` +- Prepared score seconds: `0.634422` +- Speedup: `2.779x` +- Peak RSS MB: `337.926` +- Metric parity: `True` +- Prepared endpoint result retained: `False` + +Prepared score reuses market arrays and returns `NativeEventScoreResult` rather than storing full public artifacts on the endpoint. diff --git a/benchmarks/run_phase34b_native_event_prepared_score.py b/benchmarks/run_phase34b_native_event_prepared_score.py new file mode 100644 index 0000000..bb6e84f --- /dev/null +++ b/benchmarks/run_phase34b_native_event_prepared_score.py @@ -0,0 +1,173 @@ +from __future__ import annotations + +import argparse +import json +import resource +import time +from pathlib import Path + +import numpy as np +import pandas as pd + +from quantbt import QuantBTEndpoint +from quantbt.core.orders import OrderCommand +from quantbt.core.schema import OrderSide, OrderType, TimeInForce + + +def _rss_mb() -> float: + return float(resource.getrusage(resource.RUSAGE_SELF).ru_maxrss) / 1024.0 + + +def _bars(rows: int) -> pd.DataFrame: + idx = pd.date_range("2020-01-01", periods=rows, freq="1h", tz="UTC") + x = np.arange(rows, dtype=np.float64) + close = pd.Series(100.0 + np.sin(x / 11.0) * 2.0 + x * 0.001, index=idx) + return pd.DataFrame( + { + "open": close, + "high": close + 2.0, + "low": close - 2.0, + "close": close, + "volume": 1_000.0, + }, + index=idx, + ) + + +class TimedStrategy: + def __init__(self, entry_mod: int, hold: int, qty: float): + self.entry_mod = int(entry_mod) + self.hold = int(hold) + self.qty = float(qty) + self.open_bar = -1 + + def on_bar_close(self, context): + symbol = context.symbols[0] + if context.positions[symbol] == 0.0 and context.bar_index % self.entry_mod == 0: + self.open_bar = int(context.bar_index) + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=symbol, + side=OrderSide.BUY, + order_type=OrderType.MARKET, + qty=self.qty, + tif=TimeInForce.IOC, + order_id=f"entry-{context.bar_index}", + ) + ] + if context.positions[symbol] > 0.0 and self.open_bar >= 0 and context.bar_index - self.open_bar >= self.hold: + self.open_bar = -1 + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=symbol, + side=OrderSide.SELL, + order_type=OrderType.MARKET, + qty=abs(context.positions[symbol]), + tif=TimeInForce.IOC, + reduce_only=True, + order_id=f"exit-{context.bar_index}", + ) + ] + return [] + + +def _params(trials: int): + return [ + { + "entry_mod": 5 + (i % 7), + "hold": 2 + (i % 5), + "qty": 0.1 + (i % 4) * 0.05, + } + for i in range(trials) + ] + + +def _metrics_subset(report: dict) -> dict: + return { + "sharpe": report["sharpe"], + "max_drawdown_pct": report["max_drawdown_pct"], + "profit_factor": report["profit_factor"], + "num_trades": report["num_trades"], + "final_equity": report["final_equity"], + "liquidated": report["liquidated"], + } + + +def run(rows: int, trials: int) -> dict: + df = _bars(rows) + params = _params(trials) + public_endpoint = QuantBTEndpoint.native_event_strategy( + initial_capital=50_000, + leverage=10, + use_funding=False, + fee_rate=0.0002, + report_level="audit", + ) + start = time.perf_counter() + public_reports = [] + for param in params: + result = public_endpoint.simulate(data=df, strategy=TimedStrategy(**param), symbols=["BTC"]) + public_reports.append(_metrics_subset(result.full_report(scope="full"))) + public_seconds = time.perf_counter() - start + + prepared_endpoint = QuantBTEndpoint.native_event_strategy( + initial_capital=50_000, + leverage=10, + use_funding=False, + fee_rate=0.0002, + report_level="audit", + ) + prepared = prepared_endpoint.prepare_native_event_strategy(data=df, symbols=["BTC"]) + start = time.perf_counter() + score_reports = [] + for param in params: + score = prepared.score(TimedStrategy(**param)) + score_reports.append(_metrics_subset(score.metrics)) + prepared_seconds = time.perf_counter() - start + + parity = public_reports == score_reports + return { + "rows": int(rows), + "trials": int(trials), + "public_audit_seconds": float(public_seconds), + "prepared_score_seconds": float(prepared_seconds), + "speedup": float(public_seconds / prepared_seconds) if prepared_seconds > 0.0 else np.inf, + "peak_rss_mb": float(_rss_mb()), + "metric_parity": bool(parity), + "prepared_scores": int(prepared.metadata["scores"]), + "public_last_report_level": public_endpoint.result.metadata["report_level"], + "prepared_endpoint_result_retained": prepared_endpoint.result is not None, + } + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--rows", type=int, default=1_000) + parser.add_argument("--trials", type=int, default=20) + parser.add_argument("--json-out", default="benchmarks/phase34b_native_event_prepared_score.json") + parser.add_argument("--md-out", default="benchmarks/phase34b_native_event_prepared_score.md") + args = parser.parse_args() + payload = run(rows=args.rows, trials=args.trials) + Path(args.json_out).write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n") + lines = [ + "# Phase 34B Native Event Prepared Score Benchmark", + "", + f"- Rows: `{payload['rows']}`", + f"- Trials: `{payload['trials']}`", + f"- Public audit seconds: `{payload['public_audit_seconds']:.6f}`", + f"- Prepared score seconds: `{payload['prepared_score_seconds']:.6f}`", + f"- Speedup: `{payload['speedup']:.3f}x`", + f"- Peak RSS MB: `{payload['peak_rss_mb']:.3f}`", + f"- Metric parity: `{payload['metric_parity']}`", + f"- Prepared endpoint result retained: `{payload['prepared_endpoint_result_retained']}`", + "", + "Prepared score reuses market arrays and returns `NativeEventScoreResult` rather than storing full public artifacts on the endpoint.", + ] + Path(args.md_out).write_text("\n".join(lines) + "\n") + print(json.dumps(payload, indent=2, sort_keys=True)) + + +if __name__ == "__main__": + main() diff --git a/core/__init__.py b/core/__init__.py index e921a68..7d25106 100644 --- a/core/__init__.py +++ b/core/__init__.py @@ -2,7 +2,7 @@ from .event import _engine_event_v1 from .vectorized import _engine_units_v2 from .types import BacktestResult -from .results import BacktestResultV2 +from .results import BacktestResultV2, NativeAccountingArrays, NativeEventScoreResult from .execution_contract import ( EXECUTION_CONTRACT_REGISTRY, AmbiguityPolicy, @@ -159,6 +159,8 @@ "_engine_portfolio", "BacktestResult", "BacktestResultV2", + "NativeAccountingArrays", + "NativeEventScoreResult", "BracketOrderSpec", "AccountConfig", "AlphaExecutionClassification", diff --git a/core/results.py b/core/results.py index 22cebf2..93e322a 100644 --- a/core/results.py +++ b/core/results.py @@ -7,7 +7,7 @@ from __future__ import annotations from dataclasses import dataclass, field -from typing import Dict, List, Sequence +from typing import Dict, List, Mapping, Sequence import numpy as np import pandas as pd @@ -120,6 +120,107 @@ def to_legacy(self) -> BacktestResult: ) +@dataclass(frozen=True) +class NativeAccountingArrays: + timestamps: np.ndarray + equity: np.ndarray + returns: np.ndarray + positions: np.ndarray + fees: np.ndarray + funding: np.ndarray + initial_margin: np.ndarray + maintenance_margin: np.ndarray + symbols: tuple[str, ...] + initial_capital: float + leverage: float = 1.0 + liquidated: bool = False + liquidation_bar: int = -1 + + @classmethod + def from_result(cls, result: BacktestResultV2) -> "NativeAccountingArrays": + position_cols = [f"Position_{symbol}" for symbol in result.symbols] + return cls( + timestamps=result.equity.index.view("int64").copy(), + equity=result.equity.to_numpy(dtype=np.float64, copy=True), + returns=result.returns.to_numpy(dtype=np.float64, copy=True), + positions=result.positions[position_cols].to_numpy(dtype=np.float64, copy=True), + fees=result.fees.to_numpy(dtype=np.float64, copy=True), + funding=result.funding.to_numpy(dtype=np.float64, copy=True), + initial_margin=result.margin.get("initial_margin", pd.Series(0.0, index=result.equity.index)).to_numpy( + dtype=np.float64, + copy=True, + ), + maintenance_margin=result.margin.get( + "maintenance_margin", + pd.Series(0.0, index=result.equity.index), + ).to_numpy(dtype=np.float64, copy=True), + symbols=tuple(result.symbols), + initial_capital=float(result.initial_capital), + leverage=float(result.leverage), + liquidated=bool(result.liquidated), + liquidation_bar=int(result.liquidation_bar), + ) + + @property + def datetime_index(self) -> pd.DatetimeIndex: + return pd.DatetimeIndex(self.timestamps) + + +@dataclass(frozen=True) +class NativeEventScoreResult: + accounting: NativeAccountingArrays + final_positions: np.ndarray + fill_count: int + rejection_count: int + cancellation_count: int + liquidated: bool + liquidation_bar: int + metrics: Mapping[str, float] + metadata: Mapping[str, object] = field(default_factory=dict) + + @property + def equity(self) -> np.ndarray: + return self.accounting.equity + + @property + def returns(self) -> np.ndarray: + return self.accounting.returns + + @property + def positions(self) -> np.ndarray: + return self.accounting.positions + + @property + def fees(self) -> np.ndarray: + return self.accounting.fees + + @property + def funding(self) -> np.ndarray: + return self.accounting.funding + + @property + def initial_margin(self) -> np.ndarray: + return self.accounting.initial_margin + + @property + def maintenance_margin(self) -> np.ndarray: + return self.accounting.maintenance_margin + + def full_report(self, trading_days: int = 365) -> Dict: + from ..metrics.performance import compute_performance_metrics + + return compute_performance_metrics( + timestamps=self.accounting.datetime_index, + equity=self.accounting.equity, + returns=self.accounting.returns, + positions=self.accounting.positions, + symbols=self.accounting.symbols, + initial_capital=float(self.accounting.initial_capital), + liquidated=bool(self.liquidated), + trading_days=trading_days, + ) + + @dataclass class OptionBacktestResult(BacktestResultV2): """ diff --git a/docs/endpoint.md b/docs/endpoint.md index cd4264e..d76f78a 100644 --- a/docs/endpoint.md +++ b/docs/endpoint.md @@ -1047,6 +1047,38 @@ For reactive strategies, `report_level="minimal"` intentionally omits `emitted_command_count`. Use `report_level="audit"` when a replayable command tape is required for certification. +Prepared native-event scoring: + +```python +bt = QuantBTEndpoint.native_event_strategy( + initial_capital=20_000, + leverage=5, + fee_rate=0.0005, + report_level="audit", +) + +prepared = bt.prepare_native_event_strategy( + data=df, + symbols=["ETHUSDT"], +) + +score = prepared.score( + strategy=DynamicGridStrategy(params), + trading_days=365, +) + +audit = prepared.run( + strategy=DynamicGridStrategy(params), + report_level="audit", +) +``` + +`prepared.score(...)` returns `NativeEventScoreResult`: ndarray accounting +paths plus metrics, not a public `BacktestResultV2`. It does not update +`bt.result`, so Optuna/WFO loops do not retain the previous trial's full +artifact bundle. `prepared.run(...)` returns the normal public +`BacktestResultV2` and should be used for final audit/replay exports. + Scoped cancel-all: ```python diff --git a/endpoint.py b/endpoint.py index a52aa6c..09b7886 100644 --- a/endpoint.py +++ b/endpoint.py @@ -58,7 +58,7 @@ from .core.intrabar_kernel import FillReplayTape, run_fill_replay_kernel, run_intrabar_kernel, run_intrabar_session_kernel from .core.market_tape import PreparedMarketTape, prepare_market_tape from .core.orders import OrderCommand, OrderIntent, order_intents_to_lifecycle_commands -from .core.results import BacktestResultV2, OptionBacktestResult +from .core.results import BacktestResultV2, NativeAccountingArrays, NativeEventScoreResult, OptionBacktestResult from .core.schema import AccountConfig, BasketLegSpec, BasketSpec, ExecutionConfig, InstrumentSpec, OrderSide, OrderType, TimeInForce from .core.structured_orders import ( BracketOrderSpec, @@ -296,6 +296,133 @@ def run(self, intent: IntrabarIntentTape, *, report_level: Optional[str] = None) return self.endpoint.result +@dataclass(frozen=True) +class PreparedNativeEventStrategyRunner: + """Prepared native-event reactive runner for repeated strategy scoring.""" + + endpoint: "QuantBTEndpoint" + idx: pd.DatetimeIndex + symbols: list + close_map: SeriesMap + high_map: SeriesMap + low_map: SeriesMap + opens_arr: np.ndarray + volumes_arr: np.ndarray + market_arrays: object + backend: NativeEventBackend + profile_metadata: Dict + runs: int = 0 + scores: int = 0 + + def run(self, strategy, *, report_level: Optional[str] = None) -> BacktestResultV2: + """Run the prepared strategy and return the public BacktestResultV2.""" + if strategy is None: + raise ValueError("prepared native-event runner requires strategy=...") + config = self.endpoint.config + level = report_level or config.report_level + result = self.backend.run_strategy( + datetime_index=self.idx, + strategy=strategy, + closes=self.close_map, + highs=self.high_map, + lows=self.low_map, + opens=None, + volumes=None, + funding_rate=config.funding_rate, + contract_size=config.contract_size, + leverage=config.account.leverage, + fee_rate=config.v2_fee_rate, + symbols=self.symbols, + instruments=config.instruments, + qty_step=config.qty_step, + lot_size=config.lot_size, + slot_size=config.slot_size, + min_qty=config.min_qty, + min_notional=config.min_notional, + execution_mode=config.reactive_execution_mode, + report_level=level, + audit_sink=config.audit_sink, + audit_sink_path=config.audit_sink_path, + market_arrays=self.market_arrays, + opens_arr=self.opens_arr, + volumes_arr=self.volumes_arr, + ) + result.metadata.setdefault("prepared_native_event_strategy", self.metadata) + object.__setattr__(self, "runs", self.runs + 1) + self.endpoint._store_result(result) + return self.endpoint.result + + simulate = run + + def score(self, strategy, *, trading_days: int = 365) -> NativeEventScoreResult: + """ + Run the prepared strategy with score artifact retention. + + The returned object stores ndarray accounting arrays and scalar metrics; + it intentionally does not update `endpoint.result`. + """ + if strategy is None: + raise ValueError("prepared native-event score requires strategy=...") + config = self.endpoint.config + result = self.backend.run_strategy( + datetime_index=self.idx, + strategy=strategy, + closes=self.close_map, + highs=self.high_map, + lows=self.low_map, + opens=None, + volumes=None, + funding_rate=config.funding_rate, + contract_size=config.contract_size, + leverage=config.account.leverage, + fee_rate=config.v2_fee_rate, + symbols=self.symbols, + instruments=config.instruments, + qty_step=config.qty_step, + lot_size=config.lot_size, + slot_size=config.slot_size, + min_qty=config.min_qty, + min_notional=config.min_notional, + execution_mode=config.reactive_execution_mode, + report_level="score", + audit_sink="none", + market_arrays=self.market_arrays, + opens_arr=self.opens_arr, + volumes_arr=self.volumes_arr, + ) + accounting = NativeAccountingArrays.from_result(result) + counters = dict(result.metadata.get("lifecycle_counters") or {}) + score = NativeEventScoreResult( + accounting=accounting, + final_positions=accounting.positions[-1].copy(), + fill_count=int(counters.get("fill_count", 0)), + rejection_count=int(counters.get("rejected_count", 0)), + cancellation_count=int(counters.get("canceled_count", 0)), + liquidated=bool(result.liquidated), + liquidation_bar=int(result.liquidation_bar), + metrics={}, + metadata={ + "backend": "native_event", + "engine": "event_v2_reactive_score", + "report_level": "score", + "prepared_native_event_strategy": self.metadata, + "lifecycle_counters": counters, + "artifact_plan": result.metadata.get("artifact_plan"), + }, + ) + object.__setattr__(self, "scores", self.scores + 1) + return replace(score, metrics=score.full_report(trading_days=trading_days)) + + @property + def metadata(self) -> Dict[str, object]: + return { + **self.profile_metadata, + "runs": int(self.runs), + "scores": int(self.scores), + "market_signature": self.market_arrays.signature, + } + + class QuantBTEndpoint: """ Stable notebook/service facade for all QuantBT backtest modes. @@ -411,6 +538,98 @@ def prepare_intrabar( session_tape=session_tape, ) + def prepare_native_event_strategy( + self, + *, + data=None, + closes=None, + highs=None, + lows=None, + datetime_index=None, + symbols: Optional[Sequence[str]] = None, + ) -> PreparedNativeEventStrategyRunner: + """ + Prepare native-event reactive market state once for repeated scoring. + + Normal `native_event_strategy(...).simulate(...)` remains unchanged. + This helper is for WFO/Optuna/service loops where the same market tape + is replayed many times with different strategy parameters. + """ + config = self.config + if str(config.backend).lower().strip() not in {"native_event", "auto"}: + raise ValueError("prepare_native_event_strategy requires backend='native_event' or auto") + symbol_list = list(symbols or config.symbols or (closes.keys() if closes is not None else [])) + if data is not None and not isinstance(data, dict) and not symbol_list: + symbol_list = ["asset"] + if not symbol_list: + raise ValueError("prepare_native_event_strategy requires symbols") + if data is not None and not isinstance(data, dict): + if len(symbol_list) != 1: + raise ValueError("single DataFrame native-event preparation requires exactly one symbol") + frame = _standardize_frame(data, datetime_index=datetime_index) + symbol = symbol_list[0] + idx = frame.index + close_map = {symbol: frame["close"]} + high_map = {symbol: frame.get("high", frame["close"])} + low_map = {symbol: frame.get("low", frame["close"])} + opens_arr = np.ascontiguousarray(frame[["open"]].to_numpy(dtype=np.float64)) + volumes_arr = np.ascontiguousarray(frame[["volume"]].to_numpy(dtype=np.float64)) + else: + close_map, high_map, low_map, idx, symbol_list = _normalize_symbol_data( + data=data, + closes=closes, + highs=highs, + lows=lows, + datetime_index=datetime_index, + symbols=symbol_list, + ) + opens_arr, volumes_arr = _prepared_native_event_open_volume_arrays(data, idx, symbol_list, close_map) + backend = NativeEventBackend( + NativeEventConfig( + account=config.account, + execution=config.execution, + fee_rate=config.v2_fee_rate, + use_funding=bool(config.use_funding), + report_level=config.report_level, + audit_sink=config.audit_sink, + audit_sink_path=config.audit_sink_path, + ) + ) + market = backend.prepare_market_arrays( + datetime_index=idx, + closes=close_map, + highs=high_map, + lows=low_map, + funding_rate=config.funding_rate, + symbols=symbol_list, + ) + profile = { + "mode": config.mode, + "backend": "native_event", + "event_engine_version": "v2", + "reactive_execution_mode": config.reactive_execution_mode, + "account": asdict(config.account), + "execution": asdict(config.execution), + "fee_rate": config.v2_fee_rate, + "report_level": config.report_level, + "symbols": tuple(symbol_list), + "bars": int(len(idx)), + "data_signature": market.signature, + } + return PreparedNativeEventStrategyRunner( + endpoint=self, + idx=idx, + symbols=list(symbol_list), + close_map=close_map, + high_map=high_map, + low_map=low_map, + opens_arr=opens_arr, + volumes_arr=volumes_arr, + market_arrays=market, + backend=backend, + profile_metadata=profile, + ) + @classmethod def pct_equity(cls, **kwargs) -> "QuantBTEndpoint": """ @@ -1789,9 +2008,6 @@ def _intrabar_execution_kwargs(self, symbol: str) -> Dict: slot_size=self.config.slot_size, min_qty=self.config.min_qty, min_notional=self.config.min_notional, - report_level=self.config.report_level, - audit_sink=self.config.audit_sink, - audit_sink_path=self.config.audit_sink_path, ) sizing_mode = IntrabarSizingMode(str(self.config.metadata.get("intrabar_sizing_mode", IntrabarSizingMode.UNITS.value))) fixed_notional = float(self.config.metadata.get("fixed_notional", self.config.alloc_per_trade if not isinstance(self.config.alloc_per_trade, dict) else self.config.alloc_per_trade.get(symbol, 0.0))) @@ -3765,6 +3981,24 @@ def _frames_from_symbol_maps(close_map, high_map, low_map, symbols) -> FrameMap: return frames +def _prepared_native_event_open_volume_arrays(data, idx: pd.DatetimeIndex, symbols, close_map) -> tuple[np.ndarray, np.ndarray]: + open_cols = [] + volume_cols = [] + for symbol in symbols: + close = close_map[symbol] + if isinstance(data, dict) and symbol in data and isinstance(data[symbol], pd.DataFrame): + frame = _standardize_frame(data[symbol], datetime_index=None) + open_cols.append(_align_series(frame.get("open", frame["close"]), idx).to_numpy(dtype=np.float64)) + volume_cols.append(_align_series(frame.get("volume", pd.Series(0.0, index=frame.index)), idx).to_numpy(dtype=np.float64)) + else: + open_cols.append(close.to_numpy(dtype=np.float64)) + volume_cols.append(np.zeros(len(idx), dtype=np.float64)) + return ( + np.ascontiguousarray(np.column_stack(open_cols), dtype=np.float64), + np.ascontiguousarray(np.column_stack(volume_cols), dtype=np.float64), + ) + + def _empty_nautilus_preflight_result(data, symbols, account: AccountConfig, metadata: Dict) -> BacktestResultV2: symbol_list = list(symbols) if not symbol_list: diff --git a/metrics/performance.py b/metrics/performance.py index 410f0b1..da7c377 100644 --- a/metrics/performance.py +++ b/metrics/performance.py @@ -10,7 +10,7 @@ from __future__ import annotations -from typing import Tuple, Dict +from typing import Dict, Sequence, Tuple import numpy as np import pandas as pd @@ -269,34 +269,273 @@ def rolling_drawdown(result: BacktestResult) -> pd.Series: # ── full report dict ───────────────────────────────────────────────────────── -def full_report(result: BacktestResult, trading_days: int = 365) -> Dict: +def compute_performance_metrics( + *, + timestamps: Sequence, + equity: Sequence[float], + returns: Sequence[float], + positions, + symbols: Sequence[str], + initial_capital: float, + liquidated: bool = False, + trading_days: int = 365, +) -> Dict: """ - Returns an ordered dict of all key metrics. - Suitable for programmatic use; viz/tearsheet renders it. + Shared array-first metric contract. + + This intentionally mirrors `full_report()` semantics so lightweight + prepared/native-event score paths and public `BacktestResultV2` reports use + one metric implementation. """ - lh, sh = hitrate(result) - aw, al = avg_win_loss(result) - md, ad = drawdown_duration(result) + idx = pd.DatetimeIndex(timestamps) + equity_arr = np.asarray(equity, dtype=np.float64) + returns_arr = np.asarray(returns, dtype=np.float64) + pos_arr = np.asarray(positions, dtype=np.float64) + if pos_arr.ndim == 1: + pos_arr = pos_arr.reshape(-1, 1) + if len(equity_arr) == 0: + raise ValueError("equity path cannot be empty") + if len(returns_arr) != len(equity_arr): + raise ValueError("returns must have the same length as equity") + if pos_arr.shape[0] != len(equity_arr): + raise ValueError("positions must have the same number of rows as equity") + + stats_returns = _array_returns_for_stats(idx, equity_arr, returns_arr) + annual_periods = _array_annualization_periods(idx, stats_returns, trading_days) + elapsed_years = _array_elapsed_years(idx, equity_arr, trading_days) + drawdown = _array_drawdown(equity_arr) + max_dd = float(np.nanmax(drawdown)) if len(drawdown) else 0.0 + avg_dd = float(np.nanmean(drawdown[drawdown > 0.0])) if np.any(drawdown > 0.0) else 0.0 + max_dd_duration, avg_dd_duration = _array_drawdown_duration_days(idx, equity_arr) + + final_equity = float(equity_arr[-1]) + total_ret = (final_equity - float(initial_capital)) / float(initial_capital) + cagr_value = _array_cagr(equity_arr, total_ret, elapsed_years) + sharpe_value = _array_sharpe(stats_returns, annual_periods) + sortino_value = _array_sortino(stats_returns, annual_periods) + omega_value = _array_omega(stats_returns) + pf_value = _array_profit_factor(stats_returns) + long_hr, short_hr = _array_hitrate(returns_arr, pos_arr) + avg_win, avg_loss = _array_avg_win_loss(stats_returns) + hr = (long_hr + short_hr) / 200.0 + expectancy_value = hr * avg_win + (1.0 - hr) * avg_loss return { - "initial_capital": result.initial_capital, - "final_equity": float(result.equity.iloc[-1]), - "total_return_pct": float(total_return(result) * 100), - "cagr_pct": float(cagr(result, trading_days) * 100), - "sharpe": float(sharpe(result, trading_days)), - "sortino": float(sortino(result, trading_days)), - "calmar": float(calmar(result, trading_days)), - "omega": float(omega(result)), - "max_drawdown_pct": float(max_drawdown_pct(result)), - "avg_drawdown_pct": float(avg_drawdown(result) * 100), - "max_dd_duration_days": md, - "avg_dd_duration_days": ad, - "profit_factor": float(profit_factor(result)), - "long_hitrate_pct": float(lh), - "short_hitrate_pct": float(sh), - "avg_win_pct": float(aw), - "avg_loss_pct": float(al), - "expectancy_pct": float(expectancy(result)), - "num_trades": int(number_of_trades(result)), - "liquidated": result.liquidated, + "initial_capital": float(initial_capital), + "final_equity": final_equity, + "total_return_pct": float(total_ret * 100.0), + "cagr_pct": float(cagr_value * 100.0), + "sharpe": float(sharpe_value), + "sortino": float(sortino_value), + "calmar": float(cagr_value / max_dd) if max_dd > 0.0 else 0.0, + "omega": float(omega_value), + "max_drawdown_pct": float(max_dd * 100.0), + "avg_drawdown_pct": float(avg_dd * 100.0), + "max_dd_duration_days": int(max_dd_duration), + "avg_dd_duration_days": int(avg_dd_duration), + "profit_factor": float(pf_value), + "long_hitrate_pct": float(long_hr), + "short_hitrate_pct": float(short_hr), + "avg_win_pct": float(avg_win), + "avg_loss_pct": float(avg_loss), + "expectancy_pct": float(expectancy_value), + "num_trades": int(_array_number_of_trades(pos_arr)), + "liquidated": bool(liquidated), } + + +def _array_finite_returns(values: np.ndarray) -> np.ndarray: + arr = np.asarray(values, dtype=np.float64) + return arr[np.isfinite(arr)] + + +def _array_daily_equity(idx: pd.DatetimeIndex, equity: np.ndarray) -> np.ndarray: + if len(equity) == 0: + return np.empty(0, dtype=np.float64) + if len(idx) != len(equity): + return np.asarray(equity, dtype=np.float64) + day_ns = 86_400_000_000_000 + days = idx.view("int64") // day_ns + if len(days) == 0: + return np.empty(0, dtype=np.float64) + change = np.flatnonzero(days[1:] != days[:-1]) + last_idx = np.concatenate((change, np.array([len(days) - 1], dtype=np.int64))) + return np.asarray(equity, dtype=np.float64)[last_idx] + + +def _array_returns_for_stats(idx: pd.DatetimeIndex, equity: np.ndarray, returns: np.ndarray) -> np.ndarray: + daily_equity = _array_daily_equity(idx, equity) + if len(daily_equity) >= 2: + base = daily_equity[:-1] + daily_returns = np.divide( + daily_equity[1:] - base, + base, + out=np.zeros(len(base), dtype=np.float64), + where=base != 0.0, + ) + daily_returns = _array_finite_returns(daily_returns) + if len(daily_returns) > 0: + return daily_returns + bar = _array_finite_returns(returns) + if len(bar) > 0: + return bar + if len(equity) < 2: + return np.zeros(1, dtype=np.float64) + base = equity[:-1] + out = np.divide(equity[1:] - base, base, out=np.zeros(len(base), dtype=np.float64), where=base != 0.0) + return _array_finite_returns(out) + + +def _array_annualization_periods(idx: pd.DatetimeIndex, stats_returns: np.ndarray, trading_days: int) -> float: + daily_equity_returns = len(stats_returns) > 0 + if daily_equity_returns and len(idx) >= 2: + day_ns = 86_400_000_000_000 + if len(np.unique(idx.view("int64") // day_ns)) >= 2: + return float(trading_days) + if len(idx) >= 2: + ns = idx.view("int64") + deltas = np.diff(ns).astype(np.float64) / 1_000_000_000.0 + deltas = deltas[deltas > 0.0] + if len(deltas) > 0: + median_seconds = float(np.median(deltas)) + if median_seconds > 0.0: + return float(365.25 * 24 * 60 * 60 / median_seconds) + return float(trading_days) + + +def _array_elapsed_years(idx: pd.DatetimeIndex, equity: np.ndarray, trading_days: int) -> float: + if len(equity) < 2: + return 0.0 + if len(idx) >= 2: + elapsed_days = (idx[-1] - idx[0]).total_seconds() / 86_400.0 + if elapsed_days > 0.0: + return float(elapsed_days / 365.25) + daily_equity = _array_daily_equity(idx, equity) + if len(daily_equity) >= 2: + return float(len(daily_equity) / float(trading_days)) + return float(len(equity) / float(trading_days)) + + +def _array_cagr(equity: np.ndarray, total_ret: float, years: float) -> float: + if len(equity) >= 2 and years > 0.0: + elapsed_days = years * 365.25 + if 0.0 < elapsed_days < 1.0: + return float(total_ret) + if years <= 0.0: + return 0.0 + growth = float(equity[-1] / equity[0]) + if growth <= 0.0: + return -1.0 + annual_log = np.log(growth) / years + if annual_log > 50.0: + return float(np.expm1(50.0)) + if annual_log < -50.0: + return float(np.expm1(-50.0)) + return float(np.expm1(annual_log)) + + +def _array_sharpe(r: np.ndarray, periods: float) -> float: + if len(r) < 2: + return 0.0 + sd = float(np.std(r, ddof=1)) + return float((np.mean(r) / sd) * np.sqrt(periods)) if sd > 0.0 else 0.0 + + +def _array_sortino(r: np.ndarray, periods: float, mar: float = 0.0) -> float: + downside = r[r < mar] - mar + dd = float(np.sqrt(np.mean(downside ** 2))) if len(downside) > 0 else 0.0 + mean = float(np.mean(r)) if len(r) > 0 else 0.0 + if dd == 0.0 and mean > mar: + return np.inf + return float((mean / dd) * np.sqrt(periods)) if dd > 0.0 else 0.0 + + +def _array_omega(r: np.ndarray, threshold: float = 0.0) -> float: + gain = float(np.sum(r[r > threshold] - threshold)) + loss = float(np.sum(threshold - r[r < threshold])) + return gain / loss if loss > 0.0 else np.inf + + +def _array_drawdown(equity: np.ndarray) -> np.ndarray: + peak = np.maximum.accumulate(equity) + return np.divide(peak - equity, peak, out=np.zeros_like(equity, dtype=np.float64), where=peak != 0.0) + + +def _array_drawdown_duration_days(idx: pd.DatetimeIndex, equity: np.ndarray) -> Tuple[int, int]: + daily_equity = _array_daily_equity(idx, equity) + if len(daily_equity) == 0: + return 0, 0 + peak = np.maximum.accumulate(daily_equity) + in_dd = peak != daily_equity + durations = [] + run = 0 + for value in in_dd: + if value: + run += 1 + elif run > 0: + durations.append(run) + run = 0 + if run > 0: + durations.append(run) + if not durations: + return 0, 0 + return int(max(durations)), int(np.mean(durations)) + + +def _array_hitrate(returns: np.ndarray, positions: np.ndarray) -> Tuple[float, float]: + long_hr = [] + short_hr = [] + for col in range(positions.shape[1]): + pos = positions[:, col] + long_mask = pos > 0.0 + short_mask = pos < 0.0 + long_total = int(np.sum(long_mask)) + short_total = int(np.sum(short_mask)) + long_wins = int(np.sum((returns > 0.0) & long_mask)) + short_wins = int(np.sum((returns > 0.0) & short_mask)) + long_hr.append(long_wins / long_total * 100.0 if long_total > 0 else 0.0) + short_hr.append(short_wins / short_total * 100.0 if short_total > 0 else 0.0) + return float(np.mean(long_hr)), float(np.mean(short_hr)) + + +def _array_number_of_trades(positions: np.ndarray) -> int: + if positions.size == 0: + return 0 + total = 0 + for col in range(positions.shape[1]): + pos = positions[:, col] + total += 1 + if len(pos) > 1: + total += int(np.sum(np.diff(pos) != 0.0)) + return int(total) + + +def _array_profit_factor(r: np.ndarray) -> float: + gains = float(np.sum(r[r > 0.0])) + loss = float(abs(np.sum(r[r < 0.0]))) + return gains / loss if loss > 0.0 else np.inf + + +def _array_avg_win_loss(r: np.ndarray) -> Tuple[float, float]: + wins = r[r > 0.0] + losses = r[r < 0.0] + win = float(np.mean(wins) * 100.0) if len(wins) > 0 else 0.0 + loss = float(np.mean(losses) * 100.0) if len(losses) > 0 else 0.0 + return win, loss + +def full_report(result: BacktestResult, trading_days: int = 365) -> Dict: + """ + Returns an ordered dict of all key metrics. + Suitable for programmatic use; viz/tearsheet renders it. + """ + positions = result.positions[[f"Position_{sym}" for sym in result.symbols]].to_numpy(dtype=np.float64) + return compute_performance_metrics( + timestamps=result.equity.index, + equity=result.equity.to_numpy(dtype=np.float64), + returns=result.returns.to_numpy(dtype=np.float64), + positions=positions, + symbols=result.symbols, + initial_capital=float(result.initial_capital), + liquidated=bool(result.liquidated), + trading_days=trading_days, + ) diff --git a/optimization/__init__.py b/optimization/__init__.py index 399b0b9..1184068 100644 --- a/optimization/__init__.py +++ b/optimization/__init__.py @@ -14,6 +14,7 @@ OptionPackageGenericEvaluator, OptionTrialOutput, PreparedIntrabarEvaluator, + PreparedNativeEventStrategyEvaluator, PreparedPortfolioEvaluator, PreparedSignalEvaluator, ) @@ -62,6 +63,7 @@ "OptimizationTrialRecord", "OptunaOptimizer", "PreparedIntrabarEvaluator", + "PreparedNativeEventStrategyEvaluator", "PreparedPortfolioEvaluator", "PreparedSignalEvaluator", "ReportMetricObjective", diff --git a/optimization/evaluators/__init__.py b/optimization/evaluators/__init__.py index d69a266..b3aff2c 100644 --- a/optimization/evaluators/__init__.py +++ b/optimization/evaluators/__init__.py @@ -12,6 +12,7 @@ from .generic import GenericEndpointEvaluator from .grid_dca import GridDCAGenericEvaluator, GridDCATrialOutput from .intrabar import PreparedIntrabarEvaluator +from .native_event import PreparedNativeEventStrategyEvaluator from .options import OptionPackageGenericEvaluator, OptionTrialOutput from .portfolio import PreparedPortfolioEvaluator from .signal import PreparedSignalEvaluator @@ -25,6 +26,7 @@ "OptionPackageGenericEvaluator", "OptionTrialOutput", "PreparedIntrabarEvaluator", + "PreparedNativeEventStrategyEvaluator", "PreparedPortfolioEvaluator", "PreparedSignalEvaluator", ] diff --git a/optimization/evaluators/native_event.py b/optimization/evaluators/native_event.py new file mode 100644 index 0000000..b494ec5 --- /dev/null +++ b/optimization/evaluators/native_event.py @@ -0,0 +1,32 @@ +"""Prepared native-event strategy evaluator.""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Any, Callable, Mapping + +from ..result import ObjectiveResult +from .generic import ObjectiveBuilder + + +@dataclass +class PreparedNativeEventStrategyEvaluator: + """Evaluate reactive native-event strategies through a prepared runner.""" + + runner: Any + strategy_factory: Callable[[Mapping[str, Any]], Any] + objective_builder: ObjectiveBuilder + trading_days: int = 365 + + last_result: Any = field(default=None, init=False) + last_strategy: Any = field(default=None, init=False) + + def evaluate(self, params: Mapping[str, Any]) -> ObjectiveResult: + strategy = self.strategy_factory(params) + result = self.runner.score(strategy, trading_days=self.trading_days) + objective = self.objective_builder(result, params) + if not isinstance(objective, ObjectiveResult): + raise TypeError("objective_builder must return ObjectiveResult") + self.last_strategy = strategy + self.last_result = result + return objective diff --git a/tests/test_phase34b_native_event_prepared_score.py b/tests/test_phase34b_native_event_prepared_score.py new file mode 100644 index 0000000..4e52dde --- /dev/null +++ b/tests/test_phase34b_native_event_prepared_score.py @@ -0,0 +1,143 @@ +from __future__ import annotations + +import numpy as np +import pandas as pd + +from quantbt import QuantBTEndpoint +from quantbt.core.orders import OrderCommand +from quantbt.core.schema import OrderSide, OrderType, TimeInForce +from quantbt.optimization import ObjectiveResult, PreparedNativeEventStrategyEvaluator + + +def _bars(n: int = 16) -> pd.DataFrame: + idx = pd.date_range("2024-01-01", periods=n, freq="1h", tz="UTC") + close = pd.Series(100.0 + np.sin(np.arange(n) / 2.0) * 2.0, index=idx) + return pd.DataFrame( + { + "open": close, + "high": close + 2.0, + "low": close - 2.0, + "close": close, + "volume": 1_000.0, + }, + index=idx, + ) + + +class TwoTradeStrategy: + def __init__(self, entry_bar: int = 0, exit_bar: int = 5, qty: float = 1.0): + self.entry_bar = int(entry_bar) + self.exit_bar = int(exit_bar) + self.qty = float(qty) + + def on_bar_close(self, context): + if context.bar_index == self.entry_bar: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=context.symbols[0], + side=OrderSide.BUY, + order_type=OrderType.MARKET, + qty=self.qty, + tif=TimeInForce.IOC, + order_id=f"entry-{self.entry_bar}", + ) + ] + if context.bar_index == self.exit_bar: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=context.symbols[0], + side=OrderSide.SELL, + order_type=OrderType.MARKET, + qty=self.qty, + tif=TimeInForce.IOC, + reduce_only=True, + order_id=f"exit-{self.exit_bar}", + ) + ] + return [] + + +def test_prepared_native_event_score_matches_public_audit_metrics_exactly(): + df = _bars() + endpoint = QuantBTEndpoint.native_event_strategy( + initial_capital=10_000, + leverage=10, + use_funding=False, + fee_rate=0.0002, + report_level="audit", + ) + prepared = endpoint.prepare_native_event_strategy(data=df, symbols=["BTC"]) + + score = prepared.score(TwoTradeStrategy(entry_bar=0, exit_bar=5), trading_days=365) + assert endpoint.result is None + audit = prepared.run(TwoTradeStrategy(entry_bar=0, exit_bar=5), report_level="audit") + + np.testing.assert_array_equal(score.equity, audit.equity.to_numpy(dtype=np.float64)) + np.testing.assert_array_equal(score.returns, audit.returns.to_numpy(dtype=np.float64)) + np.testing.assert_array_equal(score.positions, audit.positions[["Position_BTC"]].to_numpy(dtype=np.float64)) + np.testing.assert_array_equal(score.fees, audit.fees.to_numpy(dtype=np.float64)) + np.testing.assert_array_equal(score.funding, audit.funding.to_numpy(dtype=np.float64)) + np.testing.assert_array_equal(score.initial_margin, audit.margin["initial_margin"].to_numpy(dtype=np.float64)) + np.testing.assert_array_equal(score.maintenance_margin, audit.margin["maintenance_margin"].to_numpy(dtype=np.float64)) + + full_metrics = audit.full_report(trading_days=365, scope="full") + for key in ( + "sharpe", + "max_drawdown_pct", + "profit_factor", + "num_trades", + "final_equity", + "total_return_pct", + "liquidated", + ): + assert score.metrics[key] == full_metrics[key] + + assert score.metadata["report_level"] == "score" + assert score.fill_count == audit.metadata["lifecycle_counters"]["fill_count"] + assert score.rejection_count == audit.metadata["lifecycle_counters"]["rejected_count"] + assert prepared.metadata["scores"] == 1 + assert prepared.metadata["runs"] == 1 + + +def test_prepared_native_event_score_reuses_market_arrays_and_keeps_endpoint_result_light(): + df = _bars(24) + endpoint = QuantBTEndpoint.native_event_strategy(initial_capital=10_000, leverage=10, use_funding=False) + prepared = endpoint.prepare_native_event_strategy(data=df, symbols=["BTC"]) + signature = prepared.market_arrays.signature + + first = prepared.score(TwoTradeStrategy(entry_bar=0, exit_bar=4)) + second = prepared.score(TwoTradeStrategy(entry_bar=2, exit_bar=8)) + + assert prepared.market_arrays.signature == signature + assert prepared.metadata["scores"] == 2 + assert endpoint.result is None + assert not hasattr(first, "fills") + assert not hasattr(second, "orders") + assert first.metadata["prepared_native_event_strategy"]["market_signature"] == signature + assert second.metadata["prepared_native_event_strategy"]["market_signature"] == signature + + +def test_prepared_native_event_strategy_evaluator_uses_score_result_contract(): + df = _bars() + endpoint = QuantBTEndpoint.native_event_strategy(initial_capital=10_000, leverage=10, use_funding=False) + prepared = endpoint.prepare_native_event_strategy(data=df, symbols=["BTC"]) + + def strategy_factory(params): + return TwoTradeStrategy(entry_bar=int(params["entry_bar"]), exit_bar=int(params["exit_bar"])) + + def objective_builder(result, params): + report = result.full_report() + return ObjectiveResult(values=(float(report["sharpe"]),), metrics=report, metadata={"params": dict(params)}) + + evaluator = PreparedNativeEventStrategyEvaluator( + runner=prepared, + strategy_factory=strategy_factory, + objective_builder=objective_builder, + ) + objective = evaluator.evaluate({"entry_bar": 0, "exit_bar": 5}) + + assert isinstance(objective, ObjectiveResult) + assert evaluator.last_result.metadata["engine"] == "event_v2_reactive_score" + assert prepared.metadata["scores"] == 1 diff --git a/upgrade/implement.md b/upgrade/implement.md index 43df844..1a01964 100644 --- a/upgrade/implement.md +++ b/upgrade/implement.md @@ -6703,6 +6703,8 @@ Benchmark interpretation: ### Phase 34B - Prepared Native Event Score Path +Status: implemented on `feat/30-native-event-lifecycle`. + Scope: - Add prepared native-event strategy runner: @@ -6736,6 +6738,57 @@ Acceptance: - Optimizers can use the prepared score path without changing public endpoint behavior. +Implemented: + +- Added `NativeAccountingArrays` as the canonical ndarray accounting payload + extracted from native-event public results. +- Added `NativeEventScoreResult`: + - ndarray equity/returns/positions/fees/funding/margin views; + - lifecycle counters; + - scalar metrics; + - no public fills/orders artifact bundle. +- Added shared array-first performance metric function: + `metrics.performance.compute_performance_metrics(...)`. +- `BacktestResultV2.full_report()` and `NativeEventScoreResult.full_report()` + now use the same metric implementation through `metrics.performance`. +- Added `QuantBTEndpoint.prepare_native_event_strategy(...)`. +- Added `PreparedNativeEventStrategyRunner`: + - prepares market arrays once; + - reuses OHLC/funding/open/volume arrays; + - `.score(strategy)` returns `NativeEventScoreResult` and does not store + `endpoint.result`; + - `.run(strategy, report_level=...)` returns public `BacktestResultV2`. +- Added `PreparedNativeEventStrategyEvaluator` for the optimization framework. +- Exported the new score/result/evaluator APIs from top-level/core/optimization + namespaces. +- Updated endpoint docs with prepared native-event scoring examples. + +Validation: + +```bash +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q tests/test_phase34b_native_event_prepared_score.py +# 3 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run python benchmarks/run_phase34b_native_event_prepared_score.py --rows 600 --trials 12 +# metric_parity: true +# public_audit_seconds: 1.763319 +# prepared_score_seconds: 0.634422 +# speedup: 2.779x +# prepared_endpoint_result_retained: false + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q +# 544 passed, 1 skipped +``` + +Scope note: + +- Phase 34B still uses the existing reactive session plus static replay kernel + as the accounting source of truth. It prunes artifacts and reuses prepared + market arrays, but it is not yet the single-pass stateful kernel. +- Fully eliminating transient pandas public-result construction from score + execution belongs to Phase 34C, where the stateful kernel can emit + `NativeAccountingArrays` directly. + ### Phase 34C - Single-Pass Stateful Native Event Kernel Scope: From 2cb953a7b864e388bd4e6cd493c49456285802bf Mon Sep 17 00:00:00 2001 From: BobbyAxerol Date: Wed, 29 Jul 2026 09:43:14 +0000 Subject: [PATCH 3/4] Add native event single-pass reactive score path --- backends/native_event.py | 386 ++++++++++++++++-- .../phase34c_native_event_single_pass.json | 11 + .../phase34c_native_event_single_pass.md | 13 + .../run_phase34c_native_event_single_pass.py | 176 ++++++++ docs/endpoint.md | 25 +- endpoint.py | 10 + engines.py | 4 + .../test_phase34c_native_event_single_pass.py | 137 +++++++ upgrade/implement.md | 68 +++ 9 files changed, 798 insertions(+), 32 deletions(-) create mode 100644 benchmarks/phase34c_native_event_single_pass.json create mode 100644 benchmarks/phase34c_native_event_single_pass.md create mode 100644 benchmarks/run_phase34c_native_event_single_pass.py create mode 100644 tests/test_phase34c_native_event_single_pass.py diff --git a/backends/native_event.py b/backends/native_event.py index 2eaa0fe..d71f1cb 100644 --- a/backends/native_event.py +++ b/backends/native_event.py @@ -130,6 +130,7 @@ class NativeEventConfig: report_level: str = "audit" audit_sink: str = "memory" audit_sink_path: Optional[str] = None + reactive_kernel_mode: str = "replay_certified" def __post_init__(self) -> None: if isinstance(self.fee_rate, dict): @@ -139,6 +140,7 @@ def __post_init__(self) -> None: raise ValueError("fee_rate must be >= 0") object.__setattr__(self, "report_level", _normalize_native_event_report_level(self.report_level)) object.__setattr__(self, "audit_sink", _normalize_native_event_audit_sink(self.audit_sink)) + object.__setattr__(self, "reactive_kernel_mode", _normalize_reactive_kernel_mode(self.reactive_kernel_mode)) @dataclass(frozen=True) @@ -233,6 +235,15 @@ def _normalize_native_event_audit_sink(audit_sink: str) -> str: return sink +def _normalize_reactive_kernel_mode(reactive_kernel_mode: str) -> str: + mode = str(reactive_kernel_mode or "replay_certified").lower().strip() + aliases = {"replay": "replay_certified", "certified": "replay_certified", "stateful": "single_pass"} + mode = aliases.get(mode, mode) + if mode not in {"replay_certified", "single_pass"}: + raise ValueError("reactive_kernel_mode must be replay_certified or single_pass") + return mode + + def _native_event_artifact_plan(report_level: str) -> NativeEventArtifactPlan: level = _normalize_native_event_report_level(report_level) if level == "score": @@ -366,6 +377,18 @@ def __init__( self.processed_bar = -1 self.last_initial_margin = 0.0 self.last_maintenance_margin = 0.0 + n_bars = len(idx) + n_syms = len(symbols) + self.equity_path = np.zeros(n_bars, dtype=np.float64) + self.pos_path = np.zeros((n_bars, n_syms), dtype=np.float64) + self.fee_path = np.zeros(n_bars, dtype=np.float64) + self.turnover_path = np.zeros(n_bars, dtype=np.float64) + self.funding_path = np.zeros(n_bars, dtype=np.float64) + self.initial_margin_path = np.zeros(n_bars, dtype=np.float64) + self.maintenance_margin_path = np.zeros(n_bars, dtype=np.float64) + self.rejected_bar = np.zeros(n_bars, dtype=np.int64) + self.canceled_bar = np.zeros(n_bars, dtype=np.int64) + self._record_bar(0) def schedule(self, bar: int, commands: Sequence[OrderCommand]) -> None: if not commands or bar >= len(self.idx): @@ -413,6 +436,7 @@ def context(self, bar: int) -> NativeStrategyContext: def _process_single_bar(self, bar: int) -> None: if self.liquidated: + self._record_bar(bar) return if bar > 0: for s in range(len(self.symbols)): @@ -425,6 +449,7 @@ def _process_single_bar(self, bar: int) -> None: ) if bar > 0 and self._liquidated_intrabar(bar): self._liquidate(bar, LIQ_INTRABAR) + self._record_bar(bar) return if bar > 0 and self.use_funding and self.market_arrays.is_funding_bar[bar]: funding_cost = 0.0 @@ -438,10 +463,12 @@ def _process_single_bar(self, bar: int) -> None: * self.market_arrays.funding[bar, s] ) self.equity -= funding_cost + self.funding_path[bar] += funding_cost if bar > 0: _, close_mm = self._close_margin(bar) if close_mm > 0.0 and self.equity <= close_mm: self._liquidate(bar, LIQ_AFTER_FUNDING) + self._record_bar(bar) return self._expire_orders(bar) @@ -452,6 +479,18 @@ def _process_single_bar(self, bar: int) -> None: _, close_mm = self._close_margin(bar) if close_mm > 0.0 and self.equity <= close_mm: self._liquidate(bar, LIQ_AFTER_ORDER) + self._record_bar(bar) + + def _record_bar(self, bar: int) -> None: + if bar < 0 or bar >= len(self.idx): + return + init_margin, maint_margin = self._close_margin(bar) + self.equity_path[bar] = float(self.equity) + self.pos_path[bar, :] = self.current_pos + self.initial_margin_path[bar] = float(init_margin) + self.maintenance_margin_path[bar] = float(maint_margin) + self.last_initial_margin = float(init_margin) + self.last_maintenance_margin = float(maint_margin) def _apply_command(self, bar: int, command: OrderCommand) -> None: action = command.action @@ -561,6 +600,8 @@ def _match_orders(self, bar: int) -> None: self.equity += delta * (close - float(exec_price)) * cs - fee_cost self.current_pos[state.symbol_col] += delta + self.fee_path[bar] += fee_cost + self.turnover_path[bar] += trade_notional state.status = ORDER_STATUS_FILLED state.active = False state.waiting_parent = False @@ -633,6 +674,7 @@ def _cancel_state( state.active = False state.waiting_parent = False state.status = ORDER_STATUS_CANCELED + self.canceled_bar[bar] += 1 self._event( bar, command, @@ -652,6 +694,8 @@ def _event( target_order_id: Optional[str] = None, related_order_id: Optional[str] = None, ) -> None: + if event_name == "reject": + self.rejected_bar[bar] += 1 self.events_by_bar.setdefault(bar, []).append( NativeOrderEvent( timestamp=self.idx[bar], @@ -1243,6 +1287,7 @@ def run_strategy( min_notional: Optional[Union[float, Dict[str, float]]] = None, execution_mode: str = "fast", command_effective_phase: str = "next_bar", + reactive_kernel_mode: Optional[str] = None, report_level: Optional[str] = None, audit_sink: Optional[str] = None, audit_sink_path: Optional[str] = None, @@ -1265,6 +1310,9 @@ def run_strategy( execution_mode = str(execution_mode).lower().strip() if execution_mode not in {"fast", "audit"}: raise ValueError("execution_mode must be 'fast' or 'audit'") + kernel_mode = _normalize_reactive_kernel_mode( + self.config.reactive_kernel_mode if reactive_kernel_mode is None else reactive_kernel_mode + ) requested_report_level = self.config.report_level if report_level is None else report_level level = _normalize_native_event_report_level(requested_report_level) plan = _native_event_artifact_plan(level) @@ -1386,53 +1434,76 @@ def run_strategy( emitted.extend(scheduled) ignored_commands_after_end += ignored - final_result = self.run_order_commands( - datetime_index=idx, - commands=tuple(emitted), - closes=closes, - highs=highs, - lows=lows, - funding_rate=funding_rate, - contract_size=contract_size, - leverage=leverage, - fee_rate=fee_rate, - symbols=symbol_list, - market_arrays=market_arrays, - instruments=instruments, - qty_step=qty_step, - lot_size=lot_size, - slot_size=slot_size, - min_qty=min_qty, - min_notional=min_notional, - report_level=level, - audit_sink=audit_sink, - audit_sink_path=audit_sink_path, - ) + replay_required = kernel_mode == "replay_certified" or level in {"standard", "audit"} or execution_mode == "audit" + replay_result = None + if replay_required: + replay_result = self.run_order_commands( + datetime_index=idx, + commands=tuple(emitted), + closes=closes, + highs=highs, + lows=lows, + funding_rate=funding_rate, + contract_size=contract_size, + leverage=leverage, + fee_rate=fee_rate, + symbols=symbol_list, + market_arrays=market_arrays, + instruments=instruments, + qty_step=qty_step, + lot_size=lot_size, + slot_size=slot_size, + min_qty=min_qty, + min_notional=min_notional, + report_level=level, + audit_sink=audit_sink, + audit_sink_path=audit_sink_path, + ) + if kernel_mode == "replay_certified": + final_result = replay_result + engine_name = "event_v2_reactive_incremental" + else: + if replay_result is not None: + self._assert_reactive_session_replay_parity(session, replay_result) + final_result = self._reactive_session_result( + session=session, + symbol_list=symbol_list, + market_arrays=market_arrays, + leverages=leverages, + report_level=level, + plan=plan, + replay_result=replay_result, + audit_sink=audit_sink, + audit_sink_path=audit_sink_path, + ) + engine_name = "event_v2_reactive_single_pass" final_result.metadata.update( { - "engine": "event_v2_reactive_incremental", + "engine": engine_name, "reactive_execution_mode": execution_mode, + "reactive_kernel_mode": kernel_mode, "command_effective_phase": "next_bar", "emitted_command_tape": tuple(emitted) if plan.keep_command_tape else (), "emitted_command_tape_retained": bool(plan.keep_command_tape), "emitted_command_count": len(emitted), "ignored_commands_after_end": int(ignored_commands_after_end), "strategy_callback_count": int(callback_count), - "static_replay_available": True, + "static_replay_available": bool(replay_result is not None), + "reactive_static_replay_count": int(replay_result is not None), "reactive_context_builder": "incremental_session_v1", "reactive_incremental_compile_replays": 0, "reactive_session_liquidated": bool(session.liquidated), "reactive_session_liquidation_bar": int(session.liquidation_bar), } ) - if execution_mode == "audit": + if execution_mode == "audit" and replay_result is not None: replay_last_pos = { - symbol: float(final_result.positions[f"Position_{symbol}"].iloc[-1]) + symbol: float(replay_result.positions[f"Position_{symbol}"].iloc[-1]) for symbol in symbol_list } session_last_pos = {symbol: float(last_context.positions[symbol]) for symbol in symbol_list} final_result.metadata["reactive_audit"] = { - "final_equity_diff": float(abs(float(final_result.equity.iloc[-1]) - float(last_context.equity))), + "final_equity_diff": float(abs(float(replay_result.equity.iloc[-1]) - float(last_context.equity))), "final_position_diff": { symbol: float(abs(replay_last_pos.get(symbol, 0.0) - session_last_pos.get(symbol, 0.0))) for symbol in symbol_list @@ -1780,6 +1851,267 @@ def _apply_command_quantity_constraints( out.append(command) return tuple(out), {"changed_count": changed, "dropped_count": len(dropped), "dropped_orders": dropped} + def _reactive_session_result( + self, + *, + session: _NativeEventReactiveSession, + symbol_list: List[str], + market_arrays: PreparedMarketArrays, + leverages: np.ndarray, + report_level: str, + plan: NativeEventArtifactPlan, + replay_result: Optional[BacktestResultV2], + audit_sink: Optional[str], + audit_sink_path: Optional[str], + ) -> BacktestResultV2: + idx = session.idx + equity = pd.Series(session.equity_path.copy(), index=idx, name="equity") + returns = equity.pct_change().replace([np.inf, -np.inf], np.nan).fillna(0.0) + positions = pd.DataFrame( + {f"Position_{symbol}": session.pos_path[:, j].copy() for j, symbol in enumerate(symbol_list)}, + index=idx, + ) + closes = pd.DataFrame( + {f"Close_{symbol}": market_arrays.closes[:, j].copy() for j, symbol in enumerate(symbol_list)}, + index=idx, + ) + margin = pd.DataFrame( + { + "initial_margin": session.initial_margin_path.copy(), + "maintenance_margin": session.maintenance_margin_path.copy(), + }, + index=idx, + ) + diagnostics = pd.DataFrame( + { + "turnover": session.turnover_path.copy(), + "rejected_orders": session.rejected_bar.copy(), + "canceled_orders": session.canceled_bar.copy(), + }, + index=idx, + ) + session_fills = self._fills_from_reactive_session(session) + fill_ledger = self._compact_fill_ledger_from_session(session, symbol_list) + lifecycle_counters = { + "fill_count": int(len(session_fills)), + "event_count": int(sum(len(events) for events in session.events_by_bar.values())), + "rejected_count": int(np.sum(session.rejected_bar)), + "canceled_count": int(np.sum(session.canceled_bar)), + "filled_command_count": int(len(session_fills)), + "pending_command_count": int(sum(1 for state in session.pending if session._is_pending(state))), + "expired_event_count": int( + sum(1 for events in session.events_by_bar.values() for event in events if event.event_name == "expire") + ), + } + command_report = pd.DataFrame() + order_events = pd.DataFrame() + active_orders = pd.DataFrame() + orders = () + fills = tuple(session_fills) if plan.materialize_python_objects else () + compact_command_ledger = None + compact_order_event_ledger = None + audit_artifacts = {} + if replay_result is not None: + command_report = replay_result.metadata.get("command_report", pd.DataFrame()) + order_events = replay_result.metadata.get("order_events", pd.DataFrame()) + active_orders = replay_result.metadata.get("active_orders", pd.DataFrame()) + orders = replay_result.orders if plan.materialize_python_objects else () + fills = replay_result.fills if plan.materialize_python_objects else () + compact_command_ledger = replay_result.metadata.get("compact_command_ledger") + compact_order_event_ledger = replay_result.metadata.get("compact_order_event_ledger") + audit_artifacts = replay_result.metadata.get("audit_artifacts", {}) + + metadata = { + "backend": "native_event", + "engine": "event_v2_reactive_single_pass", + "report_level": report_level, + "artifact_plan": asdict(plan), + "audit_sink": self.config.audit_sink if audit_sink is None else _normalize_native_event_audit_sink(audit_sink), + "audit_sink_path": self.config.audit_sink_path if audit_sink_path is None else audit_sink_path, + "audit_artifacts": audit_artifacts, + "fee_rate_oneway": self._fee_rate_metadata(session.fee_rates, symbol_list), + "slippage_bps": self.config.execution.slippage_bps, + "order_report": command_report, + "command_report": command_report, + "order_events": order_events, + "active_orders": active_orders, + "compact_fill_ledger": fill_ledger if plan.keep_fill_ledger else None, + "compact_command_ledger": compact_command_ledger if plan.keep_command_terminal_state else None, + "compact_order_event_ledger": compact_order_event_ledger if plan.keep_event_ledger else None, + "quantity_constraints": session.constraints.as_dict(), + "quantity_preflight": {"changed_count": 0, "dropped_count": 0, "dropped_orders": []}, + "initial_buying_power": self.config.account.initial_capital * float(np.mean(leverages)), + "liquidation_reason": int(session.liquidation_reason), + "lifecycle_counters": lifecycle_counters, + "single_pass_accounting_source": "reactive_session_state", + "single_pass_replay_certified": bool(replay_result is not None), + } + return BacktestResultV2( + equity=equity, + returns=returns, + positions=positions, + closes=closes, + symbols=symbol_list, + initial_capital=self.config.account.initial_capital, + leverage=float(np.mean(leverages)), + liquidated=bool(session.liquidated), + liquidation_bar=int(session.liquidation_bar), + orders=orders, + fills=fills, + fees=pd.Series(session.fee_path.copy(), index=idx, name="fees"), + funding=pd.Series(session.funding_path.copy(), index=idx, name="funding"), + margin=margin, + diagnostics=diagnostics, + metadata=metadata, + ) + + @staticmethod + def _fills_from_reactive_session(session: _NativeEventReactiveSession) -> tuple[Fill, ...]: + fills: list[Fill] = [] + for bar in sorted(session.fills_by_bar): + for fill in session.fills_by_bar[bar]: + fills.append( + Fill( + timestamp=fill.timestamp, + symbol=fill.symbol, + side=fill.side, + qty=float(fill.qty), + price=float(fill.price), + fee=float(fill.fee), + order_id=fill.order_id, + metadata={ + **dict(fill.metadata), + "tag": fill.tag, + "campaign_id": fill.campaign_id, + "cycle_id": fill.cycle_id, + "level_id": fill.level_id, + "parent_order_id": fill.parent_order_id, + "oco_group_id": fill.oco_group_id, + }, + ) + ) + return tuple(fills) + + @staticmethod + def _compact_fill_ledger_from_session( + session: _NativeEventReactiveSession, + symbol_list: List[str], + ) -> CompactFillLedger: + id_map: Dict[str, int] = {} + symbol_to_col = {symbol: j for j, symbol in enumerate(symbol_list)} + bars = [] + command_index = [] + original_index = [] + order_id_code = [] + symbol_code = [] + side = [] + qty = [] + price = [] + fee = [] + fill_index = 0 + for bar in sorted(session.fills_by_bar): + for fill in session.fills_by_bar[bar]: + code = -1 + if fill.order_id: + if fill.order_id not in id_map: + id_map[fill.order_id] = len(id_map) + code = id_map[fill.order_id] + bars.append(int(bar)) + command_index.append(fill_index) + original_index.append(-1) + order_id_code.append(code) + symbol_code.append(symbol_to_col.get(fill.symbol, -1)) + side.append(fill.side.sign) + qty.append(float(fill.qty)) + price.append(float(fill.price)) + fee.append(float(fill.fee)) + fill_index += 1 + return CompactFillLedger( + bar=np.asarray(bars, dtype=np.int64), + command_index=np.asarray(command_index, dtype=np.int64), + original_index=np.asarray(original_index, dtype=np.int64), + order_id_code=np.asarray(order_id_code, dtype=np.int64), + symbol_code=np.asarray(symbol_code, dtype=np.int64), + side=np.asarray(side, dtype=np.int64), + qty=np.asarray(qty, dtype=np.float64), + price=np.asarray(price, dtype=np.float64), + fee=np.asarray(fee, dtype=np.float64), + id_values=tuple(sorted(id_map, key=id_map.get)), + symbols=tuple(symbol_list), + ) + + @staticmethod + def _compact_fill_ledger_from_fills(fills: Sequence[Fill], symbol_list: List[str]) -> CompactFillLedger: + id_map: Dict[str, int] = {} + symbol_to_col = {symbol: j for j, symbol in enumerate(symbol_list)} + bars = [] + command_index = [] + original_index = [] + order_id_code = [] + symbol_code = [] + side = [] + qty = [] + price = [] + fee = [] + for n, fill in enumerate(fills): + code = -1 + if fill.order_id: + if fill.order_id not in id_map: + id_map[fill.order_id] = len(id_map) + code = id_map[fill.order_id] + bars.append(n) + command_index.append(n) + original_index.append(-1) + order_id_code.append(code) + symbol_code.append(symbol_to_col.get(fill.symbol, -1)) + side.append(fill.side.sign) + qty.append(float(fill.qty)) + price.append(float(fill.price)) + fee.append(float(fill.fee)) + return CompactFillLedger( + bar=np.asarray(bars, dtype=np.int64), + command_index=np.asarray(command_index, dtype=np.int64), + original_index=np.asarray(original_index, dtype=np.int64), + order_id_code=np.asarray(order_id_code, dtype=np.int64), + symbol_code=np.asarray(symbol_code, dtype=np.int64), + side=np.asarray(side, dtype=np.int64), + qty=np.asarray(qty, dtype=np.float64), + price=np.asarray(price, dtype=np.float64), + fee=np.asarray(fee, dtype=np.float64), + id_values=tuple(sorted(id_map, key=id_map.get)), + symbols=tuple(symbol_list), + ) + + @staticmethod + def _assert_reactive_session_replay_parity( + session: _NativeEventReactiveSession, + replay_result: BacktestResultV2, + *, + atol: float = 1e-9, + ) -> None: + checks = { + "equity": (session.equity_path, replay_result.equity.to_numpy(dtype=np.float64)), + "fees": (session.fee_path, replay_result.fees.to_numpy(dtype=np.float64)), + "funding": (session.funding_path, replay_result.funding.to_numpy(dtype=np.float64)), + "positions": ( + session.pos_path, + replay_result.positions[[f"Position_{symbol}" for symbol in replay_result.symbols]].to_numpy(dtype=np.float64), + ), + "initial_margin": (session.initial_margin_path, replay_result.margin["initial_margin"].to_numpy(dtype=np.float64)), + "maintenance_margin": ( + session.maintenance_margin_path, + replay_result.margin["maintenance_margin"].to_numpy(dtype=np.float64), + ), + } + for name, (left, right) in checks.items(): + if not np.allclose(left, right, rtol=0.0, atol=atol, equal_nan=True): + diff = float(np.nanmax(np.abs(left - right))) + raise AssertionError(f"reactive single-pass replay parity failed for {name}: max_diff={diff}") + if bool(session.liquidated) != bool(replay_result.liquidated): + raise AssertionError("reactive single-pass replay parity failed for liquidated flag") + if int(session.liquidation_bar) != int(replay_result.liquidation_bar): + raise AssertionError("reactive single-pass replay parity failed for liquidation_bar") + def _reactive_replay( self, *, diff --git a/benchmarks/phase34c_native_event_single_pass.json b/benchmarks/phase34c_native_event_single_pass.json new file mode 100644 index 0000000..fe6d149 --- /dev/null +++ b/benchmarks/phase34c_native_event_single_pass.json @@ -0,0 +1,11 @@ +{ + "accounting_parity": true, + "peak_rss_mb": 333.76953125, + "replay_certified_seconds": 1.4313148567453027, + "replay_certified_static_replays": 12, + "rows": 600, + "single_pass_seconds": 0.7509573502466083, + "single_pass_static_replays": 0, + "speedup": 1.9059868796480526, + "trials": 12 +} diff --git a/benchmarks/phase34c_native_event_single_pass.md b/benchmarks/phase34c_native_event_single_pass.md new file mode 100644 index 0000000..d40656b --- /dev/null +++ b/benchmarks/phase34c_native_event_single_pass.md @@ -0,0 +1,13 @@ +# Phase 34C Native Event Single-Pass Benchmark + +- Rows: `600` +- Trials: `12` +- Replay-certified seconds: `1.431315` +- Single-pass seconds: `0.750957` +- Speedup: `1.906x` +- Replay-certified static replays: `12` +- Single-pass static replays: `0` +- Accounting parity: `True` +- Peak RSS MB: `333.770` + +This benchmark isolates the Phase 34C mode switch: `single_pass` materializes accounting from the reactive session for minimal/score runs and skips the final static replay. diff --git a/benchmarks/run_phase34c_native_event_single_pass.py b/benchmarks/run_phase34c_native_event_single_pass.py new file mode 100644 index 0000000..ed7ea4e --- /dev/null +++ b/benchmarks/run_phase34c_native_event_single_pass.py @@ -0,0 +1,176 @@ +from __future__ import annotations + +import argparse +import json +import resource +import time +from pathlib import Path + +import numpy as np +import pandas as pd + +from quantbt import OrderCommand, QuantBTEndpoint +from quantbt.core.schema import OrderSide, OrderType, TimeInForce + + +def _rss_mb() -> float: + return float(resource.getrusage(resource.RUSAGE_SELF).ru_maxrss) / 1024.0 + + +def _bars(rows: int) -> pd.DataFrame: + idx = pd.date_range("2020-01-01", periods=rows, freq="1h", tz="UTC") + x = np.arange(rows, dtype=np.float64) + close = pd.Series(100.0 + np.sin(x / 9.0) * 3.0 + np.cos(x / 23.0) * 1.5, index=idx) + return pd.DataFrame( + { + "open": close.shift(1).fillna(close.iloc[0]), + "high": close + 2.5, + "low": close - 2.5, + "close": close, + "volume": 1_000.0 + (x % 50.0), + }, + index=idx, + ) + + +class CyclicStrategy: + def __init__(self, entry_mod: int, hold: int, qty: float): + self.entry_mod = int(entry_mod) + self.hold = int(hold) + self.qty = float(qty) + self.open_bar = -1 + + def on_bar_close(self, context): + symbol = context.symbols[0] + if context.positions[symbol] == 0.0 and context.bar_index % self.entry_mod == 0: + self.open_bar = int(context.bar_index) + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=symbol, + side=OrderSide.BUY, + order_type=OrderType.MARKET, + qty=self.qty, + tif=TimeInForce.IOC, + order_id=f"entry-{context.bar_index}", + ) + ] + if context.positions[symbol] > 0.0 and self.open_bar >= 0 and context.bar_index - self.open_bar >= self.hold: + self.open_bar = -1 + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=symbol, + side=OrderSide.SELL, + order_type=OrderType.MARKET, + qty=abs(context.positions[symbol]), + tif=TimeInForce.IOC, + reduce_only=True, + order_id=f"exit-{context.bar_index}", + ) + ] + return [] + + +def _params(trials: int): + return [ + { + "entry_mod": 4 + (i % 9), + "hold": 2 + (i % 6), + "qty": 0.1 + (i % 5) * 0.025, + } + for i in range(trials) + ] + + +def _accounting_tuple(result) -> tuple: + return ( + tuple(np.round(result.equity.to_numpy(dtype=np.float64), 12)), + tuple(np.round(result.returns.to_numpy(dtype=np.float64), 12)), + tuple(np.round(result.positions.to_numpy(dtype=np.float64).ravel(), 12)), + tuple(np.round(result.fees.to_numpy(dtype=np.float64), 12)), + tuple(np.round(result.funding.to_numpy(dtype=np.float64), 12)), + tuple(np.round(result.margin.to_numpy(dtype=np.float64).ravel(), 12)), + bool(result.liquidated), + int(result.liquidation_bar), + ) + + +def run(rows: int, trials: int) -> dict: + df = _bars(rows) + params = _params(trials) + kwargs = dict( + initial_capital=50_000, + leverage=10, + use_funding=False, + fee_rate=0.0002, + report_level="minimal", + ) + + replay_endpoint = QuantBTEndpoint.native_event_strategy(**kwargs, reactive_kernel_mode="replay_certified") + start = time.perf_counter() + replay_fingerprints = [] + replay_static_replays = 0 + for param in params: + result = replay_endpoint.simulate(data=df, strategy=CyclicStrategy(**param), symbols=["BTC"]) + replay_fingerprints.append(_accounting_tuple(result)) + replay_static_replays += int(result.metadata.get("reactive_static_replay_count", 0)) + replay_seconds = time.perf_counter() - start + + single_endpoint = QuantBTEndpoint.native_event_strategy(**kwargs, reactive_kernel_mode="single_pass") + start = time.perf_counter() + single_fingerprints = [] + single_static_replays = 0 + for param in params: + result = single_endpoint.simulate(data=df, strategy=CyclicStrategy(**param), symbols=["BTC"]) + single_fingerprints.append(_accounting_tuple(result)) + single_static_replays += int(result.metadata.get("reactive_static_replay_count", 0)) + single_seconds = time.perf_counter() - start + + return { + "rows": int(rows), + "trials": int(trials), + "replay_certified_seconds": float(replay_seconds), + "single_pass_seconds": float(single_seconds), + "speedup": float(replay_seconds / single_seconds) if single_seconds > 0.0 else np.inf, + "replay_certified_static_replays": int(replay_static_replays), + "single_pass_static_replays": int(single_static_replays), + "accounting_parity": bool(replay_fingerprints == single_fingerprints), + "peak_rss_mb": float(_rss_mb()), + } + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--rows", type=int, default=1_000) + parser.add_argument("--trials", type=int, default=20) + parser.add_argument("--json-out", default="benchmarks/phase34c_native_event_single_pass.json") + parser.add_argument("--md-out", default="benchmarks/phase34c_native_event_single_pass.md") + args = parser.parse_args() + payload = run(rows=args.rows, trials=args.trials) + json_path = Path(args.json_out) + md_path = Path(args.md_out) + json_path.parent.mkdir(parents=True, exist_ok=True) + md_path.parent.mkdir(parents=True, exist_ok=True) + json_path.write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n") + lines = [ + "# Phase 34C Native Event Single-Pass Benchmark", + "", + f"- Rows: `{payload['rows']}`", + f"- Trials: `{payload['trials']}`", + f"- Replay-certified seconds: `{payload['replay_certified_seconds']:.6f}`", + f"- Single-pass seconds: `{payload['single_pass_seconds']:.6f}`", + f"- Speedup: `{payload['speedup']:.3f}x`", + f"- Replay-certified static replays: `{payload['replay_certified_static_replays']}`", + f"- Single-pass static replays: `{payload['single_pass_static_replays']}`", + f"- Accounting parity: `{payload['accounting_parity']}`", + f"- Peak RSS MB: `{payload['peak_rss_mb']:.3f}`", + "", + "This benchmark isolates the Phase 34C mode switch: `single_pass` materializes accounting from the reactive session for minimal/score runs and skips the final static replay.", + ] + md_path.write_text("\n".join(lines) + "\n") + print(json.dumps(payload, indent=2, sort_keys=True)) + + +if __name__ == "__main__": + main() diff --git a/docs/endpoint.md b/docs/endpoint.md index d76f78a..415d6ab 100644 --- a/docs/endpoint.md +++ b/docs/endpoint.md @@ -1011,6 +1011,7 @@ bt = QuantBTEndpoint.native_event_strategy( leverage=5, fee_rate=0.0005, reactive_execution_mode="fast", + reactive_kernel_mode="replay_certified", # replay_certified | single_pass ) result = bt.simulate( @@ -1029,16 +1030,26 @@ replay = QuantBTEndpoint.native_event_lifecycle( Reactive timing is causal: commands returned by `on_bar_close(context_t)` are retimed to bar `t+1`, so they cannot fill inside the same OHLC bar that the -strategy just observed. Phase 30E uses an incremental callback session for -speed, then replays the emitted command tape once through the certified -event-v2 lifecycle kernel for final accounting, fills, margin, liquidation and -reports. +strategy just observed. + +`reactive_kernel_mode="replay_certified"` is the conservative default. It uses +the incremental callback session to build state and then runs one certified +static event-v2 replay for the final public result. Use it for stakeholder +reports, debugging, and migration validation. + +`reactive_kernel_mode="single_pass"` materializes accounting directly from the +incremental reactive session for `report_level="minimal"` and score paths, +skipping the final static replay. For `report_level="standard"`, +`report_level="audit"`, or `reactive_execution_mode="audit"`, QuantBT still +runs the replay oracle and asserts accounting parity before returning the +single-pass result. Reactive metadata: ```python result.metadata["reactive_context_builder"] # "incremental_session_v1" result.metadata["reactive_incremental_compile_replays"] # 0 +result.metadata["reactive_static_replay_count"] # 0 for single_pass minimal/score result.metadata["emitted_command_tape"] # replayable OrderCommand tape ``` @@ -1055,6 +1066,7 @@ bt = QuantBTEndpoint.native_event_strategy( leverage=5, fee_rate=0.0005, report_level="audit", + reactive_kernel_mode="replay_certified", ) prepared = bt.prepare_native_event_strategy( @@ -1076,7 +1088,10 @@ audit = prepared.run( `prepared.score(...)` returns `NativeEventScoreResult`: ndarray accounting paths plus metrics, not a public `BacktestResultV2`. It does not update `bt.result`, so Optuna/WFO loops do not retain the previous trial's full -artifact bundle. `prepared.run(...)` returns the normal public +artifact bundle. Phase 34C makes `prepared.score(...)` use the single-pass +reactive session accounting path and skip the final replay while maintaining +parity with `prepared.run(..., report_level="audit")`. `prepared.run(...)` +returns the normal public `BacktestResultV2` and should be used for final audit/replay exports. Scoped cancel-all: diff --git a/endpoint.py b/endpoint.py index 09b7886..3128a4c 100644 --- a/endpoint.py +++ b/endpoint.py @@ -193,6 +193,7 @@ class EndpointConfig: structured_order_spec: object = None event_engine_version: str = "v1" reactive_execution_mode: str = "fast" + reactive_kernel_mode: str = "replay_certified" symbols: Optional[Sequence[str]] = None dca_kwargs: Dict = field(default_factory=dict) nautilus_config: object = None @@ -340,6 +341,7 @@ def run(self, strategy, *, report_level: Optional[str] = None) -> BacktestResult min_qty=config.min_qty, min_notional=config.min_notional, execution_mode=config.reactive_execution_mode, + reactive_kernel_mode=config.reactive_kernel_mode, report_level=level, audit_sink=config.audit_sink, audit_sink_path=config.audit_sink_path, @@ -384,6 +386,7 @@ def score(self, strategy, *, trading_days: int = 365) -> NativeEventScoreResult: min_qty=config.min_qty, min_notional=config.min_notional, execution_mode=config.reactive_execution_mode, + reactive_kernel_mode="single_pass", report_level="score", audit_sink="none", market_arrays=self.market_arrays, @@ -408,6 +411,8 @@ def score(self, strategy, *, trading_days: int = 365) -> NativeEventScoreResult: "prepared_native_event_strategy": self.metadata, "lifecycle_counters": counters, "artifact_plan": result.metadata.get("artifact_plan"), + "reactive_kernel_mode": result.metadata.get("reactive_kernel_mode"), + "static_replay_available": result.metadata.get("static_replay_available"), }, ) object.__setattr__(self, "scores", self.scores + 1) @@ -593,6 +598,7 @@ def prepare_native_event_strategy( report_level=config.report_level, audit_sink=config.audit_sink, audit_sink_path=config.audit_sink_path, + reactive_kernel_mode=config.reactive_kernel_mode, ) ) market = backend.prepare_market_arrays( @@ -608,6 +614,7 @@ def prepare_native_event_strategy( "backend": "native_event", "event_engine_version": "v2", "reactive_execution_mode": config.reactive_execution_mode, + "reactive_kernel_mode": config.reactive_kernel_mode, "account": asdict(config.account), "execution": asdict(config.execution), "fee_rate": config.v2_fee_rate, @@ -2082,6 +2089,7 @@ def _run_single(self, data, signal, signal_col, datetime_index, symbols): report_level=self.config.report_level, audit_sink=self.config.audit_sink, audit_sink_path=self.config.audit_sink_path, + reactive_kernel_mode=self.config.reactive_kernel_mode, ) markers = _intrabar_marker_columns(frame) if backend == "native_vectorized" and markers: @@ -2126,6 +2134,7 @@ def _run_orders(self, data, orders, order_commands, datetime_index, symbols): report_level=self.config.report_level, audit_sink=self.config.audit_sink, audit_sink_path=self.config.audit_sink_path, + reactive_kernel_mode=self.config.reactive_kernel_mode, ) self._store_result(self.engine.result) return self.result @@ -2162,6 +2171,7 @@ def _run_native_event_strategy(self, data, strategy, datetime_index, symbols): report_level=self.config.report_level, audit_sink=self.config.audit_sink, audit_sink_path=self.config.audit_sink_path, + reactive_kernel_mode=self.config.reactive_kernel_mode, ) self._store_result(self.engine.result) return self.result diff --git a/engines.py b/engines.py index 44336c8..ab7d42d 100644 --- a/engines.py +++ b/engines.py @@ -69,6 +69,7 @@ def __init__( strategy=None, event_engine_version: str = "v1", reactive_execution_mode: str = "fast", + reactive_kernel_mode: str = "replay_certified", report_level: str = "audit", audit_sink: str = "memory", audit_sink_path: Optional[str] = None, @@ -112,6 +113,7 @@ def __init__( self.strategy = strategy self.event_engine_version = str(event_engine_version).lower().strip() self.reactive_execution_mode = str(reactive_execution_mode).lower().strip() + self.reactive_kernel_mode = str(reactive_kernel_mode).lower().strip() self.report_level = str(report_level) self.audit_sink = str(audit_sink) self.audit_sink_path = audit_sink_path @@ -214,6 +216,7 @@ def _run_native_event(self) -> BacktestResultV2: report_level=self.report_level, audit_sink=self.audit_sink, audit_sink_path=self.audit_sink_path, + reactive_kernel_mode=self.reactive_kernel_mode, ) ) @@ -244,6 +247,7 @@ def _run_native_event(self) -> BacktestResultV2: min_qty=self.min_qty, min_notional=self.min_notional, execution_mode=self.reactive_execution_mode, + reactive_kernel_mode=self.reactive_kernel_mode, report_level=self.report_level, audit_sink=self.audit_sink, audit_sink_path=self.audit_sink_path, diff --git a/tests/test_phase34c_native_event_single_pass.py b/tests/test_phase34c_native_event_single_pass.py new file mode 100644 index 0000000..8d41587 --- /dev/null +++ b/tests/test_phase34c_native_event_single_pass.py @@ -0,0 +1,137 @@ +from __future__ import annotations + +import numpy as np +import pandas as pd + +from quantbt import OrderCommand, QuantBTEndpoint +from quantbt.core.schema import OrderSide, OrderType, TimeInForce + + +def _bars(n: int = 18) -> pd.DataFrame: + idx = pd.date_range("2024-01-01", periods=n, freq="1h", tz="UTC") + close = pd.Series(100.0 + np.sin(np.arange(n) / 3.0) * 3.0 + np.arange(n) * 0.15, index=idx) + return pd.DataFrame( + { + "open": close.shift(1).fillna(close.iloc[0]), + "high": close + 2.5, + "low": close - 2.5, + "close": close, + "volume": 1_000.0 + np.arange(n), + }, + index=idx, + ) + + +class EnterExitStrategy: + def __init__(self, entry_bar: int = 0, exit_bar: int = 6, qty: float = 1.0): + self.entry_bar = int(entry_bar) + self.exit_bar = int(exit_bar) + self.qty = float(qty) + + def on_bar_close(self, context): + if context.bar_index == self.entry_bar: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=context.symbols[0], + side=OrderSide.BUY, + order_type=OrderType.MARKET, + qty=self.qty, + tif=TimeInForce.IOC, + order_id=f"entry-{self.entry_bar}", + ) + ] + if context.bar_index == self.exit_bar: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=context.symbols[0], + side=OrderSide.SELL, + order_type=OrderType.MARKET, + qty=self.qty, + tif=TimeInForce.IOC, + reduce_only=True, + order_id=f"exit-{self.exit_bar}", + ) + ] + return [] + + +def _assert_accounting_equal(left, right) -> None: + pd.testing.assert_series_equal(left.equity, right.equity) + pd.testing.assert_series_equal(left.returns, right.returns) + pd.testing.assert_frame_equal(left.positions, right.positions) + pd.testing.assert_series_equal(left.fees, right.fees) + pd.testing.assert_series_equal(left.funding, right.funding) + pd.testing.assert_frame_equal(left.margin, right.margin) + assert left.liquidated == right.liquidated + assert left.liquidation_bar == right.liquidation_bar + + +def test_single_pass_minimal_skips_static_replay_but_matches_replay_certified_accounting(): + df = _bars() + kwargs = dict(initial_capital=10_000, leverage=10, use_funding=False, fee_rate=0.0002, report_level="minimal") + + replay = QuantBTEndpoint.native_event_strategy( + **kwargs, + reactive_kernel_mode="replay_certified", + ).simulate(data=df, strategy=EnterExitStrategy(entry_bar=0, exit_bar=6), symbols=["BTC"]) + single = QuantBTEndpoint.native_event_strategy( + **kwargs, + reactive_kernel_mode="single_pass", + ).simulate(data=df, strategy=EnterExitStrategy(entry_bar=0, exit_bar=6), symbols=["BTC"]) + + _assert_accounting_equal(single, replay) + assert single.metadata["engine"] == "event_v2_reactive_single_pass" + assert single.metadata["reactive_kernel_mode"] == "single_pass" + assert single.metadata["static_replay_available"] is False + assert single.metadata["reactive_static_replay_count"] == 0 + assert single.metadata["reactive_incremental_compile_replays"] == 0 + assert single.metadata["emitted_command_tape"] == () + + +def test_single_pass_audit_uses_replay_oracle_and_keeps_fill_bar_ledger(): + df = _bars() + result = QuantBTEndpoint.native_event_strategy( + initial_capital=10_000, + leverage=10, + use_funding=False, + fee_rate=0.0002, + report_level="audit", + reactive_execution_mode="audit", + reactive_kernel_mode="single_pass", + ).simulate(data=df, strategy=EnterExitStrategy(entry_bar=0, exit_bar=6), symbols=["BTC"]) + + ledger = result.metadata["compact_fill_ledger"] + assert result.metadata["single_pass_replay_certified"] is True + assert result.metadata["static_replay_available"] is True + assert result.metadata["reactive_static_replay_count"] == 1 + assert result.metadata["command_report"].shape[0] == 2 + assert tuple(ledger.bar.tolist()) == (1, 7) + assert result.metadata["reactive_audit"]["final_equity_diff"] == 0.0 + assert result.metadata["reactive_audit"]["final_position_diff"]["BTC"] == 0.0 + + +def test_prepared_native_event_score_uses_single_pass_and_keeps_public_run_parity(): + df = _bars(24) + endpoint = QuantBTEndpoint.native_event_strategy( + initial_capital=10_000, + leverage=10, + use_funding=False, + fee_rate=0.0002, + report_level="audit", + ) + prepared = endpoint.prepare_native_event_strategy(data=df, symbols=["BTC"]) + + score = prepared.score(EnterExitStrategy(entry_bar=2, exit_bar=9), trading_days=365) + audit = prepared.run(EnterExitStrategy(entry_bar=2, exit_bar=9), report_level="audit") + + np.testing.assert_allclose(score.equity, audit.equity.to_numpy(dtype=np.float64), rtol=0.0, atol=1e-9) + np.testing.assert_allclose(score.returns, audit.returns.to_numpy(dtype=np.float64), rtol=0.0, atol=1e-12) + np.testing.assert_allclose(score.positions, audit.positions[["Position_BTC"]].to_numpy(dtype=np.float64), rtol=0.0, atol=1e-12) + np.testing.assert_allclose(score.fees, audit.fees.to_numpy(dtype=np.float64), rtol=0.0, atol=1e-12) + np.testing.assert_allclose(score.initial_margin, audit.margin["initial_margin"].to_numpy(dtype=np.float64), rtol=0.0, atol=1e-12) + assert score.metadata["reactive_kernel_mode"] == "single_pass" + assert score.metadata["static_replay_available"] is False + assert prepared.metadata["scores"] == 1 + assert prepared.metadata["runs"] == 1 diff --git a/upgrade/implement.md b/upgrade/implement.md index 1a01964..804c081 100644 --- a/upgrade/implement.md +++ b/upgrade/implement.md @@ -6829,6 +6829,74 @@ Acceptance: - Audit can still produce full trace and optional replay certification. - 500-trial prepared run does not grow RAM with completed-trial history. +Implemented: + +- Added `NativeEventConfig.reactive_kernel_mode` with + `replay_certified` and `single_pass`. +- Kept public compatibility default at `replay_certified`. +- Added single-pass result materialization from `_NativeEventReactiveSession` + state: + - equity path; + - returns; + - position matrix; + - fee/funding arrays; + - turnover, rejection, cancellation diagnostics; + - margin paths; + - liquidation flags; + - compact fill ledger with real bar indices. +- `single_pass` skips the final static replay for `report_level="minimal"` and + score runs. +- `single_pass` still runs replay oracle for `standard`, `audit`, and + `reactive_execution_mode="audit"`, then asserts exact accounting parity. +- Added metadata: + - `reactive_kernel_mode`; + - `static_replay_available`; + - `reactive_static_replay_count`; + - `single_pass_accounting_source`; + - `single_pass_replay_certified`. +- Updated `PreparedNativeEventStrategyRunner.score(...)` to use + `reactive_kernel_mode="single_pass"` automatically. +- Threaded `reactive_kernel_mode` through endpoint, prepared runner, and + `BacktestEngineV2`. +- Preserved legacy `reactive_incremental_compile_replays == 0` semantics: + this field counts replay/compile inside callback construction, not the final + optional certification replay. + +Validation: + +```bash +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q quantbt/tests/test_phase34c_native_event_single_pass.py +# 3 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q \ + quantbt/tests/test_phase30d_native_event_reactive_runner.py \ + quantbt/tests/test_phase30e_native_event_incremental_runner.py \ + quantbt/tests/test_phase34a_native_event_artifacts.py \ + quantbt/tests/test_phase34b_native_event_prepared_score.py \ + quantbt/tests/test_phase34c_native_event_single_pass.py +# 19 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run python3 \ + benchmarks/run_phase34c_native_event_single_pass.py --rows 600 --trials 12 +# accounting_parity: true +# replay_certified_seconds: 1.431315 +# single_pass_seconds: 0.750957 +# speedup: 1.906x +# replay_certified_static_replays: 12 +# single_pass_static_replays: 0 +``` + +Scope note: + +- Phase 34C completes the practical single-pass optimization contract for + reactive strategy minimal/score loops. +- The implementation intentionally keeps the Python reactive session as the + state source and uses the existing event-v2 replay kernel as the oracle for + audit/certification. +- A deeper future rewrite could move active-order state into a true low-level + Numba step kernel, but that is no longer required for current prepared + WFO/Optuna memory and replay-reduction goals. + ### Phase 34 Final Merge Gate - Public endpoints stay source-compatible. From 4cebae2d54cbd20d2e370b9ab8d4d432c3267c2c Mon Sep 17 00:00:00 2001 From: BobbyAxerol Date: Wed, 29 Jul 2026 11:19:23 +0000 Subject: [PATCH 4/4] Harden native event phase34 dev integration --- core/results.py | 5 +- quantbt_phase34_merge_gate.py | 113 ++++++++++++++++++ ...st_phase34b_native_event_prepared_score.py | 33 ++++- upgrade/implement.md | 40 +++++++ 4 files changed, 188 insertions(+), 3 deletions(-) create mode 100644 quantbt_phase34_merge_gate.py diff --git a/core/results.py b/core/results.py index 93e322a..ac71499 100644 --- a/core/results.py +++ b/core/results.py @@ -206,7 +206,10 @@ def initial_margin(self) -> np.ndarray: def maintenance_margin(self) -> np.ndarray: return self.accounting.maintenance_margin - def full_report(self, trading_days: int = 365) -> Dict: + def full_report(self, trading_days: int = 365, scope: str = "auto") -> Dict: + if str(scope).lower().strip() not in {"auto", "full"}: + raise ValueError("NativeEventScoreResult supports scope='auto' or scope='full'") + from ..metrics.performance import compute_performance_metrics return compute_performance_metrics( diff --git a/quantbt_phase34_merge_gate.py b/quantbt_phase34_merge_gate.py new file mode 100644 index 0000000..088e393 --- /dev/null +++ b/quantbt_phase34_merge_gate.py @@ -0,0 +1,113 @@ +from __future__ import annotations + +import inspect + +import numpy as np +import pandas as pd + +import quantbt +from quantbt import EndpointConfig, OrderCommand, PreparedNativeEventStrategyRunner, QuantBTEndpoint +from quantbt.core.schema import OrderSide, OrderType, TimeInForce +from quantbt.optimization import ObjectiveResult, PreparedNativeEventStrategyEvaluator, ReportMetricObjective + + +class _GateStrategy: + def on_bar_close(self, context): + symbol = context.symbols[0] + if context.bar_index == 0: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=symbol, + side=OrderSide.BUY, + order_type=OrderType.MARKET, + qty=1.0, + tif=TimeInForce.IOC, + order_id="entry", + ) + ] + if context.bar_index == 4 and context.positions[symbol] > 0.0: + return [ + OrderCommand( + timestamp=context.timestamp, + symbol=symbol, + side=OrderSide.SELL, + order_type=OrderType.MARKET, + qty=abs(context.positions[symbol]), + tif=TimeInForce.IOC, + reduce_only=True, + order_id="exit", + ) + ] + return [] + + +def _bars() -> pd.DataFrame: + idx = pd.date_range("2024-01-01", periods=12, freq="1h", tz="UTC") + close = pd.Series(100.0 + np.sin(np.arange(len(idx)) / 2.0), index=idx) + return pd.DataFrame( + { + "open": close, + "high": close + 2.0, + "low": close - 2.0, + "close": close, + "volume": 1_000.0, + }, + index=idx, + ) + + +def _assert(condition: bool, message: str) -> None: + if not condition: + raise AssertionError(message) + + +def main() -> None: + fields = EndpointConfig.__dataclass_fields__ + _assert( + all(name in fields for name in ("reactive_kernel_mode", "audit_sink", "audit_sink_path")), + "EndpointConfig missing Phase 34 fields", + ) + print("EndpointConfig fields: PASSED") + + _assert(hasattr(QuantBTEndpoint, "prepare_native_event_strategy"), "prepare_native_event_strategy missing") + _assert(hasattr(quantbt, "PreparedNativeEventStrategyRunner"), "PreparedNativeEventStrategyRunner missing") + _assert(quantbt.PreparedNativeEventStrategyRunner is PreparedNativeEventStrategyRunner, "Prepared runner export mismatch") + print("Prepared endpoint API: PASSED") + + endpoint = QuantBTEndpoint.native_event_strategy(initial_capital=10_000, leverage=10, use_funding=False) + prepared = endpoint.prepare_native_event_strategy(data=_bars(), symbols=["BTC"]) + _assert(isinstance(prepared, PreparedNativeEventStrategyRunner), "prepared runner type mismatch") + score = prepared.score(_GateStrategy()) + _assert(score.metadata["reactive_kernel_mode"] == "single_pass", "prepared score did not use single_pass") + print("prepared.score(): PASSED") + + signature = inspect.signature(score.full_report) + _assert("scope" in signature.parameters, "NativeEventScoreResult.full_report missing scope") + score.full_report(scope="auto") + print("Score/full_report signature: PASSED") + + def strategy_factory(_params): + return _GateStrategy() + + def objective_builder(result, params): + objective = ReportMetricObjective(value_metrics=("sharpe",), scope="auto") + return objective(result, params) + + evaluator = PreparedNativeEventStrategyEvaluator( + runner=prepared, + strategy_factory=strategy_factory, + objective_builder=objective_builder, + ) + objective = evaluator.evaluate({}) + _assert(isinstance(objective, ObjectiveResult), "Prepared evaluator did not return ObjectiveResult") + print("Prepared evaluator import: PASSED") + + objective = ReportMetricObjective(value_metrics=("sharpe",), scope="auto")(score, {}) + _assert(isinstance(objective, ObjectiveResult), "ReportMetricObjective(score) failed") + print("ReportMetricObjective(score): PASSED") + print("PHASE 34 MERGE GATE: PASSED") + + +if __name__ == "__main__": + main() diff --git a/tests/test_phase34b_native_event_prepared_score.py b/tests/test_phase34b_native_event_prepared_score.py index 4e52dde..9e81403 100644 --- a/tests/test_phase34b_native_event_prepared_score.py +++ b/tests/test_phase34b_native_event_prepared_score.py @@ -1,12 +1,15 @@ from __future__ import annotations +import inspect + import numpy as np import pandas as pd -from quantbt import QuantBTEndpoint +import quantbt +from quantbt import EndpointConfig, PreparedNativeEventStrategyRunner, QuantBTEndpoint from quantbt.core.orders import OrderCommand from quantbt.core.schema import OrderSide, OrderType, TimeInForce -from quantbt.optimization import ObjectiveResult, PreparedNativeEventStrategyEvaluator +from quantbt.optimization import ObjectiveResult, PreparedNativeEventStrategyEvaluator, ReportMetricObjective def _bars(n: int = 16) -> pd.DataFrame: @@ -141,3 +144,29 @@ def objective_builder(result, params): assert isinstance(objective, ObjectiveResult) assert evaluator.last_result.metadata["engine"] == "event_v2_reactive_score" assert prepared.metadata["scores"] == 1 + + +def test_public_native_event_phase34_contract_is_available_from_quantbt(): + fields = EndpointConfig.__dataclass_fields__ + + assert hasattr(quantbt.QuantBTEndpoint, "prepare_native_event_strategy") + assert hasattr(quantbt, "PreparedNativeEventStrategyRunner") + assert PreparedNativeEventStrategyRunner is quantbt.PreparedNativeEventStrategyRunner + assert "reactive_kernel_mode" in fields + assert "audit_sink" in fields + assert "audit_sink_path" in fields + + +def test_report_metric_objective_accepts_native_event_score_result_scope_contract(): + df = _bars() + endpoint = QuantBTEndpoint.native_event_strategy(initial_capital=10_000, leverage=10, use_funding=False) + prepared = endpoint.prepare_native_event_strategy(data=df, symbols=["BTC"]) + score_result = prepared.score(TwoTradeStrategy(entry_bar=0, exit_bar=5)) + + assert "scope" in inspect.signature(score_result.full_report).parameters + objective = ReportMetricObjective(value_metrics=("sharpe",), scope="auto") + result = objective(score_result, {"entry_bar": 0, "exit_bar": 5}) + + assert isinstance(result, ObjectiveResult) + assert result.values == (score_result.metrics["sharpe"],) + assert result.metrics["num_trades"] == score_result.metrics["num_trades"] diff --git a/upgrade/implement.md b/upgrade/implement.md index 804c081..3f2e7e2 100644 --- a/upgrade/implement.md +++ b/upgrade/implement.md @@ -6910,3 +6910,43 @@ Scope note: - Benchmark report records wall time, CPU time, peak RSS, Python heap peak, NumPy allocated bytes, object count, ledger bytes, command count, fill count, report construction time, and stage timings. + +Merge regression fix on `dev`: + +- Cherry-picked Phase 34A, 34B, and 34C onto `dev` after detecting that the + native-event public endpoint integration was missing from the research + branch. +- Restored public exports and endpoint contracts: + - `PreparedNativeEventStrategyRunner`; + - `QuantBTEndpoint.prepare_native_event_strategy(...)`; + - `EndpointConfig.reactive_kernel_mode`; + - `EndpointConfig.audit_sink`; + - `EndpointConfig.audit_sink_path`. +- Added `scope="auto"` compatibility to + `NativeEventScoreResult.full_report(...)`, matching the public result + contract expected by `ReportMetricObjective`. +- Added `quantbt_phase34_merge_gate.py` so future merges can verify public + Phase 34 integration directly. +- Added regression tests covering: + - public `quantbt` imports; + - endpoint Phase 34 fields; + - prepared native-event runner availability; + - `prepared.score(...)`; + - `PreparedNativeEventStrategyEvaluator`; + - `ReportMetricObjective(score_result)`. + +Validation on `dev`: + +```bash +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run python3 quantbt_phase34_merge_gate.py +# PHASE 34 MERGE GATE: PASSED + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q \ + quantbt/tests/test_phase34a_native_event_artifacts.py \ + quantbt/tests/test_phase34b_native_event_prepared_score.py \ + quantbt/tests/test_phase34c_native_event_single_pass.py +# 11 passed + +MPLCONFIGDIR=/tmp PYTHONPATH=/root/bobby/pool_alpha poetry run pytest -q quantbt/tests +# 549 passed, 1 skipped +```