""" Background scheduler for the Saxo connection — mirrors services/var_scheduler.py. Two independent loops: - Token refresh: proactively renews the access/refresh token well before the 40-min refresh-token expiry (SIM), so the user never has to re-login as long as the server keeps running. - Snapshot poller: periodically captures an options-chain snapshot for each symbol in the configured watchlist and appends it to saxo_option_snapshots. """ from __future__ import annotations import json import logging import threading logger = logging.getLogger(__name__) _refresh_thread: threading.Thread | None = None _snapshot_thread: threading.Thread | None = None _refresh_stop = threading.Event() _snapshot_stop = threading.Event() DEFAULT_WATCHLIST = ["SPY", "QQQ", "GLD", "SLV", "USO"] REFRESH_INTERVAL_SECONDS = 10 * 60 # well under the 40-min refresh-token lifetime def get_watchlist() -> list[str]: from .database import get_config raw = get_config("saxo_watchlist") if not raw: return DEFAULT_WATCHLIST try: return json.loads(raw) except Exception: return DEFAULT_WATCHLIST def set_watchlist(symbols: list[str]): from .database import set_config set_config("saxo_watchlist", json.dumps([s.upper() for s in symbols])) def _refresh_loop(stop: threading.Event): while not stop.wait(0): try: from .saxo_auth import get_status, refresh_tokens, disconnect if get_status().get("connected"): try: refresh_tokens() logger.info("[Saxo Scheduler] Token refreshed proactively") except Exception as e: # Refresh token rejected (expired/rotated/invalid) — the stored access # token is now dead even though its expires_at hasn't passed yet. Clear # it so /api/saxo/status stops reporting a stale "connected" state, and # the user gets a clear "reconnect" signal instead of silent 401s later. logger.error(f"[Saxo Scheduler] Refresh failed, clearing dead session: {e}") disconnect() except Exception as e: logger.error(f"[Saxo Scheduler] Refresh loop error: {e}") stop.wait(timeout=REFRESH_INTERVAL_SECONDS) def get_snapshot_minutes() -> float: from .database import get_config return float(get_config("saxo_snapshot_minutes") or "5") def set_snapshot_minutes(minutes: float): from .database import set_config set_config("saxo_snapshot_minutes", str(minutes)) def run_snapshot_pass() -> dict[str, int]: """Snapshot every watchlist symbol once — shared by the periodic loop, the 'connect now' trigger (saxo_oauth.py callback), and the manual refresh button.""" from .saxo_client import snapshot_options_chain from .database import save_saxo_snapshot_rows results: dict[str, int] = {} for symbol in get_watchlist(): try: rows = snapshot_options_chain(symbol) save_saxo_snapshot_rows(rows) results[symbol] = len(rows) logger.info(f"[Saxo Scheduler] Snapshot saved: {symbol} ({len(rows)} rows)") except Exception as e: results[symbol] = 0 logger.warning(f"[Saxo Scheduler] Snapshot failed for {symbol}: {e}") return results def _snapshot_loop(stop: threading.Event): while not stop.wait(0): try: from .saxo_auth import get_status minutes = get_snapshot_minutes() if get_status().get("connected"): run_snapshot_pass() else: minutes = max(minutes, 5) # don't busy-loop while disconnected except Exception as e: logger.error(f"[Saxo Scheduler] Snapshot loop error: {e}") minutes = 5 stop.wait(timeout=minutes * 60) def start_saxo_scheduler(): global _refresh_thread, _snapshot_thread _refresh_stop.clear() _snapshot_stop.clear() if not (_refresh_thread and _refresh_thread.is_alive()): _refresh_thread = threading.Thread(target=_refresh_loop, args=(_refresh_stop,), name="saxo-refresh", daemon=True) _refresh_thread.start() if not (_snapshot_thread and _snapshot_thread.is_alive()): _snapshot_thread = threading.Thread(target=_snapshot_loop, args=(_snapshot_stop,), name="saxo-snapshot", daemon=True) _snapshot_thread.start() logger.info("[Saxo Scheduler] Started") def stop_saxo_scheduler(): _refresh_stop.set() _snapshot_stop.set()