feat: wavelets
This commit is contained in:
@@ -180,6 +180,10 @@ def startup():
|
||||
# Start Saxo OAuth token refresh + options-chain snapshot poller
|
||||
from services.saxo_scheduler import start_saxo_scheduler
|
||||
start_saxo_scheduler()
|
||||
# Start Saxo-priced Watchlist + wavelet recompute refresh (own cadence, independent of
|
||||
# the once-a-day auto_cycle — see services/wavelet_scheduler.py)
|
||||
from services.wavelet_scheduler import start_wavelet_scheduler
|
||||
start_wavelet_scheduler()
|
||||
|
||||
# One-time cleanup: collapse snapshot rows stored before save-time dedup existed
|
||||
try:
|
||||
@@ -270,6 +274,8 @@ def shutdown():
|
||||
stop_institutional_scheduler()
|
||||
from services.saxo_scheduler import stop_saxo_scheduler
|
||||
stop_saxo_scheduler()
|
||||
from services.wavelet_scheduler import stop_wavelet_scheduler
|
||||
stop_wavelet_scheduler()
|
||||
|
||||
|
||||
app.include_router(market_data.router)
|
||||
|
||||
@@ -211,6 +211,35 @@ def wavelet_reliability_endpoint(
|
||||
return result
|
||||
|
||||
|
||||
# ── Watchlist refresh scheduler — Saxo-priced quotes + wavelet recompute, own cadence ──
|
||||
# ── (services/wavelet_scheduler.py), independent of the once-a-day auto_cycle ─────────
|
||||
|
||||
class RefreshSettingsRequest(BaseModel):
|
||||
enabled: bool
|
||||
refresh_minutes: float
|
||||
|
||||
|
||||
@router.get("/refresh-settings")
|
||||
def get_refresh_settings():
|
||||
from services.wavelet_scheduler import get_settings
|
||||
return get_settings()
|
||||
|
||||
|
||||
@router.put("/refresh-settings")
|
||||
def update_refresh_settings(req: RefreshSettingsRequest):
|
||||
from services.wavelet_scheduler import set_settings, get_settings
|
||||
set_settings(req.enabled, req.refresh_minutes)
|
||||
return get_settings()
|
||||
|
||||
|
||||
@router.post("/refresh-now")
|
||||
def refresh_now():
|
||||
"""Manual immediate refresh of the whole Watchlist (Saxo-priced quotes + wavelet
|
||||
recompute) — doesn't wait for the periodic poll."""
|
||||
from services.wavelet_scheduler import run_refresh_pass
|
||||
return {"signal_rows": run_refresh_pass()}
|
||||
|
||||
|
||||
# ── Saved simulation/optimization runs ────────────────────────────────────────
|
||||
|
||||
class SimulationCreate(BaseModel):
|
||||
|
||||
76
backend/services/wavelet_scheduler.py
Normal file
76
backend/services/wavelet_scheduler.py
Normal file
@@ -0,0 +1,76 @@
|
||||
"""
|
||||
Periodic Saxo-priced Watchlist refresh + wavelet recompute — mirrors services/saxo_scheduler.py's
|
||||
pattern (own thread, own config-driven interval, `while not stop.wait(0)` so the first pass runs
|
||||
immediately on startup). Independent of services/auto_cycle.py's once-a-day cycle, which also runs
|
||||
this same computation as one of its steps but only at cycle cadence — this lets the Dashboard's
|
||||
Wavelets Signal card and Instrument Analysis's cached Wavelet tab (services.wavelet_signals.
|
||||
scan_watchlist_wavelet_signals writes both) stay current without waiting for, or manually
|
||||
triggering, a full cycle.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
import uuid
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_thread: threading.Thread | None = None
|
||||
_stop = threading.Event()
|
||||
|
||||
DEFAULT_REFRESH_MINUTES = 15
|
||||
|
||||
|
||||
def get_settings() -> dict:
|
||||
from .database import get_config
|
||||
enabled = (get_config("wavelet_refresh_enabled") or "true").lower() == "true"
|
||||
try:
|
||||
minutes = float(get_config("wavelet_refresh_minutes") or str(DEFAULT_REFRESH_MINUTES))
|
||||
except (TypeError, ValueError):
|
||||
minutes = DEFAULT_REFRESH_MINUTES
|
||||
return {"enabled": enabled, "refresh_minutes": minutes}
|
||||
|
||||
|
||||
def set_settings(enabled: bool, refresh_minutes: float) -> None:
|
||||
from .database import set_config
|
||||
set_config("wavelet_refresh_enabled", "true" if enabled else "false")
|
||||
set_config("wavelet_refresh_minutes", str(max(1.0, refresh_minutes)))
|
||||
|
||||
|
||||
def run_refresh_pass() -> int:
|
||||
"""One pass over the whole Watchlist: Saxo-first price fetch + wavelet recompute (see
|
||||
services.wavelet_signals.scan_watchlist_wavelet_signals) — the same work the daily cycle's
|
||||
own wavelet step does, just callable on its own cadence. Shared by the periodic loop and
|
||||
the manual 'refresh now' button. Returns the number of signal rows written."""
|
||||
from .wavelet_signals import compute_and_save_wavelet_signals
|
||||
run_id = f"refresh-{uuid.uuid4().hex[:10]}"
|
||||
try:
|
||||
results = compute_and_save_wavelet_signals(run_id)
|
||||
logger.info(f"[Wavelet Scheduler] Refresh pass complete: {len(results)} signal rows ({run_id})")
|
||||
return len(results)
|
||||
except Exception as e:
|
||||
logger.warning(f"[Wavelet Scheduler] Refresh pass failed: {e}")
|
||||
return 0
|
||||
|
||||
|
||||
def _loop(stop: threading.Event):
|
||||
while not stop.wait(0):
|
||||
settings = get_settings()
|
||||
if not settings["enabled"]:
|
||||
stop.wait(timeout=300) # re-check periodically in case it gets enabled without a restart
|
||||
continue
|
||||
run_refresh_pass()
|
||||
stop.wait(timeout=settings["refresh_minutes"] * 60)
|
||||
|
||||
|
||||
def start_wavelet_scheduler():
|
||||
global _thread
|
||||
_stop.clear()
|
||||
if not (_thread and _thread.is_alive()):
|
||||
_thread = threading.Thread(target=_loop, args=(_stop,), name="wavelet-refresh", daemon=True)
|
||||
_thread.start()
|
||||
logger.info("[Wavelet Scheduler] Started")
|
||||
|
||||
|
||||
def stop_wavelet_scheduler():
|
||||
_stop.set()
|
||||
Reference in New Issue
Block a user