87 lines
2.9 KiB
Python
87 lines
2.9 KiB
Python
import datetime as dt
|
|
import logging
|
|
|
|
from apscheduler.schedulers.background import BackgroundScheduler
|
|
from apscheduler.triggers.date import DateTrigger
|
|
|
|
from .conflict import poll_conflict_events
|
|
from .config import settings
|
|
from .db import SessionLocal
|
|
from .flights import poll_flights
|
|
from .ingest import fetch_all
|
|
from .markets import backfill_market_history, poll_markets
|
|
from .polymarket import poll_polymarket
|
|
|
|
log = logging.getLogger("newsatlas.scheduler")
|
|
|
|
|
|
def _run_rss_job() -> None:
|
|
session = SessionLocal()
|
|
try:
|
|
added = fetch_all(session)
|
|
log.info("RSS poll complete: %d new articles", added)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
def _run_market_job() -> None:
|
|
session = SessionLocal()
|
|
try:
|
|
added = poll_markets(session)
|
|
log.info("Market poll complete: %d instruments recorded", added)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
def _run_conflict_job() -> None:
|
|
session = SessionLocal()
|
|
try:
|
|
added = poll_conflict_events(session)
|
|
log.info("Conflict-event poll complete: %d new events", added)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
def _run_flights_job() -> None:
|
|
count = poll_flights()
|
|
log.info("Flight poll complete: %d aircraft", count)
|
|
|
|
|
|
def _run_market_backfill_job() -> None:
|
|
session = SessionLocal()
|
|
try:
|
|
added = backfill_market_history(session)
|
|
log.info("Market history backfill complete: %d historical points", added)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
def _run_polymarket_job() -> None:
|
|
count = poll_polymarket()
|
|
log.info("Polymarket poll complete: %d trending markets", count)
|
|
|
|
|
|
def start_scheduler() -> BackgroundScheduler:
|
|
scheduler = BackgroundScheduler(timezone="UTC")
|
|
scheduler.add_job(_run_rss_job, "interval", minutes=settings.rss_poll_minutes, next_run_time=None)
|
|
scheduler.add_job(_run_market_job, "interval", minutes=settings.market_poll_minutes, next_run_time=None)
|
|
scheduler.add_job(_run_conflict_job, "interval", minutes=settings.conflict_poll_minutes, next_run_time=None)
|
|
scheduler.add_job(_run_flights_job, "interval", seconds=settings.flights_poll_seconds, next_run_time=None)
|
|
scheduler.add_job(_run_polymarket_job, "interval", minutes=settings.polymarket_poll_minutes, next_run_time=None)
|
|
scheduler.start()
|
|
|
|
# Kick off an immediate first run of each recurring job in the
|
|
# background so the globe isn't empty while waiting for the first
|
|
# interval to elapse.
|
|
now = dt.datetime.utcnow()
|
|
for job in scheduler.get_jobs():
|
|
job.modify(next_run_time=now)
|
|
|
|
# One-time (non-recurring) job: backfills a week of market history so
|
|
# the Economic Incident History view isn't empty on a fresh deployment.
|
|
# Separate from the loop above — it must run exactly once, not on
|
|
# market_poll_minutes' schedule.
|
|
scheduler.add_job(_run_market_backfill_job, trigger=DateTrigger(run_date=now), id="market_backfill")
|
|
|
|
return scheduler
|