Use SnapshotStore for paper trader next_open candidate selection

PaperTradingEngine now accepts an optional SnapshotStore and uses it
for run_next_open candidate fetching, ensuring live candidate selection
matches the backtester's pre-computed scores exactly. run_reaction_close
keeps EventDetector for real-time intraday event detection. Adds
load_snapshot_store_for_session() helper with auto-refresh logic.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
main
I Luk Kim 4 months ago
parent 969dedc635
commit e415743444

@ -571,3 +571,64 @@ def run_backtest(
)
return results
def load_snapshot_store_for_session(
config: Any,
oracle_url: str,
db_dsn: str,
*,
auto_refresh: bool = True,
) -> Any:
"""Load SnapshotStore for a live paper trading session.
Mirrors the snapshot loading done in run_backtest(), but for live use.
Returns None if snapshot_id is not configured or loading fails (caller
should fall back to EventDetector).
This is a SYNC function safe to call from _make_engine() / make_engine()
before an event loop is started.
"""
import asyncio
from libs.backtest.snapshot_store import SnapshotStore
snapshot_id = getattr(config, "dataset_snapshot_id", None)
if not snapshot_id:
logger.warning("snapshot_store_no_snapshot_id", config=str(config))
return None
try:
if auto_refresh and _snapshot_needs_refresh(snapshot_id, dt.date.today()):
universe_profile = None
if "midlarge" in snapshot_id:
universe_profile = "midlarge-liquid-long-v1"
elif "midwide" in snapshot_id:
universe_profile = "midwide-liquid-long-v1"
elif "smallcap" in snapshot_id:
universe_profile = "smallcap-liquid-long-v1"
try:
asyncio.run(_refresh_snapshot(snapshot_id, universe_profile, manual=False))
except Exception as exc:
logger.warning("snapshot_store_refresh_failed", error=str(exc))
if not _snapshot_has_required_coverage(snapshot_id, dt.date.today()):
logger.warning("snapshot_store_no_coverage_after_refresh_failure")
return None
snapshot_path = _resolve_snapshot_path(snapshot_id)
if snapshot_path is None:
logger.warning("snapshot_store_path_not_found", snapshot_id=snapshot_id)
return None
from apps.backtester.run import _resolve_scoring_fn
scoring_fn = _resolve_scoring_fn(config)
return SnapshotStore.load_merged(
snapshot_dir=snapshot_path,
split_names=["train", "valid", "test"],
oracle_url=oracle_url,
db_dsn=db_dsn,
scoring_fn=scoring_fn,
)
except Exception as exc:
logger.warning("snapshot_store_load_failed", error=str(exc))
return None

@ -185,11 +185,19 @@ def _make_engine(session, args):
broker = _get_broker()
from apps.paper_trader.event_detector import EventDetector
from apps.paper_trader.engine import PaperTradingEngine
from apps.backtester.run import load_manifest, resolve_config
from apps.paper_trader.backtest_sim import load_snapshot_store_for_session
oracle_url = os.environ.get("ORACLE_URL") or os.environ.get("STOCK_ORACLE_URL", "http://localhost:8000")
db_dsn = os.environ.get("DB_DSN") or os.environ.get("POSTGRES_DSN", "")
detector = EventDetector(db_dsn=db_dsn, oracle_url=oracle_url)
state = _get_state_manager(args.db)
return PaperTradingEngine(session=session, broker=broker, state=state, event_detector=detector)
manifest = load_manifest(session.config_path)
config = resolve_config(manifest)
snapshot_store = load_snapshot_store_for_session(config, oracle_url, db_dsn)
return PaperTradingEngine(
session=session, broker=broker, state=state, event_detector=detector,
snapshot_store=snapshot_store,
)
def cmd_run_close(args: argparse.Namespace) -> None:

