From 2b6cea57b27a9fc148f5354a1651b9bde5495192 Mon Sep 17 00:00:00 2001 From: I Luk Kim Date: Wed, 22 Apr 2026 23:33:23 -0700 Subject: [PATCH] Paper trader Phase 1 fixes: multi-session isolation, pipeline halt, snapshot refresh unblock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - v7.356 config: swap dataset_snapshot_id from manual_only ftb_fix_v2 to auto_full_rebuild base canonical so paper trader can refresh snapshot (root cause of processed_events=0 for 30 days) - Multi-session order isolation (1.A.2/1.A.3): tag client_order_id with pt-{session_id[:8]}-{uuid} prefix on all entry orders; _cancel_stale_orders filters by own session prefix so one session no longer ghost-cancels another's orders on shared Alpaca account - Pipeline halt on failure (1.B.1): _run_pipeline returns bool and stops on first subprocess failure instead of silently progressing with stale data - Daemon restart window skip (2.2): run_open/run_close only marked completed if processed_phases DB confirms prior execution — no more trading-less days after mid-day restart - event_parser: periodic batch commits every 500 docs (hypothesis fix for 3h hangs; unverified — may just be slow serial Oracle calls) - Tests updated for _verify_order_fill tuple return + new cross-session isolation test; all 23 paper_trader unit tests green Co-Authored-By: Claude Sonnet 4.6 --- apps/paper_trader/alpaca_broker.py | 8 +++- apps/paper_trader/auto.py | 43 ++++++++++++++++-- apps/paper_trader/engine.py | 23 ++++++++-- apps/pipeline/event_parser/main.py | 6 ++- ...max_long_v7.356_composed_gld_compound.json | 2 +- .../unit/paper_trader/test_reconciliation.py | 44 +++++++++++++------ 6 files changed, 101 insertions(+), 25 deletions(-) diff --git a/apps/paper_trader/alpaca_broker.py b/apps/paper_trader/alpaca_broker.py index 225cdb9..6600da8 100644 --- a/apps/paper_trader/alpaca_broker.py +++ b/apps/paper_trader/alpaca_broker.py @@ -30,6 +30,7 @@ class Order: status: str filled_avg_price: float | None filled_qty: int + client_order_id: str | None = None @dataclass @@ -115,7 +116,7 @@ class AlpacaBroker: # Orders # ------------------------------------------------------------------ # - def submit_market_buy(self, symbol: str, qty: int) -> Order: + def submit_market_buy(self, symbol: str, qty: int, client_order_id: str | None = None) -> Order: from alpaca.trading.requests import MarketOrderRequest from alpaca.trading.enums import OrderSide, TimeInForce @@ -124,11 +125,12 @@ class AlpacaBroker: qty=qty, side=OrderSide.BUY, time_in_force=TimeInForce.DAY, + **({"client_order_id": client_order_id} if client_order_id else {}), ) order = self._trading.submit_order(req) return self._to_order(order) - def submit_moc_buy(self, symbol: str, qty: int) -> Order: + def submit_moc_buy(self, symbol: str, qty: int, client_order_id: str | None = None) -> Order: """Submit a Market-on-Close buy order (fills at today's closing price).""" from alpaca.trading.requests import MarketOrderRequest from alpaca.trading.enums import OrderSide, TimeInForce @@ -138,6 +140,7 @@ class AlpacaBroker: qty=qty, side=OrderSide.BUY, time_in_force=TimeInForce.CLS, + **({"client_order_id": client_order_id} if client_order_id else {}), ) order = self._trading.submit_order(req) return self._to_order(order) @@ -376,6 +379,7 @@ class AlpacaBroker: status=str(order.status.value if hasattr(order.status, "value") else order.status), filled_avg_price=float(order.filled_avg_price) if order.filled_avg_price else None, filled_qty=int(float(order.filled_qty or 0)), + client_order_id=str(order.client_order_id) if getattr(order, "client_order_id", None) else None, ) @staticmethod diff --git a/apps/paper_trader/auto.py b/apps/paper_trader/auto.py index 618ab5a..efc4cd1 100644 --- a/apps/paper_trader/auto.py +++ b/apps/paper_trader/auto.py @@ -250,9 +250,13 @@ def _run_cmd(cmd: list[str], dry_run: bool, pipeline: bool = False) -> bool: return False -def _run_pipeline(cmds: list[list[str]], dry_run: bool) -> None: +def _run_pipeline(cmds: list[list[str]], dry_run: bool) -> bool: + """Run pipeline commands in order; abort and return False on first failure.""" for cmd in cmds: - _run_cmd(cmd, dry_run, pipeline=True) + if not _run_cmd(cmd, dry_run, pipeline=True): + _console.print(f" [red]Pipeline halted — downstream steps skipped[/]") + return False + return True def _run_paper(command: str, sessions: list[str], db: str, dry_run: bool) -> None: @@ -330,6 +334,24 @@ def _get_active_sessions(db: str) -> list[str]: return [] +def _all_sessions_ran_phase(sessions: list[str], db: str, date: dt.date, phase: str) -> bool: + """Return True if every session has a processed_phases record for (date, phase).""" + if not sessions: + return False + try: + from apps.paper_trader.state import StateManager + state = StateManager(db) + all_sessions = state.list_sessions() + session_map = {s.session_name: s.session_id for s in all_sessions} + return all( + state.is_phase_processed(session_map[name], date, phase) + for name in sessions + if name in session_map + ) + except Exception: + return False + + def run_auto(sessions: list[str], db: str, dry_run: bool) -> None: resolved = sessions or _get_active_sessions(db) if not resolved: @@ -380,11 +402,24 @@ def run_auto(sessions: list[str], db: str, dry_run: bool) -> None: next_td = _next_trading_day(today) _log(f"Non-trading day. Next trading day: [bold]{next_td}[/]") else: - # Mark events that already passed when starting mid-day as skipped + # On mid-day start: skip past pipeline events (catchup thread handles them). + # For trading phases (run_open/run_close), only skip if they actually ran + # today per processed_phases — prevents silently losing a trading window + # when the daemon restarts after a crash. + _phase_map = {"run_open": "next_open", "run_close": "reaction_close"} for ev in SCHEDULE: - if _et_dt_for(today, ev) <= now_et: + if _et_dt_for(today, ev) > now_et: + continue + if ev.kind in ("pipeline_pre", "pipeline_post"): completed.add(ev.name) _log(f"[dim]Skipping past event: {ev.description}[/]") + else: + db_phase = _phase_map.get(ev.kind, "") + if db_phase and _all_sessions_ran_phase(resolved, db, today, db_phase): + completed.add(ev.name) + _log(f"[dim]Skipping past event (already ran): {ev.description}[/]") + else: + _log(f"[yellow]Mid-day start: {ev.description} not yet run — will execute[/]") _print_schedule(now_et, completed) diff --git a/apps/paper_trader/engine.py b/apps/paper_trader/engine.py index d16d9e4..188991d 100644 --- a/apps/paper_trader/engine.py +++ b/apps/paper_trader/engine.py @@ -305,11 +305,20 @@ class PaperTradingEngine: # ------------------------------------------------------------------ # def _cancel_stale_orders(self) -> list[str]: - """Cancel all open orders. Daily system — any leftover is stale.""" + """Cancel open orders belonging to this session. Daily system — any leftover is stale. + + Uses client_order_id prefix (`pt-{session_id[:8]}-`) to distinguish this + session's orders from other sessions sharing the same Alpaca account. + Orders without a matching prefix are left untouched. + """ cancelled: list[str] = [] + session_prefix = f"pt-{self._session.session_id[:8]}-" try: open_orders = self._broker.list_orders("open") for order in open_orders: + coid = order.client_order_id or "" + if not coid.startswith(session_prefix): + continue # belongs to another session or is a parking order try: self._broker.cancel_order(order.id) cancelled.append(f"{order.symbol}:{order.id}") @@ -884,7 +893,9 @@ class PaperTradingEngine: }) continue try: - order = self._broker.submit_market_buy(candidate.symbol, plan.shares) + from uuid import uuid4 + coid = f"pt-{session_id[:8]}-{uuid4().hex[:8]}" + order = self._broker.submit_market_buy(candidate.symbol, plan.shares, client_order_id=coid) logger.info( "paper_engine_buy_submitted", symbol=candidate.symbol, @@ -1076,7 +1087,9 @@ class PaperTradingEngine: }) continue try: - order = self._broker.submit_market_buy(candidate.symbol, plan.shares) + from uuid import uuid4 + coid = f"pt-{session_id[:8]}-{uuid4().hex[:8]}" + order = self._broker.submit_market_buy(candidate.symbol, plan.shares, client_order_id=coid) except Exception as exc: self._state.record_processed_event( session_id, candidate.event_id, today.isoformat(), @@ -2824,7 +2837,9 @@ class PaperTradingEngine: rejected.append({"symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": "market_closed"}) continue try: - order = order_fn(candidate.symbol, plan.shares) + from uuid import uuid4 + coid = f"pt-{session_id[:8]}-{uuid4().hex[:8]}" + order = order_fn(candidate.symbol, plan.shares, client_order_id=coid) 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)) diff --git a/apps/pipeline/event_parser/main.py b/apps/pipeline/event_parser/main.py index 846bf73..7ea766d 100644 --- a/apps/pipeline/event_parser/main.py +++ b/apps/pipeline/event_parser/main.py @@ -111,7 +111,11 @@ async def run_event_parser( for sm in sm_result.scalars(): ticker_map[sm.symbol_id] = sm.ticker - for doc in docs: + for _batch_i, doc in enumerate(docs): + if _batch_i > 0 and _batch_i % 500 == 0: + await session.commit() + logger.info("event_parser_batch_commit", processed=_batch_i, total=stats["seen"]) + if not doc.accession_no: stats["invalid"] += 1 doc.parsed_status = "failed" diff --git a/configs/experiments/return_max_long_v7.356_composed_gld_compound.json b/configs/experiments/return_max_long_v7.356_composed_gld_compound.json index de8e330..4942379 100644 --- a/configs/experiments/return_max_long_v7.356_composed_gld_compound.json +++ b/configs/experiments/return_max_long_v7.356_composed_gld_compound.json @@ -1,6 +1,6 @@ { "experiment_name": "return_max_long_v7.356_composed_gld_compound", - "dataset_snapshot_id": "midlarge-liquid-long-v1_bucketfix_full_audit_canonical_ftb_fix_v2", + "dataset_snapshot_id": "midlarge-liquid-long-v1_bucketfix_full_audit_canonical", "description": "v7.118 + 6 per-engine exit/veto optimizations. 12개 엔진. Standalone SQS 92.4 (493% return).\nBEST conservative: tqqq_calm_v2 + sleeves = 3029.38% CW, SQS 89.2\nBEST aggressive: tqqq_active_v2 + sleeves = 3245.24% CW, SQS 89.2\nParking v3: gate_vol 0.35, TQQQ vol 0.22, temp 1.0, entropy 1.45 (1.15→1.45)\nLineage: v7.70→v7.110→v7.113→v7.114→v7.115→v7.116→v7.117→v7.118→v7.119", "base_config": "configs/backtest/return_max_long_v1.json", "overrides": { diff --git a/tests/unit/paper_trader/test_reconciliation.py b/tests/unit/paper_trader/test_reconciliation.py index 12d8c22..3e65010 100644 --- a/tests/unit/paper_trader/test_reconciliation.py +++ b/tests/unit/paper_trader/test_reconciliation.py @@ -270,9 +270,10 @@ class TestCancelStaleOrders: def test_cancels_open_orders(self, session, mock_broker, state_manager): engine = _make_engine(session, mock_broker, state_manager) + prefix = f"pt-{session.session_id[:8]}-" mock_broker.list_orders.return_value = [ - Order(id="ord1", symbol="AAPL", qty=10, side="buy", status="open", filled_avg_price=None, filled_qty=0), - Order(id="ord2", symbol="TSLA", qty=5, side="buy", status="open", filled_avg_price=None, filled_qty=0), + Order(id="ord1", symbol="AAPL", qty=10, side="buy", status="open", filled_avg_price=None, filled_qty=0, client_order_id=f"{prefix}aaa"), + Order(id="ord2", symbol="TSLA", qty=5, side="buy", status="open", filled_avg_price=None, filled_qty=0, client_order_id=f"{prefix}bbb"), ] cancelled = engine._cancel_stale_orders() assert len(cancelled) == 2 @@ -280,13 +281,28 @@ class TestCancelStaleOrders: def test_handles_cancel_failure(self, session, mock_broker, state_manager): engine = _make_engine(session, mock_broker, state_manager) + prefix = f"pt-{session.session_id[:8]}-" mock_broker.list_orders.return_value = [ - Order(id="ord1", symbol="AAPL", qty=10, side="buy", status="open", filled_avg_price=None, filled_qty=0), + Order(id="ord1", symbol="AAPL", qty=10, side="buy", status="open", filled_avg_price=None, filled_qty=0, client_order_id=f"{prefix}aaa"), ] mock_broker.cancel_order.side_effect = Exception("API error") cancelled = engine._cancel_stale_orders() assert cancelled == [] # failed to cancel + def test_skips_other_session_orders(self, session, mock_broker, state_manager): + """Orders from other sessions (different prefix) must not be cancelled.""" + engine = _make_engine(session, mock_broker, state_manager) + prefix = f"pt-{session.session_id[:8]}-" + mock_broker.list_orders.return_value = [ + Order(id="ord1", symbol="AAPL", qty=10, side="buy", status="open", filled_avg_price=None, filled_qty=0, client_order_id=f"{prefix}aaa"), + Order(id="ord2", symbol="TSLA", qty=5, side="buy", status="open", filled_avg_price=None, filled_qty=0, client_order_id="pt-othersess-bbb"), + Order(id="ord3", symbol="GOOG", qty=2, side="buy", status="open", filled_avg_price=None, filled_qty=0, client_order_id=None), + ] + cancelled = engine._cancel_stale_orders() + assert len(cancelled) == 1 + assert "AAPL" in cancelled[0] + assert mock_broker.cancel_order.call_count == 1 + def test_handles_list_orders_failure(self, session, mock_broker, state_manager): engine = _make_engine(session, mock_broker, state_manager) mock_broker.list_orders.side_effect = Exception("API down") @@ -304,9 +320,9 @@ class TestVerifyOrderFill: id="ord1", symbol="AAPL", qty=10, side="buy", status="filled", filled_avg_price=101.5, filled_qty=10, ) - result = engine._verify_order_fill("ord1", "AAPL", timeout_sec=0.1) - assert result is not None - assert result.filled_avg_price == 101.5 + order, reason = engine._verify_order_fill("ord1", "AAPL", timeout_sec=0.1) + assert order is not None + assert order.filled_avg_price == 101.5 def test_rejected_order(self, session, mock_broker, state_manager): engine = _make_engine(session, mock_broker, state_manager) @@ -314,8 +330,9 @@ class TestVerifyOrderFill: id="ord1", symbol="AAPL", qty=10, side="buy", status="rejected", filled_avg_price=None, filled_qty=0, ) - result = engine._verify_order_fill("ord1", "AAPL", timeout_sec=0.1) - assert result is None + order, reason = engine._verify_order_fill("ord1", "AAPL", timeout_sec=0.1) + assert order is None + assert "rejected" in reason def test_timeout(self, session, mock_broker, state_manager): engine = _make_engine(session, mock_broker, state_manager) @@ -323,8 +340,9 @@ class TestVerifyOrderFill: id="ord1", symbol="AAPL", qty=10, side="buy", status="new", filled_avg_price=None, filled_qty=0, ) - result = engine._verify_order_fill("ord1", "AAPL", timeout_sec=0.1) - assert result is None + order, reason = engine._verify_order_fill("ord1", "AAPL", timeout_sec=0.1) + assert order is None + assert "timeout" in reason def test_api_error_retries(self, session, mock_broker, state_manager): engine = _make_engine(session, mock_broker, state_manager) @@ -333,9 +351,9 @@ class TestVerifyOrderFill: Order(id="ord1", symbol="AAPL", qty=10, side="buy", status="filled", filled_avg_price=102.0, filled_qty=10), ] - result = engine._verify_order_fill("ord1", "AAPL", timeout_sec=2.0) - assert result is not None - assert result.filled_avg_price == 102.0 + order, reason = engine._verify_order_fill("ord1", "AAPL", timeout_sec=2.0) + assert order is not None + assert order.filled_avg_price == 102.0 # ── _check_kill_switch ──────────────────────────────────────────────────────