You cannot select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

634 lines
26 KiB
Python

This file contains ambiguous Unicode characters!

This file contains ambiguous Unicode characters that may be confused with others in your current locale. If your use case is intentional and legitimate, you can safely ignore this warning. Use the Escape button to highlight these characters.

"""ORB paper trading in-process service layer.
Provides:
- ORBAutoScheduler: asyncio-based intraday scheduler
- orb_auto_scheduler: module-level singleton
- State persistence to .orb_auto_state.json
Schedule is computed dynamically from strategy params (orb_minutes,
order_timeout_minutes, sim_bar_minutes) rather than being hardcoded.
Daily event sequence:
9:30 + orb_minutes → orb_detect (daily bars + 5-min ORB bars → candidates)
orb_end + 1 min → breakout (Oracle snapshot price check, every 1 min)
...repeat until order_timeout_minutes elapses
orb_end + N×sim_bar → stop_check (fetch 5-min bars, aggregate, check stops)
...repeat until 15:50
15:55 → eod_exit
16:00 → post_close
"""
from __future__ import annotations
import asyncio
import datetime as dt
import json
import logging
import traceback
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from zoneinfo import ZoneInfo
log = logging.getLogger(__name__)
_TZ_ET = ZoneInfo("America/New_York")
_TZ_PHOENIX = ZoneInfo("America/Phoenix") # UTC-7 always (no DST)
_DEFAULT_DB = "data/paper/orb.db"
_MARKET_OPEN = dt.time(9, 30) # ET
_EOD_EXIT = dt.time(15, 55) # ET
_POST_CLOSE = dt.time(16, 0) # ET
# ── Dynamic schedule builder ──────────────────────────────────────────────────
def build_schedule(
date: dt.date,
orb_minutes: int = 10,
order_timeout_minutes: int = 45,
sim_bar_minutes: int = 90,
) -> list[dict[str, Any]]:
"""Build today's event list from strategy parameters.
Args:
date: The trading date.
orb_minutes: Length of the ORB window (e.g. 10 → window is 9:309:40).
order_timeout_minutes: How many minutes after ORB end to keep looking for
breakouts (one check per minute).
sim_bar_minutes: Interval between stop-management checks after breakout
window closes (matches the sim_bar_minutes strategy param).
Returns:
List of event dicts sorted chronologically. Each dict has:
name, label, kind, et_dt (aware datetime in ET)
Note: breakout checks run every sim_bar_minutes (matching the backtest bar
aggregation interval), up to order_timeout_minutes after *market open*
(not after ORB end) — consistent with orb_simulator.py timeout semantics.
"""
mkt_open = dt.datetime(date.year, date.month, date.day, 9, 30, tzinfo=_TZ_ET)
orb_end = mkt_open + dt.timedelta(minutes=orb_minutes)
eod = dt.datetime(date.year, date.month, date.day, 15, 55, tzinfo=_TZ_ET)
post = dt.datetime(date.year, date.month, date.day, 16, 0, tzinfo=_TZ_ET)
events: list[dict[str, Any]] = []
# ── ORB window monitoring: 9:30 → orb_end-1, every minute (no-op) ─────────
for i in range(orb_minutes):
t = mkt_open + dt.timedelta(minutes=i)
events.append({
"name": f"orb_monitor_{i + 1}",
"label": f"ORB 윈도우 ({t.strftime('%H:%M')} ET, {i + 1}/{orb_minutes}분)",
"kind": "orb_monitor",
"et_dt": t,
})
# ── ORB detection at window close (fetch bars + rank candidates) ──────────
events.append({
"name": "orb_detect",
"label": f"ORB 감지 ({orb_end.strftime('%H:%M')} ET)",
"kind": "orb_detect",
"et_dt": orb_end,
})
# ── Breakout checks: every sim_bar_minutes from first bar close to timeout ──
# Timeout is measured from market open (matches orb_simulator.py semantics).
timeout_dt = mkt_open + dt.timedelta(minutes=order_timeout_minutes)
i = 1
t = orb_end + dt.timedelta(minutes=sim_bar_minutes)
while t <= timeout_dt:
if t >= eod:
break
events.append({
"name": f"breakout_{i}",
"label": f"브레이크아웃 ({t.strftime('%H:%M')} ET, +{i * sim_bar_minutes}분)",
"kind": "breakout",
"et_dt": t,
})
t += dt.timedelta(minutes=sim_bar_minutes)
i += 1
# ── Stop checks: every sim_bar_minutes from orb_end ───────────────────────
t = orb_end + dt.timedelta(minutes=sim_bar_minutes)
idx = 1
while t < eod:
events.append({
"name": f"stop_{idx}",
"label": f"스톱 체크 ({t.strftime('%H:%M')} ET, +{sim_bar_minutes}분)",
"kind": "stop_check",
"et_dt": t,
})
t += dt.timedelta(minutes=sim_bar_minutes)
idx += 1
# ── EOD + post-close ──────────────────────────────────────────────────────
events.append({"name": "eod_exit", "label": "EOD 청산 (15:55 ET)",
"kind": "eod_exit", "et_dt": eod})
events.append({"name": "post_close", "label": "마감 후 스냅샷 (16:00 ET)",
"kind": "post_close", "et_dt": post})
return sorted(events, key=lambda e: e["et_dt"])
def _load_session_params(db_path: str, session_name: str) -> dict[str, Any]:
"""Load strategy params (orb_minutes, sim_bar_minutes, order_timeout_minutes)
from the session's config YAML. Returns defaults on any error.
"""
defaults = {"orb_minutes": 10, "sim_bar_minutes": 90, "order_timeout_minutes": 45}
try:
import yaml
from apps.orb_trader.state import ORBStateManager
session = ORBStateManager(db_path).get_session(session_name)
if session is None:
return defaults
raw = yaml.safe_load(Path(session.config_path).read_text()) or {}
orb = raw.get("orb_strategy", {})
return {
"orb_minutes": orb.get("orb_minutes", defaults["orb_minutes"]),
"sim_bar_minutes": orb.get("sim_bar_minutes", defaults["sim_bar_minutes"]),
"order_timeout_minutes": orb.get("order_timeout_minutes", defaults["order_timeout_minutes"]),
}
except Exception:
return defaults
# ── State file helpers ────────────────────────────────────────────────────────
def _state_file_path(db_path: str) -> Path:
return Path(db_path).parent / ".orb_auto_state.json"
def load_orb_saved_state(db_path: str) -> dict[str, Any] | None:
sf = _state_file_path(db_path)
if sf.exists():
try:
return json.loads(sf.read_text())
except Exception:
pass
return None
# ── ORBAutoScheduler ──────────────────────────────────────────────────────────
class ORBAutoScheduler:
"""In-process ORB intraday auto-trader.
Schedule is built dynamically each trading day from the sessions' strategy
params (orb_minutes, order_timeout_minutes, sim_bar_minutes).
"""
def __init__(self) -> None:
self._task: asyncio.Task | None = None # type: ignore[type-arg]
self._sessions: list[str] = []
self._db_path: str = _DEFAULT_DB
self._dry_run: bool = False
self._log_lines: list[str] = []
self._completed: set[str] = set()
self._today_schedule: list[dict[str, Any]] = [] # built each trading day
self._engines: dict[str, Any] = {}
self._engine_date: str = ""
@property
def running(self) -> bool:
return self._task is not None and not self._task.done()
def _log_file_path(self) -> Path:
return Path(self._db_path).parent / "orb_scheduler.log"
def _load_persisted_log(self, max_lines: int = 3000) -> list[str]:
"""Load existing log lines from file (survives server restarts).
Trims to last max_lines if the file has grown too large.
"""
try:
lf = self._log_file_path()
if lf.exists():
lines = lf.read_text(encoding="utf-8").splitlines()
if len(lines) > max_lines:
# Keep last max_lines; rewrite trimmed file
lines = lines[-max_lines:]
try:
lf.write_text("\n".join(lines) + "\n", encoding="utf-8")
except Exception:
pass
return lines
except Exception:
pass
return []
def start(self, sessions: list[str], db_path: str, dry_run: bool = False) -> None:
if self.running:
raise RuntimeError("ORBAutoScheduler already running")
self._sessions = sessions
self._db_path = db_path
self._dry_run = dry_run
self._log_lines = self._load_persisted_log() # restore previous logs
self._completed = set()
self._today_schedule = []
self._engines = {}
self._engine_date = ""
self._save_state()
self._task = asyncio.create_task(self._run_loop())
def stop(self) -> None:
if self._task and not self._task.done():
self._task.cancel()
self._clear_state()
def shutdown(self) -> None:
"""Server shutdown — cancel task but keep state for auto-restart."""
if self._task and not self._task.done():
self._task.cancel()
def get_log(self, lines: int = 200) -> str:
return "\n".join(self._log_lines[-lines:])
def get_log_tail(self, lines: int = 80) -> list[str]:
return self._log_lines[-lines:]
@property
def log_line_count(self) -> int:
return len(self._log_lines)
def get_status(self) -> dict[str, Any]:
now_et = self._now_et()
today = now_et.date()
schedule_view = []
for ev in self._today_schedule:
ev_dt = ev["et_dt"]
past = ev_dt <= now_et
wait = (ev_dt - now_et).total_seconds()
schedule_view.append({
"name": ev["name"],
"kind": ev["kind"],
"session": ev.get("session", ""),
"label": ev["label"],
"et_time": ev_dt.strftime("%H:%M ET"),
"et_iso": ev_dt.isoformat(),
"past": past,
"done": ev["name"] in self._completed,
"wait_secs": max(0, wait),
})
return {
"running": self.running,
"sessions": self._sessions,
"dry_run": self._dry_run,
"schedule": schedule_view,
"log_tail": list(self._log_lines), # full log, not truncated
"log_line_count": self.log_line_count,
}
# ── State persistence ──────────────────────────────────────────────────────
def _save_state(self) -> None:
try:
_state_file_path(self._db_path).write_text(json.dumps({
"running": True,
"sessions": self._sessions,
"dry_run": self._dry_run,
"db_path": self._db_path,
}))
except Exception:
pass
def _clear_state(self) -> None:
try:
sf = _state_file_path(self._db_path)
if sf.exists():
sf.unlink()
except Exception:
pass
# ── Helpers ────────────────────────────────────────────────────────────────
def _log(self, msg: str) -> None:
ts = datetime.now(tz=_TZ_PHOENIX).strftime("%H:%M MST")
line = f"{ts} {msg}"
self._log_lines.append(line)
log.info("[ORBScheduler] %s", msg)
# Persist to file so logs survive server restarts
try:
with self._log_file_path().open("a", encoding="utf-8") as f:
f.write(line + "\n")
except Exception:
pass
def _now_et(self) -> dt.datetime:
return dt.datetime.now(tz=_TZ_ET)
def _is_trading_day(self, date: dt.date) -> bool:
try:
from libs.common.time_utils import is_trading_day
return is_trading_day(date)
except Exception:
return date.weekday() < 5
def _next_trading_day(self, from_date: dt.date) -> dt.date:
check = from_date + dt.timedelta(days=1)
for _ in range(14):
if self._is_trading_day(check):
return check
check += dt.timedelta(days=1)
raise RuntimeError("No trading day found in next 14 days")
@staticmethod
def _fmt_countdown(seconds: float) -> str:
if seconds <= 0:
return "now"
h = int(seconds // 3600)
m = int((seconds % 3600) // 60)
s = int(seconds % 60)
if h > 0:
return f"{h}h {m:02d}m"
if m > 0:
return f"{m}m {s:02d}s"
return f"{s}s"
def _get_active_sessions(self) -> list[str]:
try:
from apps.orb_trader.state import ORBStateManager
return [
s.session_name
for s in ORBStateManager(self._db_path).list_sessions()
if s.status == "active"
]
except Exception:
return []
def _build_today_schedule(self, date: dt.date, sessions: list[str]) -> list[dict[str, Any]]:
"""Build per-session schedules and merge into one sorted timeline.
Each event is tagged with its session so the loop can dispatch
the right engine for the right session at the right time.
Sessions that no longer exist in the DB are skipped with a warning.
"""
from apps.orb_trader.state import ORBStateManager
state_mgr = ORBStateManager(self._db_path)
combined: list[dict[str, Any]] = []
valid_sessions: list[str] = []
for session_name in sessions:
# Verify session still exists in DB
session_obj = state_mgr.get_session(session_name)
if session_obj is None:
self._log(f" WARNING: session '{session_name}' not found in DB — skipping")
continue
valid_sessions.append(session_name)
params = _load_session_params(self._db_path, session_name)
events = build_schedule(date, **params)
for ev in events:
ev["session"] = session_name
ev["name"] = f"{session_name}:{ev['name']}"
combined.extend(events)
self._log(
f" {session_name}: orb={params['orb_minutes']}min, "
f"timeout={params['order_timeout_minutes']}min, "
f"bar={params['sim_bar_minutes']}min → {len(events)} events"
)
if not valid_sessions and sessions:
# All specified sessions are gone — fall back to all active sessions
active = self._get_active_sessions()
self._log(
f" All specified sessions missing; falling back to active sessions: "
f"{', '.join(active) or 'none'}"
)
return self._build_today_schedule(date, active)
combined.sort(key=lambda e: (e["et_dt"], e["session"]))
self._log(f"Total schedule: {len(combined)} events across {len(valid_sessions)} session(s)")
return combined
def _get_or_create_engine(self, session_name: str, date_str: str) -> Any:
if self._engine_date != date_str:
self._engines = {}
self._engine_date = date_str
if session_name not in self._engines:
try:
from apps.orb_trader.state import ORBStateManager
from apps.orb_trader.engine import make_orb_engine
state_mgr = ORBStateManager(self._db_path)
session = state_mgr.get_session(session_name)
if session is None:
return None
self._engines[session_name] = make_orb_engine(session, self._db_path)
except Exception as exc:
self._log(f" ERROR creating engine for {session_name}: {exc}")
return None
return self._engines[session_name]
# ── Engine operation runners ───────────────────────────────────────────────
async def _run_trading(self, kind: str, sessions: list[str], date_str: str) -> None:
# ORB 윈도우 모니터링은 no-op (범위 형성 중, 장중 데이터는 orb_detect에서 일괄 fetch)
if kind == "orb_monitor":
self._log(f" ORB 윈도우 모니터링 중...")
return
from apps.orb_trader.state import ORBStateManager
state_mgr = ORBStateManager(self._db_path)
for session_name in sessions:
session = state_mgr.get_session(session_name)
if session is None or session.status != "active":
self._log(f" {session_name}: skipped (not active)")
continue
if self._dry_run:
self._log(f" [DRY] {kind}{session_name}")
continue
engine = self._get_or_create_engine(session_name, date_str)
if engine is None:
continue
self._log(f"{kind}{session_name}")
try:
def _run_sync(e=engine, k=kind, d=date_str) -> dict[str, Any]:
import asyncio as _asyncio
loop = _asyncio.new_event_loop()
_asyncio.set_event_loop(loop)
try:
if k == "orb_detect":
return e.run_orb_detection(d)
elif k == "breakout":
return e.run_breakout_check(d)
elif k == "stop_check":
return e.run_stop_check(d)
elif k == "eod_exit":
return e.run_eod_exit(d)
elif k == "post_close":
return e.run_post_close(d)
return {}
finally:
loop.close()
summary = await asyncio.to_thread(_run_sync)
self._log(f"{session_name}: {summary}")
except Exception as exc:
tb = traceback.format_exc()
self._log(f" ERROR {session_name}: {exc}")
log.error("ORB engine error: %s\n%s", exc, tb)
# ── Run-now: manually trigger detection for a late-added session ─────────
async def run_session_now(self, session_name: str) -> dict[str, Any]:
"""Run ORB detection immediately for a session that missed the morning window.
After detection, injects the session's remaining today-events into the
active schedule so breakout checks, stop checks, and EOD exit still fire.
Returns a summary dict or {"error": "<reason>"} on failure.
"""
from apps.orb_trader.state import ORBStateManager
now_et = self._now_et()
today = now_et.date()
date_str = today.isoformat()
state_mgr = ORBStateManager(self._db_path)
session = state_mgr.get_session(session_name)
if session is None:
return {"error": f"Session '{session_name}' not found"}
# Already ran if daily_state.phase is set (engine wrote it during orb_detect)
daily = state_mgr.get_daily_state(session.session_id, date_str)
if daily.phase:
return {"error": f"Session '{session_name}' already ran today (phase={daily.phase})"}
self._log(f"🔄 지금 시작: {session_name} — ORB 감지 실행 중 (현재 시세 기준)...")
# Run detection now (blocking in thread so we don't stall the event loop)
await self._run_trading("orb_detect", [session_name], date_str)
# Inject remaining future events into the live schedule
params = _load_session_params(self._db_path, session_name)
all_events = build_schedule(today, **params)
for ev in all_events:
ev["session"] = session_name
ev["name"] = f"{session_name}:{ev['name']}"
existing_names = {e["name"] for e in self._today_schedule}
injected = 0
for ev in all_events:
if ev["et_dt"] <= now_et:
# Past events: mark as completed so the loop skips them
self._completed.add(ev["name"])
elif ev["name"] not in existing_names:
self._today_schedule.append(ev)
injected += 1
self._today_schedule.sort(key=lambda e: (e["et_dt"], e.get("session", "")))
self._log(f"{session_name}: {injected} 이벤트 추가됨 (오늘 남은 일정 포함)")
return {"session": session_name, "injected": injected}
# ── Main scheduler loop ────────────────────────────────────────────────────
async def _run_loop(self) -> None:
resolved = self._sessions or self._get_active_sessions()
if not resolved:
self._log("No active ORB sessions. Stopping.")
return
self._log(f"ORB auto-scheduler started — sessions: {', '.join(resolved)}")
if self._dry_run:
self._log("DRY RUN — orders will not be placed")
last_schedule_date: dt.date | None = None
try:
while True:
now_et = self._now_et()
today = now_et.date()
date_str = today.isoformat()
# ── New trading day: rebuild schedule ─────────────────────────
if last_schedule_date != today:
self._completed.clear()
last_schedule_date = today
if not self._sessions:
fresh = self._get_active_sessions()
if set(fresh) != set(resolved):
self._log(f"Sessions refreshed: {', '.join(fresh) or 'none'}")
resolved = fresh
self._log(
f"━━━ {today.strftime('%a %Y-%m-%d')} "
f"— sessions: {', '.join(resolved) or 'none'} ━━━"
)
if not self._is_trading_day(today):
self._today_schedule = []
next_td = self._next_trading_day(today)
self._log(f"Non-trading day. Next: {next_td}")
else:
self._today_schedule = self._build_today_schedule(today, resolved)
# Skip events already past
for ev in self._today_schedule:
if ev["et_dt"] <= now_et:
self._completed.add(ev["name"])
self._log(f"Past (skipped): {ev['label']}")
if not self._is_trading_day(today):
await asyncio.sleep(1800)
continue
if not resolved:
await asyncio.sleep(300)
continue
pending = [ev for ev in self._today_schedule if ev["name"] not in self._completed]
if not pending:
next_td = self._next_trading_day(today)
# Wake up just before the ORB detect event of next trading day
# (schedule isn't built yet, assume market open 9:30 ET)
wake_et = dt.datetime(
next_td.year, next_td.month, next_td.day, 9, 25, tzinfo=_TZ_ET
)
wait = (wake_et - now_et).total_seconds()
self._log(
f"All done today. Sleeping until "
f"{wake_et.strftime('%I:%M %p ET')} on {next_td} "
f"({self._fmt_countdown(wait)})"
)
await asyncio.sleep(min(wait, 3600))
continue
next_ev = pending[0]
next_et = next_ev["et_dt"]
wait = (next_et - now_et).total_seconds()
if wait > 90:
label = next_ev.get("label", next_ev["name"])
sess = next_ev.get("session", "")
self._log(f"Next: [{sess}] {label}{self._fmt_countdown(wait)}")
await asyncio.sleep(min(wait - 60, 600))
continue
if wait > 0:
await asyncio.sleep(wait)
# ── Execute all events at this time slot (per session) ─────────
batch = [ev for ev in pending if ev["et_dt"] == next_et]
for ev in batch:
session = ev.get("session", "")
self._log(f"▶ [{session}] {ev['label']}")
try:
await self._run_trading(ev["kind"], [session] if session else resolved, date_str)
except Exception as exc:
self._log(f" ERROR [{session}]: {exc}")
self._completed.add(ev["name"])
except asyncio.CancelledError:
self._log("ORB auto-scheduler stopped.")
raise
# ── Module-level singleton ────────────────────────────────────────────────────
orb_auto_scheduler = ORBAutoScheduler()