@ -80,11 +80,13 @@ class PaperTradingEngine:
broker: AlpacaBroker,
state: StateManager,
event_detector: EventDetector,
snapshot_store: "SnapshotStore | None" = None,
) -> None:
self._session = session
self._broker = broker
self._state = state
self._detector = event_detector
self._snapshot_store = snapshot_store
manifest = load_manifest(session.config_path)
self._config: BacktestConfig = resolve_config(manifest)
@ -129,6 +131,7 @@ class PaperTradingEngine:
self._lookback_injected: bool = False
# Overlay shock brake cooldown (in-memory, session-scoped)
self._parking_brake_cooldown_remaining: int = 0
self._parking_brake_skip_buy_today: bool = False # skip same-day re-buy after brake fires
def _get_candidate_capital_bucket_id(self, candidate: Candidate) -> str | None:
return candidate.engine_capital_bucket_id
@ -382,7 +385,7 @@ class PaperTradingEngine:
)
return report
def _verify_order_fill(self, order_id: str, symbol: str, timeout_sec: float = 2.0) -> Order | None:
def _verify_order_fill(self, order_id: str, symbol: str, timeout_sec: float = 15.0) -> Order | None:
"""Poll broker to verify order fill. Returns filled Order or None."""
deadline = time.monotonic() + timeout_sec
while time.monotonic() < deadline:
@ -740,14 +743,6 @@ class PaperTradingEngine:
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 == "insufficient_cash":
# Attempt to free parking cash before giving up
needed = plan.shares * float(candidate.entry_price_est) if plan.shares else float(
@ -777,6 +772,10 @@ class PaperTradingEngine:
)
if plan.skip_reason:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=plan.skip_reason,
)
rejected.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -800,6 +799,10 @@ class PaperTradingEngine:
symbol=candidate.symbol,
error=str(exc),
)
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=f"order_failed:{exc}",
)
rejected.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -811,6 +814,10 @@ class PaperTradingEngine:
# Verify fill
verified = self._verify_order_fill(order.id, candidate.symbol)
if verified is None:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason="order_not_filled",
)
rejected.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -841,6 +848,9 @@ class PaperTradingEngine:
status="open",
),
)
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(), "entered",
)
trade_risk_state = candidate_portfolio_state.sizing_equity or candidate_portfolio_state.equity
trade_risk = trade_risk_state * (
@ -923,14 +933,11 @@ class PaperTradingEngine:
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:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=plan.skip_reason,
)
rejected.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -942,6 +949,10 @@ class PaperTradingEngine:
try:
order = self._broker.submit_market_buy(candidate.symbol, plan.shares)
except Exception as exc:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=f"order_failed:{exc}",
)
rejected.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -952,6 +963,10 @@ class PaperTradingEngine:
verified = self._verify_order_fill(order.id, candidate.symbol)
if verified is None:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason="order_not_filled",
)
rejected.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -981,6 +996,9 @@ class PaperTradingEngine:
status="open",
),
)
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(), "entered",
)
entries.append({
"symbol": candidate.symbol,
"event_type": candidate.event_type,
@ -1250,14 +1268,25 @@ class PaperTradingEngine:
return True
# Signal 4: near-SMA buffer — exit when QQQ within sma_buffer% ABOVE SMA10 (pre-emptive)
# Requires vol5/vol20 ∈ (0.25, rv_ratio_upper) — "barely elevated" pre-crash signature.
# Upper bound filters out high-vol days (regular gate handles those) and false alarms.
sma_buffer = getattr(risk, "cash_parking_overlay_shock_brake_sma_buffer", 0.0)
if sma_buffer > 0 and n >= 11:
qqq_close_now = qqq_closes[-1]
qqq_sma10 = sum(qqq_closes[-10:]) / 10
if qqq_sma10 > 0:
sma_gap = (qqq_close_now - qqq_sma10) / qqq_sma10
if 0 < sma_gap < sma_buffer:
return True
qqq_sma10_pt = sum(qqq_closes[-10:]) / 10
vol5_pt = _vol(5)
vol20_pt = _vol(20)
if (
qqq_sma10_pt > 0
and vol5_pt is not None and vol20_pt is not None and vol20_pt > 0
):
rv = vol5_pt / vol20_pt
rv_upper = getattr(risk, "cash_parking_overlay_shock_brake_rv_ratio_upper", 0.0)
vol_ok = rv > 0.25 and (rv_upper <= 0 or rv < rv_upper)
if vol_ok:
sma_gap = (qqq_close_now - qqq_sma10_pt) / qqq_sma10_pt
if 0 < sma_gap < sma_buffer:
return True
return False
@ -1320,6 +1349,7 @@ class PaperTradingEngine:
# Set cooldown for next active parking state (will be created on re-buy)
# We store cooldown on a session-level attribute for now
self._parking_brake_cooldown_remaining = risk.cash_parking_overlay_shock_brake_cooldown_days
self._parking_brake_skip_buy_today = True # skip same-day re-buy
return True
# --- Dwell cap check ---
@ -1580,6 +1610,11 @@ class PaperTradingEngine:
return # For now, no actual top-up execution (matches backtester behavior of holding)
# --- New parking position ---
# Skip same-day re-buy if brake fired today; next day gate re-evaluates fresh.
if self._parking_brake_skip_buy_today:
self._parking_brake_skip_buy_today = False
return
target = self._parking_evaluate_gate(today)
# Overlay brake cooldown: suppress TQQQ overlay re-entry during cooldown
@ -1755,9 +1790,23 @@ class PaperTradingEngine:
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 self._snapshot_store is not None:
from libs.backtest.calendar import get_trading_days
start_lb = self._lookback_start_date(today)
lookback_rows = []
for lb_date in get_trading_days(start_lb, today):
if lb_date >= today:
continue
for row in self._snapshot_store.get_candidates_for_date(lb_date):
row = dict(row)
row["is_lookback_entry"] = True
tdays = get_trading_days(lb_date, today)
row["lookback_days_elapsed"] = max(0, len(tdays) - 1)
lookback_rows.append(row)
else:
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(
@ -1782,9 +1831,12 @@ class PaperTradingEngine:
# tomorrow with entry_convention='next_open_after_reaction_close'.
# ENB (after-close with reaction_close convention, entry_date=reaction_date) is
# correctly included because event_date != reaction_date (not same_day).
all_rows = await self._detector.get_candidates_for_date(
today, self._config, convention=None
)
if self._snapshot_store is not None:
all_rows = self._snapshot_store.get_candidates_for_date(today)
else:
all_rows = await self._detector.get_candidates_for_date(
today, self._config, convention=None
)
next_open_rows = [
r for r in all_rows
if not (self._is_same_day_event(r) and r.get("entry_convention") == "reaction_close")
@ -2216,11 +2268,6 @@ class PaperTradingEngine:
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 == "insufficient_cash":
# Attempt to free parking cash before giving up
needed = plan.shares * float(candidate.entry_price_est) if plan.shares else float(candidate.entry_price_est)
@ -2246,6 +2293,10 @@ class PaperTradingEngine:
)
if plan.skip_reason:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=plan.skip_reason,
)
rejected.append({
"symbol": candidate.symbol, "event_type": candidate.event_type,
"score": candidate.score, "reason": plan.skip_reason,
@ -2260,6 +2311,10 @@ class PaperTradingEngine:
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:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=gap_reason,
)
rejected.append({
"symbol": candidate.symbol, "event_type": candidate.event_type,
"score": candidate.score, "reason": gap_reason,
@ -2270,6 +2325,10 @@ class PaperTradingEngine:
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))
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason=f"order_failed:{exc}",
)
rejected.append({"symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": f"order_failed:{exc}"})
continue
@ -2278,6 +2337,10 @@ class PaperTradingEngine:
if not is_moc:
verified = self._verify_order_fill(order.id, candidate.symbol)
if verified is None:
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(),
"rejected", skip_reason="order_not_filled",
)
rejected.append({"symbol": candidate.symbol, "event_type": candidate.event_type, "score": candidate.score, "reason": "order_not_filled"})
continue
fill_price = verified.filled_avg_price or plan.entry_price_limit
@ -2297,6 +2360,9 @@ class PaperTradingEngine:
status="open",
),
)
self._state.record_processed_event(
session_id, candidate.event_id, today.isoformat(), "entered",
)
trade_risk_state = candidate_portfolio_state.sizing_equity or candidate_portfolio_state.equity
trade_risk = trade_risk_state * (
candidate.engine_per_trade_risk_pct or self._config.risk.per_trade_risk_pct

@ -48,6 +48,8 @@ def make_engine(session: Any, db_path: str) -> Any:
from apps.paper_trader.engine import PaperTradingEngine
from apps.paper_trader.event_detector import EventDetector
from apps.paper_trader.state import StateManager
from apps.backtester.run import load_manifest, resolve_config
from apps.paper_trader.backtest_sim import load_snapshot_store_for_session
broker = AlpacaBroker.from_env()
oracle_url = (
@ -57,8 +59,12 @@ def make_engine(session: Any, db_path: str) -> Any:
db_dsn = os.environ.get("DB_DSN") or os.environ.get("POSTGRES_DSN", "")
state = StateManager(db_path)
detector = EventDetector(db_dsn=db_dsn, oracle_url=oracle_url)
manifest = load_manifest(session.config_path)
config = resolve_config(manifest)
snapshot_store = load_snapshot_store_for_session(config, oracle_url, db_dsn)
return PaperTradingEngine(
session=session, broker=broker, state=state, event_detector=detector
session=session, broker=broker, state=state, event_detector=detector,
snapshot_store=snapshot_store,
)
@ -581,11 +587,27 @@ class AutoScheduler:
now_et = self._now_et()
today = now_et.date()
# New day → reset completed set
# New day → reset completed set and refresh session list
if last_schedule_date != today:
self._completed.clear()
last_schedule_date = today
self._log(f"━━━ {today.strftime('%a %Y-%m-%d')} ━━━")
# Re-read active sessions from DB so new sessions are picked up
# and removed/closed sessions are dropped automatically.
if not self._sessions: # only auto-refresh if not pinned via --session
fresh = self._get_active_sessions()
if fresh != resolved:
added = set(fresh) - set(resolved)
removed = set(resolved) - set(fresh)
if added:
self._log(f"Sessions added: {', '.join(sorted(added))}")
if removed:
self._log(f"Sessions removed: {', '.join(sorted(removed))}")
resolved = fresh
if not resolved:
self._log("No active sessions. Waiting for next day.")
self._log(f"━━━ {today.strftime('%a %Y-%m-%d')} — sessions: {', '.join(resolved) or 'none'} ━━━")
if not self._is_trading_day(today):
next_td = self._next_trading_day(today)

Loading…
Cancel
Save