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.

202 lines
7.4 KiB
Python

"""
APScheduler-based daily/weekly SEC ingest scheduler.
Jobs:
form4_daily_ingest — Tue-Sat 09:00 ET: previous business day Form 4
activist_13dg_daily_ingest — same schedule: SC 13D/G index-only ingest
form4_weekly_reindex — Sat 03:00 ET: current quarter company.idx rescan
activist_13dg_enrich — every 30 min: enrich index_only rows with ownership_pct
Wire up via start_sec_ingest_scheduler() / stop_sec_ingest_scheduler() in app lifespan.
"""
import logging
from datetime import date, datetime, timedelta, timezone
logger = logging.getLogger(__name__)
_scheduler = None
def _get_scheduler():
global _scheduler
if _scheduler is None:
try:
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from app.core.config import settings
_scheduler = AsyncIOScheduler(timezone=settings.SEC_INGEST_TIMEZONE)
except ImportError:
logger.warning("apscheduler not installed; SEC ingest scheduling disabled")
return None
return _scheduler
def _prev_business_day() -> date:
"""Return the most recent business day (Mon-Fri) relative to today."""
today = datetime.now(timezone.utc).date()
delta = timedelta(days=1)
# If today is Mon, go back to Fri; otherwise go back 1 day
candidate = today - delta
while candidate.weekday() >= 5: # 5=Sat, 6=Sun
candidate -= delta
return candidate
async def _run_form4_daily_ingest() -> None:
"""Ingest prior business day Form 4 entries from SEC daily full-index."""
from app.core.database import AsyncSessionLocal
from app.services.sec_full_index_service import SECFullIndexService
from app.services.insider_transaction_service import InsiderTransactionService
target_date = _prev_business_day()
date_str = target_date.strftime("%Y%m%d")
logger.info(f"[SEC Ingest] Form 4 daily ingest for {date_str}")
index_svc = SECFullIndexService()
txn_svc = InsiderTransactionService()
entries = await index_svc.fetch_daily_form4_entries(date_str)
if not entries:
logger.info(f"[SEC Ingest] No Form 4 entries in daily index for {date_str}")
return
async with AsyncSessionLocal() as db:
try:
inserted = await txn_svc.index_form4_from_index_entries(db, entries)
logger.info(f"[SEC Ingest] Form 4 daily: {inserted} new transactions for {date_str}")
except Exception as e:
logger.error(f"[SEC Ingest] Form 4 daily ingest failed ({date_str}): {e}")
async def _run_activist_daily_ingest() -> None:
"""Ingest prior business day SC 13D/G entries from SEC daily full-index."""
from app.core.database import AsyncSessionLocal
from app.services.sec_full_index_service import SECFullIndexService
from app.services.activist_ownership_service import ActivistOwnershipService
target_date = _prev_business_day()
date_str = target_date.strftime("%Y%m%d")
logger.info(f"[SEC Ingest] Activist 13D/G daily ingest for {date_str}")
index_svc = SECFullIndexService()
activist_svc = ActivistOwnershipService()
entries = await index_svc.fetch_daily_activist_entries(date_str)
if not entries:
logger.info(f"[SEC Ingest] No 13D/G entries in daily index for {date_str}")
return
async with AsyncSessionLocal() as db:
try:
inserted = await activist_svc.ingest_from_index_entries(db, entries)
logger.info(f"[SEC Ingest] Activist 13D/G daily: {inserted} new rows for {date_str}")
except Exception as e:
logger.error(f"[SEC Ingest] Activist daily ingest failed ({date_str}): {e}")
async def _run_form4_weekly_reindex() -> None:
"""Re-scan current quarter's company.idx to catch corrections and amendments."""
from app.core.database import AsyncSessionLocal
from app.services.sec_full_index_service import SECFullIndexService
from app.services.insider_transaction_service import InsiderTransactionService
today = datetime.now(timezone.utc).date()
quarter = (today.month - 1) // 3 + 1
logger.info(f"[SEC Ingest] Form 4 weekly reindex {today.year}/Q{quarter}")
index_svc = SECFullIndexService()
txn_svc = InsiderTransactionService()
entries = await index_svc.fetch_quarterly_form4_entries(today.year, quarter)
if not entries:
return
async with AsyncSessionLocal() as db:
try:
inserted = await txn_svc.index_form4_from_index_entries(db, entries)
logger.info(f"[SEC Ingest] Form 4 weekly reindex: {inserted} new transactions")
except Exception as e:
logger.error(f"[SEC Ingest] Form 4 weekly reindex failed: {e}")
async def _run_activist_enrich() -> None:
"""Enrich a batch of index_only activist rows with ownership_pct / shares_owned."""
from app.core.database import AsyncSessionLocal
from app.services.activist_ownership_service import ActivistOwnershipService
activist_svc = ActivistOwnershipService()
async with AsyncSessionLocal() as db:
try:
enriched = await activist_svc.enrich_pending(db)
if enriched:
logger.info(f"[SEC Ingest] Activist enrich: {enriched} rows enriched")
except Exception as e:
logger.error(f"[SEC Ingest] Activist enrich job failed: {e}")
def start_sec_ingest_scheduler() -> None:
"""Start the SEC ingest scheduler. Call from FastAPI lifespan startup."""
sched = _get_scheduler()
if sched is None:
return
try:
from apscheduler.triggers.cron import CronTrigger
from apscheduler.triggers.interval import IntervalTrigger
sched.add_job(
_run_form4_daily_ingest,
trigger=CronTrigger(day_of_week="tue-sat", hour=9, minute=0),
id="form4_daily_ingest",
replace_existing=True,
max_instances=1,
misfire_grace_time=600,
coalesce=True,
)
sched.add_job(
_run_activist_daily_ingest,
trigger=CronTrigger(day_of_week="tue-sat", hour=9, minute=5),
id="activist_13dg_daily_ingest",
replace_existing=True,
max_instances=1,
misfire_grace_time=600,
coalesce=True,
)
sched.add_job(
_run_form4_weekly_reindex,
trigger=CronTrigger(day_of_week="sat", hour=3, minute=0),
id="form4_weekly_reindex",
replace_existing=True,
max_instances=1,
misfire_grace_time=600,
coalesce=True,
)
sched.add_job(
_run_activist_enrich,
trigger=IntervalTrigger(minutes=30),
id="activist_13dg_enrich",
replace_existing=True,
max_instances=1,
misfire_grace_time=120,
coalesce=True,
)
sched.start()
logger.info(
"SEC ingest scheduler started: "
"Form 4 daily @ Tue-Sat 09:00 ET, "
"13D/G daily @ Tue-Sat 09:05 ET, "
"Form 4 weekly reindex @ Sat 03:00 ET, "
"Activist enrich every 30 min"
)
except Exception as e:
logger.error(f"SEC ingest scheduler start failed: {e}")
def stop_sec_ingest_scheduler() -> None:
"""Stop the scheduler. Call from FastAPI lifespan shutdown."""
sched = _get_scheduler()
if sched and sched.running:
sched.shutdown(wait=False)
logger.info("SEC ingest scheduler stopped")