Server no longer holds portfolios. Holdings live in the browser (localStorage); the server publishes an anonymous ticker_universe and a gzipped /api/universe payload identical for every authenticated user, so access patterns can't betray which tickers a user holds. AI commentary is generated ephemerally from the browser-supplied pie and the cost ledger row records no positions. Migrations 0009-0011 added the universe table and dropped positions / portfolio_snapshots / portfolios. Authentication is now e-mail OTP only. Migration 0010 dropped password_hash and email_verified (every active session is by construction proof of email control). The /signup endpoint is gone; signup and login share a single email-entry page. Email rendering is HTML+plain-text multipart with a shared brand palette (app/branding.py) asserted in sync with the CSS by a drift-detection test. LLM provider defaults to DeepSeek-direct (cheaper, api.deepseek.com) with OpenRouter as automatic fallback if DeepSeek fails. ai_log_job and indicator_summary_job now iterate the two tones (NOVICE, INTERMEDIATE) per cycle so the dashboard's tone toggle is instant; PROMPT_VERSION bumped to 6 with an educational anti-TA / anti-gambling stance baked into _CORE. NOVICE mode renders a curated glossary inline (CBOE VIX, yield curve, HY OAS, etc.) with JS-positioned tooltips that survive viewport edges and sticky bars. Model name and tokens hidden from the user UI; still recorded in StrategicLog.model and AICall for admin. Layout adds a sticky top nav, a sticky bottom markets bar (one chip per exchange with status LED + headline index + 1d change), and Phase H feedback reporting is queued in tasks/todo.md. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
77 lines
2.8 KiB
Python
77 lines
2.8 KiB
Python
"""Scheduler container entrypoint. Runs APScheduler with 5 cron jobs, each
|
|
guarded by a MariaDB advisory lock (in job_lifecycle). Waits for the DB to be
|
|
reachable, then schedules and blocks forever."""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import signal
|
|
|
|
from apscheduler.schedulers.asyncio import AsyncIOScheduler
|
|
from apscheduler.triggers.cron import CronTrigger
|
|
|
|
from app.db import get_engine
|
|
from app.logging import configure_logging, get_logger
|
|
from app.jobs import (
|
|
market_job, news_job, ai_log_job, rollup_job,
|
|
indicator_summary_job, universe_flush_job,
|
|
)
|
|
|
|
|
|
log = get_logger("scheduler")
|
|
|
|
|
|
async def _wait_for_db(retries: int = 60, delay: float = 1.0) -> None:
|
|
engine = get_engine()
|
|
for i in range(retries):
|
|
try:
|
|
async with engine.connect() as conn:
|
|
await conn.execute(__import__("sqlalchemy").text("SELECT 1"))
|
|
return
|
|
except Exception as e:
|
|
log.warning("scheduler.db_wait", attempt=i + 1, error=str(e)[:120])
|
|
await asyncio.sleep(delay)
|
|
raise RuntimeError("DB never became reachable")
|
|
|
|
|
|
async def main() -> None:
|
|
configure_logging()
|
|
log.info("scheduler.starting")
|
|
await _wait_for_db()
|
|
|
|
sched = AsyncIOScheduler(timezone="UTC")
|
|
sched.add_job(market_job.run, CronTrigger(minute=5), name="market_job", id="market_job")
|
|
sched.add_job(news_job.run, CronTrigger(minute=10), name="news_job", id="news_job")
|
|
# portfolio_job removed in Phase G — server no longer holds holdings.
|
|
sched.add_job(indicator_summary_job.run, CronTrigger(minute=7), name="indicator_summary_job", id="indicator_summary_job")
|
|
sched.add_job(ai_log_job.run, CronTrigger(minute=20), name="ai_log_job", id="ai_log_job")
|
|
sched.add_job(rollup_job.run, CronTrigger(hour=0, minute=5), name="rollup_job", id="rollup_job")
|
|
# Phase G: flush the Redis ticker-add buffer every 5 minutes (xx:01,
|
|
# xx:06, ...). The 1-min offset gives the bucket boundary time to
|
|
# close before we read the previous one.
|
|
sched.add_job(universe_flush_job.run,
|
|
CronTrigger(minute="1-59/5"),
|
|
name="universe_flush_job", id="universe_flush_job")
|
|
sched.add_job(universe_flush_job.evict_run,
|
|
CronTrigger(hour=0, minute=15),
|
|
name="universe_evict_job", id="universe_evict_job")
|
|
sched.start()
|
|
log.info("scheduler.started", jobs=[j.id for j in sched.get_jobs()])
|
|
|
|
# Stay alive until SIGTERM.
|
|
stop_event = asyncio.Event()
|
|
|
|
def _stop(*_):
|
|
log.info("scheduler.stopping")
|
|
stop_event.set()
|
|
|
|
loop = asyncio.get_running_loop()
|
|
for sig in (signal.SIGTERM, signal.SIGINT):
|
|
loop.add_signal_handler(sig, _stop)
|
|
|
|
await stop_event.wait()
|
|
sched.shutdown(wait=False)
|
|
log.info("scheduler.stopped")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main())
|