""" ml/execution.py — Advanced execution optimization. Capabilities: 1. Spread spike avoidance (skip entry when spread > threshold) 2. Refined time-of-day filters with intraday session scoring 3. Volume-weighted entry timing 4. Microstructure filters (tick-level proxies) 5. Partial fill simulation with realistic fill rates 6. Latency modeling (configurable delay) 7. Pullback confirmation entries (wait for better price) """ from __future__ import annotations import logging import numpy as np import pandas as pd logger = logging.getLogger("trading_system.ml.execution") class ExecutionOptimizer: """Execution quality improvements. Filters out entries during: - First/last 15 minutes (spread spikes) - Low-liquidity windows (volume < 30% of session avg) - High-impact macro windows (configurable) - Microstructure anomalies (extreme range, thin volume) """ def __init__( self, skip_first_min: int = 15, skip_last_min: int = 120, # stop 2h before close (2:00 PM ET) min_volume_ratio: float = 0.30, max_spread_atr_ratio: float = 0.5, latency_bars: int = 1, partial_fill_rate: float = 0.95, pullback_atr_frac: float = 0.2, pullback_max_wait: int = 3, ): self.skip_first_min = skip_first_min self.skip_last_min = skip_last_min self.min_volume_ratio = min_volume_ratio self.max_spread_atr_ratio = max_spread_atr_ratio self.latency_bars = latency_bars self.partial_fill_rate = partial_fill_rate self.pullback_atr_frac = pullback_atr_frac self.pullback_max_wait = pullback_max_wait def compute_execution_mask( self, df: pd.DataFrame, atr: pd.Series, ) -> pd.Series: """Return boolean mask: True = OK to enter, False = skip. Args: df: OHLCV DataFrame with DatetimeIndex. atr: ATR series aligned to df. Returns: pd.Series[bool] indexed like df. """ ok = pd.Series(True, index=df.index) # Time-of-day filter if df.index.tzinfo is not None: et_index = df.index.tz_convert("US/Eastern") else: try: et_index = df.index.tz_localize("UTC").tz_convert("US/Eastern") except Exception: et_index = df.index et_minutes = et_index.hour * 60 + et_index.minute market_open = 9 * 60 + 30 # 9:30 ET market_close = 16 * 60 # 16:00 ET # Skip first N minutes ok &= et_minutes >= (market_open + self.skip_first_min) # Skip last N minutes ok &= et_minutes <= (market_close - self.skip_last_min) # Volume filter: skip low-liquidity bars volume = df["volume"].astype(float) vol_session_avg = volume.rolling(78).mean() # ~1 day on 5m bars vol_ratio = volume / vol_session_avg.replace(0, np.nan) ok &= vol_ratio.fillna(1.0) >= self.min_volume_ratio # Spread proxy: (high - low) / ATR bar_range = (df["high"].astype(float) - df["low"].astype(float)) spread_ratio = bar_range / atr.replace(0, np.nan) ok &= spread_ratio.fillna(0) <= self.max_spread_atr_ratio * 3 # Microstructure: skip if bar volume is extremely thin (< 100 shares) ok &= volume >= 100 return ok def compute_session_quality(self, df: pd.DataFrame) -> pd.Series: """Score each bar's execution quality (0-1) based on time-of-day pattern. Higher scores during mid-session liquid periods, lower near open/close. """ if df.index.tzinfo is not None: et_index = df.index.tz_convert("US/Eastern") else: try: et_index = df.index.tz_localize("UTC").tz_convert("US/Eastern") except Exception: et_index = df.index minutes = et_index.hour * 60 + et_index.minute market_open = 9 * 60 + 30 market_close = 16 * 60 # Minutes into session session_min = (minutes - market_open).values.astype(float) session_len = market_close - market_open # 390 # Parabolic quality curve: best in middle, worst at edges x = np.clip(session_min / session_len, 0, 1) quality = 4 * x * (1 - x) # peak 1.0 at midday return pd.Series(quality, index=df.index).clip(0, 1) def apply_latency(self, signal_idx: int, max_idx: int) -> int: """Delay signal execution by latency_bars to simulate real-world latency.""" return min(signal_idx + self.latency_bars, max_idx) def simulate_partial_fill( self, desired_shares: float, bar_volume: float, participation_rate: float = 0.02, ) -> float: """Simulate partial fills based on bar volume. Args: desired_shares: Requested position size in shares. bar_volume: Volume of the entry bar. participation_rate: Max fraction of bar volume we'll take. Returns: Filled shares (may be less than desired). """ max_from_volume = bar_volume * participation_rate fillable = min(desired_shares, max_from_volume) return fillable * self.partial_fill_rate def pullback_entry_price( self, side: str, signal_price: float, subsequent_lows: np.ndarray, subsequent_highs: np.ndarray, atr: float, ) -> tuple[float, int] | None: """Attempt a pullback entry within max_wait bars. For buys: wait for price to dip below signal_price - pullback_frac * ATR. For sells: wait for price to rise above signal_price + pullback_frac * ATR. Returns: (entry_price, bars_waited) or None if pullback never triggered. """ target_offset = self.pullback_atr_frac * atr max_bars = min(self.pullback_max_wait, len(subsequent_lows)) for i in range(max_bars): if side == "buy": if subsequent_lows[i] <= signal_price - target_offset: return signal_price - target_offset, i + 1 else: if subsequent_highs[i] >= signal_price + target_offset: return signal_price + target_offset, i + 1 return None # Pullback didn't happen, use market entry def adjust_entry_price( self, side: str, open_price: float, high: float, low: float, atr: float, ) -> float: """Apply realistic slippage model based on bar characteristics. For buys: entry slightly above open (adverse fill) For sells: entry slightly below open (adverse fill) Slippage scales with volatility. """ base_slip_pct = 0.02 / 100 bar_range = high - low vol_slip = 0.0 if atr > 0: vol_factor = min(bar_range / atr, 2.0) vol_slip = open_price * base_slip_pct * vol_factor slippage = open_price * base_slip_pct + vol_slip if side == "buy": return open_price + slippage else: return open_price - slippage