"""PaperTradingEngine: daily processing loop for paper trading. Strategy decisions (WHAT to buy/sell) are made locally using backtest logic. Order execution (HOW to execute) is done via Alpaca Paper Trading API. """ from __future__ import annotations import datetime as dt from typing import Any from libs.backtest.allocator import build_planned_order from libs.backtest.domain import ( BacktestConfig, Candidate, DailyPortfolioState, ExecutionConfig, OpenPosition, PlannedOrder, PositionStatus, ) from libs.backtest.execution import simulate_exit, update_trailing_stop from libs.backtest.manifests import load_manifest, resolve_config from libs.backtest.selector import select_candidates from libs.common.logging import get_logger from apps.paper_trader.alpaca_broker import AlpacaBroker, AccountInfo, Position from apps.paper_trader.event_detector import EventDetector from apps.paper_trader.state import ( DailySnapshotRow, SessionRow, StateManager, StrategyStateRow, ) logger = get_logger(__name__) # Kill-switch threshold (matches backtest) _KILL_SWITCH_DRAWDOWN_PCT = 25.0 def _is_reaction_close_entry(candidate_json: str) -> bool: """Return True if the position was entered at the reaction-day CLOSE (MOC order). Same-day events (timing_class == "same_day") enter via MOC; after-close events enter at the next open. This distinction matters for exit checking: MOC entries must not have their stop checked against the entry bar's intraday low/high, because the position did not exist during that intraday period. """ try: import json cand = json.loads(candidate_json) return cand.get("timing_class") == "same_day" except Exception: return False class PaperTradingEngine: """Daily processing loop. Mirrors BacktestRunner._simulate_day() for live use.""" def __init__( self, session: SessionRow, broker: AlpacaBroker, state: StateManager, event_detector: EventDetector, ) -> None: self._session = session self._broker = broker self._state = state self._detector = event_detector manifest = load_manifest(session.config_path) self._config: BacktestConfig = resolve_config(manifest) # Shared attention filtering service (matches BacktestRunner) from libs.backtest.attention import AttentionFilterService oracle_url = event_detector._oracle_url if hasattr(event_detector, '_oracle_url') else "" self._attention_service = AttentionFilterService( oracle_url=oracle_url, scoring_model=self._config.signal.scoring_model, ) # ------------------------------------------------------------------ # # Main entry point # ------------------------------------------------------------------ # async def run_daily(self, target_date: dt.date | None = None, force: bool = False) -> dict[str, Any]: """Process one trading day. Returns a summary dict for the CLI to display.""" today = target_date or dt.date.today() session_id = self._session.session_id # 1. Idempotency: skip if already processed (unless forced) if not force and self._state.is_date_processed(session_id, today): logger.info("paper_engine_already_processed", date=today.isoformat()) return {"date": today, "status": "already_processed"} # 2. Check trading day from libs.common.time_utils import is_trading_day if not is_trading_day(today): logger.info("paper_engine_non_trading_day", date=today.isoformat()) return {"date": today, "status": "non_trading_day"} # 3. Fetch Alpaca state account = self._broker.get_account() alpaca_positions = self._broker.list_positions() held_symbols = [p.symbol for p in alpaca_positions] # 4. Load local strategy states strategy_states = { ss.symbol: ss for ss in self._state.get_open_strategy_states(session_id) } # 5. Fetch price bars for held positions (last 30 days) bars_by_symbol: dict[str, dict[dt.date, dict]] = {} if held_symbols: bar_start = today - dt.timedelta(days=30) bars_by_symbol = self._broker.get_bars_as_dict(held_symbols, bar_start, today) # ============================================================ # EXIT PHASE # ============================================================ exits: list[dict[str, Any]] = [] session_st = self._state.get_session_state(session_id) net_pnl_today = 0.0 for alpaca_pos in alpaca_positions: sym = alpaca_pos.symbol ss = strategy_states.get(sym) if ss is None: logger.debug("paper_engine_no_local_state", symbol=sym) continue ss.days_held += 1 sym_bars = bars_by_symbol.get(sym, {}) bar = sym_bars.get(today) if bar is None: logger.warning("paper_engine_no_bar", symbol=sym, date=today.isoformat()) self._state.update_strategy_state( session_id, sym, days_held=ss.days_held ) continue # Convert to OpenPosition for backtest logic open_pos = self._to_open_position(alpaca_pos, ss) # Update trailing stop effective_exec = self._resolve_execution_config(ss) if effective_exec.trailing_model: update_trailing_stop( open_pos, bar, trailing_model=effective_exec.trailing_model, warmup_days=effective_exec.trailing_warmup_days, ) ss.current_stop = open_pos.current_stop ss.peak_price = open_pos.peak_price # Check exit filled_trade = simulate_exit(open_pos, bar, effective_exec, today) if filled_trade is not None: try: close_qty = None if filled_trade.shares < alpaca_pos.qty: close_qty = filled_trade.shares self._broker.close_position(sym, qty=close_qty, fill_price=filled_trade.exit_price) logger.info( "paper_engine_exit", symbol=sym, reason=filled_trade.exit_reason.value, pnl=filled_trade.net_pnl, ) except Exception as exc: logger.error("paper_engine_close_failed", symbol=sym, error=str(exc)) continue self._state.close_strategy_state(session_id, sym) self._state.record_trade( session_id=session_id, symbol=sym, entry_date=ss.entry_date, exit_date=today.isoformat(), entry_price=alpaca_pos.avg_entry_price, exit_price=filled_trade.exit_price, exit_reason=filled_trade.exit_reason.value, shares=filled_trade.shares, net_pnl=filled_trade.net_pnl, r_multiple=filled_trade.r_multiple, holding_days=ss.days_held, ) net_pnl_today += filled_trade.net_pnl # Update consecutive losses / cooldown if filled_trade.net_pnl < 0: session_st.consecutive_losses += 1 streak = self._config.risk.cooldown_after_loss_streak if streak > 0 and session_st.consecutive_losses >= streak: session_st.cooldown_remaining = self._config.risk.cooldown_days session_st.consecutive_losses = 0 else: session_st.consecutive_losses = 0 exits.append({ "symbol": sym, "reason": filled_trade.exit_reason.value, "pnl": filled_trade.net_pnl, "r_multiple": filled_trade.r_multiple, "shares": filled_trade.shares, "exit_price": filled_trade.exit_price, }) else: # No exit — persist updated trailing state self._state.update_strategy_state( session_id, sym, days_held=ss.days_held, current_stop=ss.current_stop, peak_price=ss.peak_price, ) # Decrement cooldown if session_st.cooldown_remaining > 0: session_st.cooldown_remaining -= 1 # Reset daily risk usage session_st.daily_new_risk_used = 0.0 # ============================================================ # ENTRY PHASE # ============================================================ entries: list[dict[str, Any]] = [] rejected: list[dict[str, Any]] = [] # Refresh account/positions after exits account = self._broker.get_account() alpaca_positions_after_exits = self._broker.list_positions() strategy_states_after_exits = { ss.symbol: ss for ss in self._state.get_open_strategy_states(session_id) } # Get event candidates for today — reaction_close convention only candidate_rows = await self._detector.get_candidates_for_date( today, self._config, convention="reaction_close" ) # Filter already-processed events n_before = len(candidate_rows) candidate_rows = [ r for r in candidate_rows if not self._state.has_processed_event(session_id, str(r.get("event_id", ""))) ] if n_before != len(candidate_rows): logger.debug( "paper_engine_candidates_after_dedup", before=n_before, after=len(candidate_rows), ) # Run selection pipeline for each enabled engine open_positions = self._to_open_positions(alpaca_positions_after_exits, strategy_states_after_exits) portfolio_state = self._build_portfolio_state(account, alpaca_positions_after_exits, today) # Fetch macro data macro_data = await self._fetch_macro(today) engine_daily_risk_used: dict[str, float] = {} engines = self._config.get_active_strategy_engines() logger.debug( "paper_engine_selection_input", date=today.isoformat(), candidate_rows=len(candidate_rows), engines=len(engines), symbols=[r.get("symbol") for r in candidate_rows], ) if candidate_rows: sample = candidate_rows[0] logger.debug( "paper_engine_sample_row", symbol=sample.get("symbol"), event_type=sample.get("event_type"), event_direction=sample.get("event_direction"), filing_time_bucket=sample.get("filing_time_bucket"), entry_price_est=sample.get("entry_price_est"), avg_dollar_volume=sample.get("avg_dollar_volume"), avg_dollar_volume_20d=sample.get("avg_dollar_volume_20d"), event_close=sample.get("event_close"), close_location=sample.get("close_location"), gap_size=sample.get("gap_size"), reaction_day_return=sample.get("reaction_day_return"), market_cap_proxy=sample.get("market_cap_proxy"), execution_date=str(sample.get("execution_date")), event_timestamp=str(sample.get("event_timestamp")), ) if engines: for engine_cfg in engines: engine_candidates = select_candidates( raw_rows=candidate_rows, universe_config=self._config.universe, signal_config=self._config.signal, event_type_profiles=self._config.event_type_profiles or {}, strategy_engine=engine_cfg, excluded_event_ids={ ss.event_id for ss in strategy_states_after_exits.values() }, excluded_symbols={p.symbol for p in alpaca_positions_after_exits if p.symbol in strategy_states_after_exits}, ) # Attention filtering (matches BacktestRunner) engine_candidates = self._attention_service.apply_filters( engine_candidates, engine_cfg, self._config.signal, ) engine_risk_used = engine_daily_risk_used.get(engine_cfg.engine_id, 0.0) for candidate in engine_candidates: plan = build_planned_order( candidate=candidate, portfolio_state=portfolio_state, open_positions=open_positions, config=self._config, cooldown_remaining=session_st.cooldown_remaining, macro_data=macro_data, engine_daily_new_risk_used=engine_risk_used, ) self._state.record_processed_event( session_id, candidate.event_id, today.isoformat(), "rejected" if plan.skip_reason else "entered", skip_reason=plan.skip_reason, ) if plan.skip_reason: rejected.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": plan.skip_reason, }) continue # Submit market buy via Alpaca try: order = self._broker.submit_market_buy(candidate.symbol, plan.shares) logger.info( "paper_engine_buy_submitted", symbol=candidate.symbol, qty=plan.shares, order_id=order.id, ) except Exception as exc: logger.error( "paper_engine_buy_failed", symbol=candidate.symbol, error=str(exc), ) rejected.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": f"order_failed:{exc}", }) continue # Save local strategy state self._state.save_strategy_state( session_id, StrategyStateRow( session_id=session_id, symbol=candidate.symbol, event_id=candidate.event_id, 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=plan.entry_price_limit, days_held=0, trade_direction=candidate.trade_direction, candidate_json=candidate.model_dump_json(), plan_json=plan.model_dump_json(), status="open", ), ) trade_risk = portfolio_state.equity * ( candidate.engine_per_trade_risk_pct or self._config.risk.per_trade_risk_pct ) engine_daily_risk_used[engine_cfg.engine_id] = engine_risk_used + trade_risk session_st.daily_new_risk_used += trade_risk # Refresh portfolio state after each entry open_positions = self._to_open_positions( alpaca_positions_after_exits, strategy_states_after_exits ) # Add the new virtual position to open_positions for gate checks new_open = self._virtual_open_position(candidate, plan, today) open_positions.append(new_open) portfolio_state = DailyPortfolioState( date=portfolio_state.date, equity=portfolio_state.equity, cash_available=max( 0.0, portfolio_state.cash_available - plan.entry_price_limit * plan.shares, ), gross_exposure=portfolio_state.gross_exposure + plan.entry_price_limit * plan.shares, net_exposure=portfolio_state.net_exposure + plan.entry_price_limit * plan.shares, reserved_risk_budget=portfolio_state.reserved_risk_budget, unrealized_pnl=portfolio_state.unrealized_pnl, realized_pnl=portfolio_state.realized_pnl, open_positions=[p.position_id for p in open_positions], daily_new_risk_used=session_st.daily_new_risk_used, peak_equity=portfolio_state.peak_equity, current_drawdown_pct=portfolio_state.current_drawdown_pct, ) entries.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "shares": plan.shares, "entry_price": plan.entry_price_limit, "stop": plan.stop_price, "target": plan.target_price, "order_id": order.id, }) else: # No engines defined — use flat candidate selection all_candidates = select_candidates( raw_rows=candidate_rows, universe_config=self._config.universe, signal_config=self._config.signal, event_type_profiles=self._config.event_type_profiles or {}, excluded_event_ids={ss.event_id for ss in strategy_states_after_exits.values()}, excluded_symbols={p.symbol for p in alpaca_positions_after_exits if p.symbol in strategy_states_after_exits}, ) for candidate in all_candidates: plan = build_planned_order( candidate=candidate, portfolio_state=portfolio_state, open_positions=open_positions, config=self._config, cooldown_remaining=session_st.cooldown_remaining, macro_data=macro_data, ) self._state.record_processed_event( session_id, candidate.event_id, today.isoformat(), "rejected" if plan.skip_reason else "entered", skip_reason=plan.skip_reason, ) if plan.skip_reason: rejected.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": plan.skip_reason, }) continue try: order = self._broker.submit_market_buy(candidate.symbol, plan.shares) except Exception as exc: rejected.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": f"order_failed:{exc}", }) continue self._state.save_strategy_state( session_id, StrategyStateRow( session_id=session_id, symbol=candidate.symbol, event_id=candidate.event_id, 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=plan.entry_price_limit, days_held=0, trade_direction=candidate.trade_direction, candidate_json=candidate.model_dump_json(), plan_json=plan.model_dump_json(), status="open", ), ) entries.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "shares": plan.shares, "entry_price": plan.entry_price_limit, "stop": plan.stop_price, "target": plan.target_price, "order_id": order.id, }) # ============================================================ # Daily snapshot # ============================================================ account_final = self._broker.get_account() peak_equity = self._state.get_peak_equity(session_id, self._session.initial_equity) peak_equity = max(peak_equity, account_final.equity) prev_snapshots = self._state.list_snapshots(session_id) prev_equity = prev_snapshots[-1]["equity"] if prev_snapshots else account_final.equity baseline_equity = prev_snapshots[0]["equity"] if prev_snapshots else account_final.equity total_pnl = account_final.equity - baseline_equity drawdown_pct = max(0.0, (peak_equity - account_final.equity) / peak_equity * 100) if peak_equity > 0 else 0.0 self._state.save_daily_snapshot( DailySnapshotRow( session_id=session_id, date=today.isoformat(), equity=account_final.equity, cash=account_final.cash, market_value=account_final.long_market_value, daily_pnl=account_final.equity - prev_equity, total_pnl=total_pnl, drawdown_pct=drawdown_pct, open_position_count=len(self._broker.list_positions()), ) ) # Persist session state (cooldown, etc.) session_st.last_processed_date = today.isoformat() self._state.update_session_state(session_st) self._state.mark_date_processed(session_id, today) summary = { "date": today, "status": "processed", "exits": exits, "entries": entries, "rejected": rejected, "candidates_detected": len(candidate_rows), "account": { "equity": account_final.equity, "cash": account_final.cash, "market_value": account_final.long_market_value, "total_pnl": total_pnl, "drawdown_pct": drawdown_pct, }, } logger.info( "paper_engine_day_done", date=today.isoformat(), exits=len(exits), entries=len(entries), rejected=len(rejected), ) return summary # ------------------------------------------------------------------ # # Phased execution: reaction_close / next_open / monitor # ------------------------------------------------------------------ # async def run_reaction_close( self, target_date: dt.date | None = None, force: bool = False ) -> dict[str, Any]: """장 마감 직전 (~3:40 PM ET): same-day 이벤트 후보 → MOC 매수 주문. 파이프라인 없이도 호출 가능. 당일 DB에 이미 적재된 이벤트를 사용. """ today = target_date or dt.date.today() session_id = self._session.session_id phase = "reaction_close" if not force and self._state.is_phase_processed(session_id, today, phase): logger.info("paper_engine_already_processed", date=today.isoformat(), phase=phase) return {"date": today, "status": "already_processed", "phase": phase} from libs.common.time_utils import is_trading_day if not is_trading_day(today): return {"date": today, "status": "non_trading_day", "phase": phase} account = self._broker.get_account() alpaca_positions = self._broker.list_positions() session_st = self._state.get_session_state(session_id) strategy_states = { ss.symbol: ss for ss in self._state.get_open_strategy_states(session_id) } all_rows = await self._detector.get_candidates_for_date( today, self._config, convention="reaction_close" ) same_day_rows = [r for r in all_rows if self._is_same_day_event(r)] macro_data = await self._fetch_macro(today) entries, rejected = await self._process_entries( today, same_day_rows, account, alpaca_positions, strategy_states, session_st, macro_data, self._broker.submit_moc_buy, ) session_st.last_processed_date = today.isoformat() self._state.update_session_state(session_st) self._state.mark_phase_processed(session_id, today, phase) logger.info( "paper_engine_reaction_close_done", date=today.isoformat(), same_day_candidates=len(same_day_rows), entries=len(entries), rejected=len(rejected), ) return { "date": today, "status": "processed", "phase": phase, "entries": entries, "rejected": rejected, "candidates_detected": len(same_day_rows), "account": {"equity": account.equity, "cash": account.cash, "market_value": account.long_market_value}, } async def run_next_open( self, target_date: dt.date | None = None, force: bool = False ) -> dict[str, Any]: """장 시작 직후 (~9:30 AM ET): 전날 바로 exit 판단 + after-close 이벤트 → 시장가 매수. 파이프라인이 전날 저녁 실행됐다고 가정. """ today = target_date or dt.date.today() session_id = self._session.session_id phase = "next_open" if not force and self._state.is_phase_processed(session_id, today, phase): logger.info("paper_engine_already_processed", date=today.isoformat(), phase=phase) return {"date": today, "status": "already_processed", "phase": phase} from libs.common.time_utils import is_trading_day if not is_trading_day(today): return {"date": today, "status": "non_trading_day", "phase": phase} account = self._broker.get_account() alpaca_positions = self._broker.list_positions() session_st = self._state.get_session_state(session_id) strategy_states = { ss.symbol: ss for ss in self._state.get_open_strategy_states(session_id) } # Exit: 전날 (bar_date = today - 1) 종가 기준으로 exit 판단 prev_date = today - dt.timedelta(days=1) exits = await self._process_exits( today, prev_date, alpaca_positions, strategy_states, session_st ) # Refresh Alpaca state after exits account = self._broker.get_account() alpaca_positions = self._broker.list_positions() strategy_states = { ss.symbol: ss for ss in self._state.get_open_strategy_states(session_id) } # Entry: after-close 이벤트만 — next_open_after_reaction_close convention all_rows = await self._detector.get_candidates_for_date( today, self._config, convention="next_open_after_reaction_close" ) after_close_rows = [r for r in all_rows if not self._is_same_day_event(r)] macro_data = await self._fetch_macro(today) entries, rejected = await self._process_entries( today, after_close_rows, account, alpaca_positions, strategy_states, session_st, macro_data, self._broker.submit_market_buy, ) summary = self._finalize_day(today, session_st, exits, entries, rejected, len(all_rows)) summary["phase"] = phase self._state.mark_phase_processed(session_id, today, phase) return summary async def run_monitor(self, interval_sec: int = 60) -> None: """장중 실시간 모니터링: stop/target 조건 충족 시 즉시 청산. Ctrl+C 로 종료. 별도 터미널에서 실행 권장. """ import asyncio as _asyncio session_id = self._session.session_id logger.info("paper_engine_monitor_start", session=session_id, interval_sec=interval_sec) while True: try: alpaca_positions = self._broker.list_positions() strategy_states = { ss.symbol: ss for ss in self._state.get_open_strategy_states(session_id) } for pos in alpaca_positions: ss = strategy_states.get(pos.symbol) if ss is None: continue price = pos.current_price # Trailing peak 업데이트 if price > ss.peak_price: ss.peak_price = price self._state.update_strategy_state( session_id, pos.symbol, peak_price=price ) # Stop 조건 if price < ss.current_stop: logger.info( "paper_engine_monitor_stop_hit", symbol=pos.symbol, price=price, stop=ss.current_stop, ) self._monitor_close(pos, ss, "STOP_INTRADAY", price) # Target 조건 elif ss.target_price and price >= ss.target_price: logger.info( "paper_engine_monitor_target_hit", symbol=pos.symbol, price=price, target=ss.target_price, ) self._monitor_close(pos, ss, "TARGET_INTRADAY", price) except Exception as exc: logger.error("paper_engine_monitor_error", error=str(exc)) await _asyncio.sleep(interval_sec) def _monitor_close(self, pos: Any, ss: Any, reason: str, price: float) -> None: """모니터링 루프에서 포지션 청산 처리.""" session_id = self._session.session_id try: self._broker.close_position(pos.symbol, fill_price=price) self._state.close_strategy_state(session_id, pos.symbol) self._state.record_trade( session_id=session_id, symbol=pos.symbol, entry_date=ss.entry_date, exit_date=dt.date.today().isoformat(), entry_price=pos.avg_entry_price, exit_price=price, exit_reason=reason, shares=pos.qty, net_pnl=(price - pos.avg_entry_price) * pos.qty, r_multiple=0.0, holding_days=ss.days_held, ) except Exception as exc: logger.error("paper_engine_monitor_close_failed", symbol=pos.symbol, error=str(exc)) # ------------------------------------------------------------------ # # Shared helpers for phased execution # ------------------------------------------------------------------ # @staticmethod def _is_same_day_event(row: dict[str, Any]) -> bool: """event_date == reaction_date 이면 same-day (종가 진입) 이벤트.""" def _pd(v: Any) -> dt.date | None: if isinstance(v, dt.datetime): return v.date() if isinstance(v, dt.date): return v if isinstance(v, str): try: return dt.date.fromisoformat(v[:10]) except ValueError: return None return None ed = _pd(row.get("event_date")) rd = _pd(row.get("reaction_date")) return ed is not None and rd is not None and ed == rd async def _process_exits( self, today: dt.date, bar_date: dt.date, alpaca_positions: list[Any], strategy_states: dict[str, Any], session_st: Any, ) -> list[dict[str, Any]]: """보유 포지션에 대해 bar_date 기준 exit 로직 실행.""" session_id = self._session.session_id exits: list[dict[str, Any]] = [] held_symbols = [p.symbol for p in alpaca_positions] if not held_symbols: session_st.daily_new_risk_used = 0.0 return exits bar_start = bar_date - dt.timedelta(days=30) bars_by_symbol = self._broker.get_bars_as_dict(held_symbols, bar_start, bar_date) for alpaca_pos in alpaca_positions: sym = alpaca_pos.symbol ss = strategy_states.get(sym) if ss is None: logger.debug("paper_engine_no_local_state", symbol=sym) continue ss.days_held += 1 # Skip exit check for reaction_close (same-day) positions on their entry bar. # These positions were opened at the CLOSE of bar_date, so the intraday # bar data (low/high) precedes the actual entry and must not trigger stops. # next_open positions are NOT skipped — they entered at the OPEN so the # full day bar is valid for exit checking. if ss.entry_date == bar_date.isoformat() and _is_reaction_close_entry(ss.candidate_json): self._state.update_strategy_state( session_id, sym, days_held=ss.days_held, current_stop=ss.current_stop, peak_price=ss.peak_price, ) continue sym_bars = bars_by_symbol.get(sym, {}) available = [d for d in sym_bars if d <= bar_date] bar = sym_bars[max(available)] if available else None if bar is None: logger.warning("paper_engine_no_bar", symbol=sym, date=bar_date.isoformat()) self._state.update_strategy_state(session_id, sym, days_held=ss.days_held) continue open_pos = self._to_open_position(alpaca_pos, ss) effective_exec = self._resolve_execution_config(ss) if effective_exec.trailing_model: update_trailing_stop( open_pos, bar, trailing_model=effective_exec.trailing_model, warmup_days=effective_exec.trailing_warmup_days, ) ss.current_stop = open_pos.current_stop ss.peak_price = open_pos.peak_price filled_trade = simulate_exit(open_pos, bar, effective_exec, bar_date) if filled_trade is not None: is_partial = filled_trade.shares < alpaca_pos.qty try: self._broker.close_position(sym, qty=filled_trade.shares if is_partial else None, fill_price=filled_trade.exit_price) logger.info( "paper_engine_exit", symbol=sym, reason=filled_trade.exit_reason.value, pnl=filled_trade.net_pnl, partial=is_partial, ) except Exception as exc: logger.error("paper_engine_close_failed", symbol=sym, error=str(exc)) continue if is_partial: # T1 partial exit: keep position tracked with breakeven stop self._state.update_strategy_state( session_id, sym, days_held=ss.days_held, current_stop=open_pos.current_stop, # set to entry_price by simulate_exit peak_price=ss.peak_price, status="partial", ) else: self._state.close_strategy_state(session_id, sym) self._state.record_trade( session_id=session_id, symbol=sym, entry_date=ss.entry_date, exit_date=today.isoformat(), entry_price=alpaca_pos.avg_entry_price, exit_price=filled_trade.exit_price, exit_reason=filled_trade.exit_reason.value, shares=filled_trade.shares, net_pnl=filled_trade.net_pnl, r_multiple=filled_trade.r_multiple, holding_days=ss.days_held, ) if filled_trade.net_pnl < 0: session_st.consecutive_losses += 1 streak = self._config.risk.cooldown_after_loss_streak if streak > 0 and session_st.consecutive_losses >= streak: session_st.cooldown_remaining = self._config.risk.cooldown_days session_st.consecutive_losses = 0 else: session_st.consecutive_losses = 0 exits.append({ "symbol": sym, "reason": filled_trade.exit_reason.value, "pnl": filled_trade.net_pnl, "r_multiple": filled_trade.r_multiple, "shares": filled_trade.shares, "exit_price": filled_trade.exit_price, }) else: self._state.update_strategy_state( session_id, sym, days_held=ss.days_held, current_stop=ss.current_stop, peak_price=ss.peak_price, ) if session_st.cooldown_remaining > 0: session_st.cooldown_remaining -= 1 session_st.daily_new_risk_used = 0.0 return exits async def _process_entries( self, today: dt.date, candidate_rows: list[dict[str, Any]], account: Any, alpaca_positions: list[Any], strategy_states: dict[str, Any], session_st: Any, macro_data: dict[str, Any], order_fn: Any, ) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]: """후보군에 대해 진입 판단 + 주문 제출. order_fn = submit_market_buy | submit_moc_buy.""" session_id = self._session.session_id entries: list[dict[str, Any]] = [] rejected: list[dict[str, Any]] = [] candidate_rows = [ r for r in candidate_rows if not self._state.has_processed_event(session_id, str(r.get("event_id", ""))) ] open_positions = self._to_open_positions(alpaca_positions, strategy_states) portfolio_state = self._build_portfolio_state(account, alpaca_positions, today) engines = self._config.get_active_strategy_engines() engine_daily_risk_used: dict[str, float] = {} logger.debug( "paper_engine_selection_input", date=today.isoformat(), candidate_rows=len(candidate_rows), engines=len(engines), symbols=[r.get("symbol") for r in candidate_rows], ) if candidate_rows: s = candidate_rows[0] logger.debug( "paper_engine_sample_row", symbol=s.get("symbol"), event_type=s.get("event_type"), event_direction=s.get("event_direction"), filing_time_bucket=s.get("filing_time_bucket"), entry_price_est=s.get("entry_price_est"), avg_dollar_volume=s.get("avg_dollar_volume"), avg_dollar_volume_20d=s.get("avg_dollar_volume_20d"), event_close=s.get("event_close"), close_location=s.get("close_location"), gap_size=s.get("gap_size"), reaction_day_return=s.get("reaction_day_return"), market_cap_proxy=s.get("market_cap_proxy"), execution_date=str(s.get("execution_date")), event_timestamp=str(s.get("event_timestamp")), ) engine_list = engines if engines else [None] for engine_cfg in engine_list: if engine_cfg is not None: engine_candidates = select_candidates( raw_rows=candidate_rows, universe_config=self._config.universe, signal_config=self._config.signal, event_type_profiles=self._config.event_type_profiles or {}, strategy_engine=engine_cfg, excluded_event_ids={ss.event_id for ss in strategy_states.values()}, excluded_symbols={p.symbol for p in alpaca_positions if p.symbol in strategy_states}, ) # Attention filtering (matches BacktestRunner) engine_candidates = self._attention_service.apply_filters( engine_candidates, engine_cfg, self._config.signal, ) engine_risk_used = engine_daily_risk_used.get(engine_cfg.engine_id, 0.0) else: engine_candidates = select_candidates( raw_rows=candidate_rows, universe_config=self._config.universe, signal_config=self._config.signal, event_type_profiles=self._config.event_type_profiles or {}, excluded_event_ids={ss.event_id for ss in strategy_states.values()}, excluded_symbols={p.symbol for p in alpaca_positions if p.symbol in strategy_states}, ) engine_risk_used = 0.0 for candidate in engine_candidates: plan = build_planned_order( candidate=candidate, portfolio_state=portfolio_state, open_positions=open_positions, config=self._config, cooldown_remaining=session_st.cooldown_remaining, macro_data=macro_data, engine_daily_new_risk_used=engine_risk_used if engine_cfg else 0.0, ) self._state.record_processed_event( session_id, candidate.event_id, today.isoformat(), "rejected" if plan.skip_reason else "entered", skip_reason=plan.skip_reason, ) if plan.skip_reason: rejected.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "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 try: order = order_fn(candidate.symbol, plan.shares) logger.info("paper_engine_buy_submitted", symbol=candidate.symbol, qty=plan.shares, order_id=order.id) except Exception as exc: logger.error("paper_engine_buy_failed", symbol=candidate.symbol, error=str(exc)) rejected.append({"symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": f"order_failed:{exc}"}) continue self._state.save_strategy_state( session_id, StrategyStateRow( session_id=session_id, symbol=candidate.symbol, event_id=candidate.event_id, 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=plan.entry_price_limit, days_held=0, trade_direction=candidate.trade_direction, candidate_json=candidate.model_dump_json(), plan_json=plan.model_dump_json(), status="open", ), ) trade_risk = portfolio_state.equity * ( candidate.engine_per_trade_risk_pct or self._config.risk.per_trade_risk_pct ) if engine_cfg: engine_daily_risk_used[engine_cfg.engine_id] = engine_risk_used + trade_risk session_st.daily_new_risk_used += trade_risk open_positions = self._to_open_positions(alpaca_positions, strategy_states) open_positions.append(self._virtual_open_position(candidate, plan, today)) portfolio_state = DailyPortfolioState( date=portfolio_state.date, equity=portfolio_state.equity, cash_available=max(0.0, portfolio_state.cash_available - plan.entry_price_limit * plan.shares), gross_exposure=portfolio_state.gross_exposure + plan.entry_price_limit * plan.shares, net_exposure=portfolio_state.net_exposure + plan.entry_price_limit * plan.shares, reserved_risk_budget=portfolio_state.reserved_risk_budget, unrealized_pnl=portfolio_state.unrealized_pnl, realized_pnl=portfolio_state.realized_pnl, open_positions=[p.position_id for p in open_positions], daily_new_risk_used=session_st.daily_new_risk_used, peak_equity=portfolio_state.peak_equity, current_drawdown_pct=portfolio_state.current_drawdown_pct, ) entries.append({ "symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "shares": plan.shares, "entry_price": plan.entry_price_limit, "stop": plan.stop_price, "target": plan.target_price, "order_id": order.id, }) return entries, rejected def _finalize_day( self, today: dt.date, session_st: Any, exits: list[dict[str, Any]], entries: list[dict[str, Any]], rejected: list[dict[str, Any]], candidates_detected: int, ) -> dict[str, Any]: """일일 스냅샷 저장 + summary dict 반환.""" session_id = self._session.session_id # 세션 소유 포지션만 집계 (Alpaca 전체 계좌가 아닌 세션 기준) alpaca_positions_final = self._broker.list_positions() session_symbols_final = { ss.symbol for ss in self._state.get_open_strategy_states(session_id) } session_positions_final = [p for p in alpaca_positions_final if p.symbol in session_symbols_final] session_market_value_final = sum(p.market_value for p in session_positions_final) session_unrealized_pl_final = sum(p.unrealized_pl for p in session_positions_final) # 세션 equity = initial_equity + 전체 실현 P&L + 현재 미실현 P&L total_realized_pnl = sum( t.get("net_pnl", 0.0) for t in self._state.list_trades(session_id) ) session_equity_final = self._session.initial_equity + total_realized_pnl + session_unrealized_pl_final session_cash_final = max(0.0, session_equity_final - session_market_value_final) prev_snapshots = self._state.list_snapshots(session_id) prev_equity = prev_snapshots[-1]["equity"] if prev_snapshots else self._session.initial_equity peak_equity = self._state.get_peak_equity(session_id, self._session.initial_equity) peak_equity = max(peak_equity, session_equity_final) total_pnl = session_equity_final - self._session.initial_equity drawdown_pct = ( max(0.0, (peak_equity - session_equity_final) / peak_equity * 100) if peak_equity > 0 else 0.0 ) self._state.save_daily_snapshot( DailySnapshotRow( session_id=session_id, date=today.isoformat(), equity=session_equity_final, cash=session_cash_final, market_value=session_market_value_final, daily_pnl=session_equity_final - prev_equity, total_pnl=total_pnl, drawdown_pct=drawdown_pct, open_position_count=len(session_positions_final), ) ) session_st.last_processed_date = today.isoformat() self._state.update_session_state(session_st) logger.info( "paper_engine_day_done", date=today.isoformat(), exits=len(exits), entries=len(entries), rejected=len(rejected), ) return { "date": today, "status": "processed", "exits": exits, "entries": entries, "rejected": rejected, "candidates_detected": candidates_detected, "account": { "equity": session_equity_final, "cash": session_cash_final, "market_value": session_market_value_final, "total_pnl": total_pnl, "drawdown_pct": drawdown_pct, }, } # ------------------------------------------------------------------ # # Conversion helpers # ------------------------------------------------------------------ # def _to_open_position( self, alpaca_pos: Position, ss: StrategyStateRow ) -> OpenPosition: """Convert Alpaca position + local state to backtest OpenPosition.""" candidate = Candidate.model_validate_json(ss.candidate_json) plan = PlannedOrder.model_validate_json(ss.plan_json) pos_status = ( PositionStatus.PARTIALLY_EXITED if ss.status == "partial" else PositionStatus.ENTERED ) return OpenPosition( position_id=ss.order_id or ss.symbol, plan=plan, entry_date=dt.date.fromisoformat(ss.entry_date), entry_price=alpaca_pos.avg_entry_price, entry_fill_slippage_bps=0.0, current_stop=ss.current_stop, target_price=ss.target_price, peak_price=ss.peak_price, shares_open=alpaca_pos.qty, shares_total=alpaca_pos.qty, days_held=ss.days_held, status=pos_status, ) def _to_open_positions( self, alpaca_positions: list[Position], strategy_states: dict[str, StrategyStateRow], ) -> list[OpenPosition]: result: list[OpenPosition] = [] for alpaca_pos in alpaca_positions: ss = strategy_states.get(alpaca_pos.symbol) if ss is None: continue result.append(self._to_open_position(alpaca_pos, ss)) return result def _virtual_open_position( self, candidate: Candidate, plan: PlannedOrder, entry_date: dt.date ) -> OpenPosition: """Create a virtual OpenPosition for gate-checking after a new entry.""" return OpenPosition( position_id=f"virtual_{candidate.symbol}", plan=plan, entry_date=entry_date, entry_price=plan.entry_price_limit, entry_fill_slippage_bps=0.0, current_stop=plan.stop_price, target_price=plan.target_price, peak_price=plan.entry_price_limit, shares_open=plan.shares, shares_total=plan.shares, days_held=0, status=PositionStatus.ENTERED, ) def _build_portfolio_state( self, account: AccountInfo, alpaca_positions: list[Position], date: dt.date, ) -> DailyPortfolioState: """세션별 독립 equity/cash 기준으로 포트폴리오 상태 계산. Alpaca 계좌는 여러 세션이 공유하므로 account.equity/cash를 직접 쓰면 안 됨. 대신 이 세션 고유의 equity(SQLite 스냅샷 기준)와 이 세션이 보유한 포지션만 사용한다. """ session_id = self._session.session_id # 이 세션 소유 포지션만 (SQLite strategy_states 기준) session_symbols = { ss.symbol for ss in self._state.get_open_strategy_states(session_id) } session_positions = [p for p in alpaca_positions if p.symbol in session_symbols] session_market_value = sum(p.market_value for p in session_positions) session_unrealized_pl = sum(p.unrealized_pl for p in session_positions) # MockBroker (backtest): broker IS the session, use actual cash directly. # AlpacaBroker (live): multiple sessions may share account, derive from snapshot. from apps.paper_trader.mock_broker import MockBroker if isinstance(self._broker, MockBroker): session_cash = max(0.0, account.cash) session_equity = session_cash + session_market_value else: snapshots = self._state.list_snapshots(session_id) session_equity = ( snapshots[-1]["equity"] if snapshots else self._session.initial_equity ) session_cash = max(0.0, session_equity - session_market_value) peak_equity = self._state.get_peak_equity(session_id, self._session.initial_equity) peak_equity = max(peak_equity, session_equity) drawdown_pct = ( max(0.0, (peak_equity - session_equity) / peak_equity * 100) if peak_equity > 0 else 0.0 ) session_st = self._state.get_session_state(session_id) return DailyPortfolioState( date=date, equity=session_equity, cash_available=session_cash, gross_exposure=session_market_value, net_exposure=session_market_value, reserved_risk_budget=0.0, unrealized_pnl=session_unrealized_pl, realized_pnl=0.0, open_positions=[p.symbol for p in session_positions], daily_new_risk_used=session_st.daily_new_risk_used, peak_equity=peak_equity, current_drawdown_pct=drawdown_pct, ) def _resolve_execution_config(self, ss: StrategyStateRow) -> ExecutionConfig: """Get effective ExecutionConfig using shared function. Delegates to libs.backtest.execution.build_effective_execution_config() for consistency with BacktestRunner. """ from libs.backtest.execution import build_effective_execution_config plan = PlannedOrder.model_validate_json(ss.plan_json) return build_effective_execution_config(plan.candidate, self._config) # ------------------------------------------------------------------ # # Macro data # ------------------------------------------------------------------ # async def _fetch_macro(self, date: dt.date) -> dict[str, Any]: """Fetch SPY/QQQ macro data for regime filtering.""" try: sma_period = self._config.risk.macro_sma_period start = date - dt.timedelta(days=sma_period * 2 + 10) symbols = ["SPY", "QQQ"] # Fast path: use bars_cache from EventDetector (backtest mode) bars_cache = getattr(self._detector, "_bars_cache", None) if bars_cache is not None: macro: dict[str, Any] = {} for sym in symbols: all_bars = bars_cache.get(sym, {}) closes = [ float(all_bars[d]["close"]) for d in sorted(all_bars.keys()) if start <= d <= date ] if closes: key_prefix = sym.lower() macro[f"{key_prefix}_close"] = closes[-1] if len(closes) >= sma_period: macro[f"{key_prefix}_sma_{sma_period}"] = sum(closes[-sma_period:]) / sma_period return macro from libs.oracle_client import OracleClient, PriceService async with OracleClient(base_url=self._detector._oracle_url) as client: svc = PriceService(client) async def _fetch_macro_sym(sym: str) -> tuple[str, list[Any]]: try: resp = await svc.get_daily_bars(sym, start=start.isoformat(), end=date.isoformat()) return sym, resp.bars except Exception: return sym, [] results = await __import__("asyncio").gather( *(_fetch_macro_sym(sym) for sym in symbols) ) macro = {} for sym, bars in results: if not bars: continue key_prefix = sym.lower() closes = [float(b.close) for b in bars] if closes: macro[f"{key_prefix}_close"] = closes[-1] if len(closes) >= sma_period: macro[f"{key_prefix}_sma_{sma_period}"] = sum(closes[-sma_period:]) / sma_period return macro except Exception as exc: logger.warning("paper_engine_macro_fetch_failed", error=str(exc)) return {}