diff --git a/apps/paper_trader/backtest_sim.py b/apps/paper_trader/backtest_sim.py index 0401ff8..3878ab3 100644 --- a/apps/paper_trader/backtest_sim.py +++ b/apps/paper_trader/backtest_sim.py @@ -37,7 +37,10 @@ def run_backtest_session_sync( parking_preset: str | None = None, idle_alpha_preset: str | None = None, form4_sleeve_preset: str | None = None, + ownership_sleeve_preset: str | None = None, + risk_off_alpha_sleeve_preset: str | None = None, snapshot_id_override: str | None = None, + fixed_capital_sizing: bool = False, ) -> dict[str, Any]: """Run a single strategy using BacktestRunner (same as research backtester). @@ -47,6 +50,7 @@ def run_backtest_session_sync( from apps.backtester.run import ( BacktestRunner, _build_merged_snapshot_store, + _compute_max_effective_mhd, _extend_store_to_requested_window, load_manifest, resolve_config, @@ -65,12 +69,28 @@ def run_backtest_session_sync( if form4_sleeve_preset: config.form4_capture_sleeve_preset = form4_sleeve_preset config.apply_form4_capture_sleeve_preset() + if ownership_sleeve_preset: + config.ownership_capture_sleeve_preset = ownership_sleeve_preset + config.apply_ownership_capture_sleeve_preset() + if risk_off_alpha_sleeve_preset: + config.risk_off_alpha_sleeve_preset = risk_off_alpha_sleeve_preset + config.apply_risk_off_alpha_sleeve_preset() + if fixed_capital_sizing: + config.risk.fixed_capital_sizing = True # Use merged store (train+valid+test) to cover the full date range. store = _build_merged_snapshot_store(manifest, config, snapshot_dir_override=None) - # Slice to requested date range - store = store.slice_by_date_range(start_date, end_date) + # Slice to requested date range. + # When lookback_entry_enabled, extend slice start backward so pre-start + # events survive and BacktestRunner._collect_lookback_candidates() can find them. + if config.execution.lookback_entry_enabled: + max_mhd = _compute_max_effective_mhd(config) + import datetime as _dt + lookback_start = start_date - _dt.timedelta(days=max_mhd * 2) + store = store.slice_by_date_range(lookback_start, end_date) + else: + store = store.slice_by_date_range(start_date, end_date) store = _extend_store_to_requested_window( store=store, config=config, @@ -409,6 +429,8 @@ def run_backtest( console=None, parking_preset: str | None = None, idle_alpha_preset: str | None = None, + form4_sleeve_preset: str | None = None, + ownership_sleeve_preset: str | None = None, snapshot_id_override: str | None = None, auto_refresh: bool = True, ) -> list[dict[str, Any]]: @@ -479,6 +501,8 @@ def run_backtest( end_date=end_date, parking_preset=parking_preset, idle_alpha_preset=idle_alpha_preset, + form4_sleeve_preset=form4_sleeve_preset, + ownership_sleeve_preset=ownership_sleeve_preset, snapshot_id_override=snapshot_id_override, ) diff --git a/apps/paper_trader/engine.py b/apps/paper_trader/engine.py index 67bd09d..1aaedca 100644 --- a/apps/paper_trader/engine.py +++ b/apps/paper_trader/engine.py @@ -99,6 +99,12 @@ class PaperTradingEngine: if session.form4_sleeve_preset: self._config.form4_capture_sleeve_preset = session.form4_sleeve_preset self._config.apply_form4_capture_sleeve_preset() + if session.ownership_sleeve_preset: + self._config.ownership_capture_sleeve_preset = session.ownership_sleeve_preset + self._config.apply_ownership_capture_sleeve_preset() + if getattr(session, "risk_off_alpha_sleeve_preset", None): + self._config.risk_off_alpha_sleeve_preset = session.risk_off_alpha_sleeve_preset + self._config.apply_risk_off_alpha_sleeve_preset() # Shared attention filtering service (matches BacktestRunner) from libs.backtest.attention import AttentionFilterService @@ -119,6 +125,8 @@ class PaperTradingEngine: self._capital_bucket_specs.get(bucket_id, 0.0), float(allocation), ) + # Lookback entry: fired once per daemon session on the first run_next_open call + self._lookback_injected: bool = False def _get_candidate_capital_bucket_id(self, candidate: Candidate) -> str | None: return candidate.engine_capital_bucket_id @@ -1559,6 +1567,29 @@ class PaperTradingEngine: if self._config.risk.cash_parking_enabled: parking_sold_today = self._parking_check_and_sell(session_id, today) + # Lookback entry: on first run_next_open, pick up events from previous days + # that are still within their holding window (fires once per daemon session). + lookback_entries: list[dict[str, Any]] = [] + lookback_rejected: list[dict[str, Any]] = [] + if self._config.execution.lookback_entry_enabled and not self._lookback_injected: + self._lookback_injected = True + lookback_rows = await self._detector.get_candidates_for_lookback( + today, self._lookback_start_date(today), self._config + ) + if lookback_rows: + macro_data_lb = await self._fetch_macro(today) + lookback_entries, lookback_rejected = await self._process_entries( + today, lookback_rows, account, alpaca_positions, strategy_states, + session_st, macro_data_lb, self._broker.submit_market_buy, + entry_timing="next_open", + ) + # Refresh state after lookback entries + alpaca_positions = self._broker.list_positions() + strategy_states = { + ss.symbol: ss + for ss in self._state.get_open_strategy_states(session_id) + } + # Entry: all events with entry_date == today, across both conventions. # Mirrors BacktestRunner: next_open engines call get_candidates_for_date(date) # which returns ALL rows with execution_date==date regardless of entry_convention. @@ -1583,6 +1614,8 @@ class PaperTradingEngine: session_st, macro_data, self._broker.submit_market_buy, entry_timing="next_open", ) + entries = lookback_entries + entries + rejected = lookback_rejected + rejected # CASH PARKING: buy with remaining idle cash after entries if self._config.risk.cash_parking_enabled and not parking_sold_today: @@ -1673,6 +1706,12 @@ class PaperTradingEngine: # Shared helpers for phased execution # ------------------------------------------------------------------ # + def _lookback_start_date(self, today: dt.date) -> dt.date: + """Return the earliest date to search for lookback events (calendar-day buffer).""" + from apps.backtester.run import _compute_max_effective_mhd + max_mhd = _compute_max_effective_mhd(self._config) + return today - dt.timedelta(days=max_mhd * 2) + @staticmethod def _is_same_day_event(row: dict[str, Any]) -> bool: """event_date == reaction_date 이면 same-day (종가 진입) 이벤트.""" @@ -2030,16 +2069,20 @@ class PaperTradingEngine: "score": candidate.score, "reason": plan.skip_reason, }) continue - # Gap cap check for next_open entries (matches BacktestRunner) - from libs.backtest.execution import check_next_open_gap_cap - today_bar = self._broker.get_bar(candidate.symbol) if hasattr(self._broker, 'get_bar') else None - gap_reason = check_next_open_gap_cap(candidate, today_bar) - if gap_reason: - rejected.append({ - "symbol": candidate.symbol, "event_type": candidate.event_type, - "score": candidate.score, "reason": gap_reason, - }) - continue + # Gap cap check for next_open entries (matches BacktestRunner). + # Skipped for lookback entries: multi-day price drift vs reaction-day + # close is not comparable to an overnight gap. + is_lookback = bool(candidate.features.get("is_lookback_entry", False)) + if not is_lookback: + from libs.backtest.execution import check_next_open_gap_cap + today_bar = self._broker.get_bar(candidate.symbol) if hasattr(self._broker, 'get_bar') else None + gap_reason = check_next_open_gap_cap(candidate, today_bar) + if gap_reason: + rejected.append({ + "symbol": candidate.symbol, "event_type": candidate.event_type, + "score": candidate.score, "reason": gap_reason, + }) + continue try: order = order_fn(candidate.symbol, plan.shares) logger.info("paper_engine_buy_submitted", symbol=candidate.symbol, qty=plan.shares, order_id=order.id) @@ -2059,6 +2102,7 @@ class PaperTradingEngine: else: fill_price = plan.entry_price_limit # MOC: actual price unknown until close + initial_days_held = int(candidate.features.get("lookback_days_elapsed", 0)) self._state.save_strategy_state( session_id, StrategyStateRow( @@ -2066,7 +2110,7 @@ class PaperTradingEngine: engine_id=candidate.engine_id, order_id=order.id, entry_date=today.isoformat(), stop_price=plan.stop_price, target_price=plan.target_price, current_stop=plan.stop_price, peak_price=fill_price, - days_held=0, trade_direction=candidate.trade_direction, + days_held=initial_days_held, trade_direction=candidate.trade_direction, candidate_json=candidate.model_dump_json(), plan_json=plan.model_dump_json(), status="open", ), diff --git a/apps/paper_trader/event_detector.py b/apps/paper_trader/event_detector.py index 761f6d1..e1b9046 100644 --- a/apps/paper_trader/event_detector.py +++ b/apps/paper_trader/event_detector.py @@ -35,37 +35,24 @@ class EventDetector: @staticmethod def _compute_score(row: dict[str, Any], config: BacktestConfig) -> float: """Compute score using the config's scoring_model — matches backtester.""" + import libs.backtest.scoring as _scoring model = config.signal.scoring_model - if model == "return_max_long_v5": - from libs.backtest.scoring import compute_return_max_long_score_v5 - return compute_return_max_long_score_v5(row) - elif model == "return_max_long_v7": - from libs.backtest.scoring import compute_return_max_long_score_v7 - return compute_return_max_long_score_v7(row) - elif model == "return_max_long_v8": - from libs.backtest.scoring import compute_return_max_long_score_v8 - return compute_return_max_long_score_v8(row) - elif model == "return_max_long_v9": - from libs.backtest.scoring import compute_return_max_long_score_v9 - return compute_return_max_long_score_v9(row) - elif model == "return_max_long_v9g": - from libs.backtest.scoring import compute_return_max_long_score_v9g - return compute_return_max_long_score_v9g(row) - elif model == "return_max_long_v10": - from libs.backtest.scoring import compute_return_max_long_score_v10 - return compute_return_max_long_score_v10(row) - elif model == "return_max_long_v11": - from libs.backtest.scoring import compute_return_max_long_score_v11 - return compute_return_max_long_score_v11(row) - elif model == "return_max_long_v11g": - from libs.backtest.scoring import compute_return_max_long_score_v11g - return compute_return_max_long_score_v11g(row) - elif model == "pead": - from libs.backtest.scoring import compute_pead_score - return compute_pead_score(row) - else: - from libs.backtest.scoring import compute_entry_score - return compute_entry_score(row) + + # Derive function name from model name to stay in sync with new models + # without requiring manual updates here. + # "return_max_long_v13e" → compute_return_max_long_score_v13e + # "return_max_long_v1" → compute_return_max_long_score (legacy, no suffix) + # "pead" / "patient_drift" / "microstructure" → compute_{model}_score + fn: Any = None + if model.startswith("return_max_long_"): + suffix = model[len("return_max_long_"):] + fn_name = "compute_return_max_long_score" if suffix == "v1" else f"compute_return_max_long_score_{suffix}" + fn = getattr(_scoring, fn_name, None) + if fn is None: + fn = getattr(_scoring, f"compute_{model}_score", None) + if fn is None: + fn = _scoring.compute_entry_score + return fn(row) async def get_candidates_for_date( self, @@ -90,17 +77,78 @@ class EventDetector: if not raw_rows: logger.debug("event_detector_no_events", date=execution_date.isoformat()) return [] + enriched = await self._enrich_raw_rows(raw_rows, bar_end_date=execution_date, config=config) + logger.info( + "event_detector_candidates_ready", + date=execution_date.isoformat(), + count=len(enriched), + ) + return enriched + + async def get_candidates_for_lookback( + self, + today: dt.date, + start_date: dt.date, + config: BacktestConfig, + ) -> list[dict[str, Any]]: + """Return enriched candidate rows for events with entry_date in [start_date, today). + Used on daemon startup to pick up events from previous trading days that are + still within their max_holding_days window. Each returned row has: + - is_lookback_entry=True + - lookback_days_elapsed=N (trading days since the original execution_date) + """ + raw_rows = await self._fetch_events_for_date_range(start_date, today) + if not raw_rows: + logger.debug( + "event_detector_no_lookback_events", + start=start_date.isoformat(), + end=today.isoformat(), + ) + return [] + enriched = await self._enrich_raw_rows(raw_rows, bar_end_date=today, config=config) + # Annotate with lookback metadata + from libs.backtest.calendar import get_trading_days + _elapsed_cache: dict[dt.date, int] = {} + for row in enriched: + raw_exec = row.get("execution_date") or row.get("entry_date") + if raw_exec: + exec_date = raw_exec if isinstance(raw_exec, dt.date) else dt.date.fromisoformat(str(raw_exec)) + if exec_date not in _elapsed_cache: + tdays = get_trading_days(exec_date, today) + _elapsed_cache[exec_date] = max(0, len(tdays) - 1) + row["is_lookback_entry"] = True + row["lookback_days_elapsed"] = _elapsed_cache[exec_date] + logger.info( + "event_detector_lookback_candidates_ready", + start=start_date.isoformat(), + today=today.isoformat(), + count=len(enriched), + ) + return enriched + + async def _enrich_raw_rows( + self, + raw_rows: list[dict[str, Any]], + bar_end_date: dt.date, + config: BacktestConfig, + ) -> list[dict[str, Any]]: + """Enrich raw DB rows with bars, market features, and scores. + + Shared by get_candidates_for_date and get_candidates_for_lookback. + bar_end_date is the upper bound for bar fetches (usually execution_date or today). + """ unique_symbols = sorted({ str(r.get("symbol", "")).upper() for r in raw_rows if r.get("symbol") }) - # Fetch bars for the last 30 days for ADV + ATR computation - bar_start = execution_date - dt.timedelta(days=45) + # Fetch 120 days of bars — enough for 60d features (entropy, Hurst, gravitational pull) + # which need ~65 trading days (~91 calendar days) before the reaction date. + bar_start = bar_end_date - dt.timedelta(days=120) bars_by_symbol, avg_dvol, atr_by_symbol = await self._fetch_enrichment_data( - unique_symbols, bar_start, execution_date + unique_symbols, bar_start, bar_end_date ) company_info = await self._fetch_company_info(unique_symbols) screener_mcaps = await self._fetch_screener_market_caps(unique_symbols) @@ -217,6 +265,57 @@ class EventDetector: ) continue + # Compute tier2/tier3 features from bars if missing from DB. + # The DB FeatureSnapshot only stores basic NLP/event features; technical + # features (entropy, Hurst, gravitational pull, etc.) are computed only in + # the Parquet enrichment pipeline. Reproduce them here from Oracle bars. + # Use event_date (pre-event baseline) to match the Parquet pipeline exactly. + # Fall back to reaction_date if event_date is unavailable or not a trading day. + sym_bars_for_features = bars_by_symbol.get(sym, {}) + _event_date_raw = enriched.get("event_date") + _ed_candidate = _parse_date(_event_date_raw) + # event_date must be present in bars (i.e., a trading day) to use it + if _ed_candidate and _ed_candidate in sym_bars_for_features: + rd_for_features = _ed_candidate + else: + rd_for_features = _parse_date(enriched.get("reaction_date")) + if rd_for_features and sym_bars_for_features: + from libs.features.market_features import ( + avg_dollar_volume_20d as _adv20d_fn, + pre_event_bb_position as _bb_pos_fn, + pre_event_entropy as _entropy_fn, + pre_event_gravitational_pull as _grav_pull_fn, + pre_event_hurst as _hurst_fn, + pre_event_market_temperature as _mkt_temp_fn, + ) + from libs.oracle_client.models import PriceBar as _OraclePriceBar + + _price_bars = [ + _OraclePriceBar( + date=d.isoformat(), + open=float(b.get("open", 0)), + high=float(b.get("high", 0)), + low=float(b.get("low", 0)), + close=float(b.get("close", 0)), + volume=int(b.get("volume", 0)), + ) + for d, b in sorted(sym_bars_for_features.items()) + ] + _rd_str = rd_for_features.isoformat() + + if enriched.get("avg_dollar_volume_20d") is None: + enriched["avg_dollar_volume_20d"] = _adv20d_fn(_price_bars, _rd_str) + if enriched.get("pre_event_bb_position") is None: + enriched["pre_event_bb_position"] = _bb_pos_fn(_price_bars, _rd_str) + if enriched.get("pre_event_hurst_60d") is None: + enriched["pre_event_hurst_60d"] = _hurst_fn(_price_bars, _rd_str) + if enriched.get("pre_event_entropy_60d") is None: + enriched["pre_event_entropy_60d"] = _entropy_fn(_price_bars, _rd_str) + if enriched.get("pre_event_gravitational_pull") is None: + enriched["pre_event_gravitational_pull"] = _grav_pull_fn(_price_bars, _rd_str) + if enriched.get("pre_event_market_temperature") is None: + enriched["pre_event_market_temperature"] = _mkt_temp_fn(_price_bars, _rd_str) + # Compute score using config's scoring model for consistency with # BacktestRunner. Only use config model when event_v1 features are # present (parse_confidence_overall etc.), otherwise the model's hard @@ -229,11 +328,6 @@ class EventDetector: enriched_rows.append(enriched) - logger.info( - "event_detector_candidates_ready", - date=execution_date.isoformat(), - count=len(enriched_rows), - ) return enriched_rows # ------------------------------------------------------------------ # @@ -319,6 +413,7 @@ class EventDetector: "event_timestamp": ts, "reaction_date": label.reaction_date, "entry_date": label.entry_date, + "entry_convention": label.entry_convention, "execution_date": execution_date, "label_status": label.label_status, # Merged features from ALL FeatureSnapshots for this event @@ -340,6 +435,98 @@ class EventDetector: self._db_unavailable = True return [] + async def _fetch_events_for_date_range( + self, + start_date: dt.date, + end_date: dt.date, + ) -> list[dict[str, Any]]: + """Query EventLabel rows where entry_date in [start_date, end_date). + + Used for lookback entry: surfaces events that fired before the daemon + started but are still within their max_holding_days window. + """ + if self._db_unavailable: + return [] + try: + from sqlalchemy import select + from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine + + from libs.db.models import Event, EventLabel, FeatureSnapshot, SymbolMaster + + engine = create_async_engine( + self._db_dsn, echo=False, + connect_args={"timeout": 5}, + ) + async_session = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) + + _UTC = __import__("zoneinfo").ZoneInfo("UTC") + + async with async_session() as session: + stmt = ( + select(Event, EventLabel, FeatureSnapshot, SymbolMaster) + .join(EventLabel, Event.event_id == EventLabel.event_id) + .join(FeatureSnapshot, Event.event_id == FeatureSnapshot.event_id) + .outerjoin(SymbolMaster, Event.symbol_id == SymbolMaster.symbol_id) + .where(EventLabel.entry_date >= start_date) + .where(EventLabel.entry_date < end_date) + .where(EventLabel.label_status.in_(["ok", "truncated", "pending"])) + .order_by(FeatureSnapshot.created_at_utc.asc()) + ) + rows = (await session.execute(stmt)).all() + + await engine.dispose() + + event_meta: dict[str, tuple] = {} + event_features: dict[str, dict] = {} + + for event, label, snapshot, sym in rows: + eid = event.event_id + if sym is None or not sym.ticker: + continue + if eid not in event_meta: + event_meta[eid] = (event, label, sym) + event_features[eid] = {} + event_features[eid].update(snapshot.feature_json or {}) + + result: list[dict[str, Any]] = [] + for eid, (event, label, sym) in event_meta.items(): + ts = event.filed_at_utc + if ts is None and event.event_date is not None: + ts = dt.datetime.combine( + event.event_date, dt.time(21, 0), tzinfo=_UTC + ) + row: dict[str, Any] = { + "event_id": eid, + "symbol": sym.ticker, + "issuer_id": event.issuer_id, + "event_type": event.event_type or "", + "event_direction": event.event_direction or "", + "event_date": event.event_date, + "event_timestamp": ts, + "reaction_date": label.reaction_date, + "entry_date": label.entry_date, + "entry_convention": label.entry_convention, + "execution_date": label.entry_date, # use entry_date as execution_date + "label_status": label.label_status, + **event_features[eid], + } + result.append(row) + + logger.debug( + "event_detector_db_range_fetched", + start=start_date.isoformat(), + end=end_date.isoformat(), + count=len(result), + ) + return result + + except Exception as exc: + if not self._db_unavailable: + err_msg = f"{type(exc).__name__}: {exc}" if str(exc) else type(exc).__name__ + logger.warning("event_detector_db_range_fetch_failed", error=err_msg) + self._db_unavailable = True + return [] + # ------------------------------------------------------------------ # # Enrichment helpers # ------------------------------------------------------------------ #