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.
201 lines
7.3 KiB
Python
201 lines
7.3 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.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
|
|
|
|
# index_form4_from_index_entries opens its own short-lived sessions for
|
|
# dedup and each flush, so we don't hold a connection across the long
|
|
# HTTP loop (which previously dropped asyncpg mid-job).
|
|
try:
|
|
inserted = await txn_svc.index_form4_from_index_entries(None, 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.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
|
|
|
|
try:
|
|
inserted = await txn_svc.index_form4_from_index_entries(None, 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")
|