""" Bootstrap SC 13D/G activist ownership events for the last N quarters. Skips Form 4 (already bootstrapped separately). Usage (inside the container): python scripts/bootstrap_13dg.py [--quarters N] Idempotent: already-indexed accessions are skipped via upsert ON CONFLICT DO NOTHING. """ import asyncio import fcntl import logging import os import sys from datetime import datetime, timezone sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") logger = logging.getLogger("bootstrap_13dg") _LOCK_FILE = "/tmp/bootstrap_13dg.lock" def _quarters_to_process(n: int): now = datetime.now(timezone.utc) current_q = (now.month - 1) // 3 + 1 results = [] year, q = now.year, current_q for _ in range(n): results.append((year, q)) q -= 1 if q == 0: q = 4 year -= 1 return results async def bootstrap(n_quarters: int = 8) -> None: from app.core.database import AsyncSessionLocal from app.services.sec_full_index_service import SECFullIndexService from app.services.activist_ownership_service import ActivistOwnershipService quarters = _quarters_to_process(n_quarters) logger.info(f"Bootstrapping 13D/G for {n_quarters} quarters: {quarters}") index_svc = SECFullIndexService() activist_svc = ActivistOwnershipService() _BATCH = 5000 # entries per sub-batch within a quarter (prevents OOM on 20k+ quarters) total_inserted = 0 for year, quarter in quarters: logger.info(f"Processing {year}/Q{quarter} ...") try: activist_entries = await index_svc.fetch_quarterly_activist_entries(year, quarter) logger.info(f" 13D/G entries from company.idx: {len(activist_entries)}") quarter_inserted = 0 for batch_start in range(0, len(activist_entries), _BATCH): batch = activist_entries[batch_start:batch_start + _BATCH] async with AsyncSessionLocal() as db: inserted = await activist_svc.ingest_from_index_entries(db, batch) quarter_inserted += inserted # Clear HTTP cache between batches to prevent OOM on large quarters activist_svc._http._text_cache.clear() activist_svc._http._json_cache.clear() import gc; gc.collect() logger.info(f" batch {batch_start// _BATCH + 1}: {inserted} inserted") total_inserted += quarter_inserted logger.info(f" 13D/G index-only inserted: {quarter_inserted}") except Exception as e: logger.warning(f" {year}/Q{quarter} failed: {e}") finally: activist_svc._http._text_cache.clear() activist_svc._http._json_cache.clear() import gc; gc.collect() logger.info(f"Bootstrap complete: {total_inserted} total 13D/G events inserted across {n_quarters} quarters.") if __name__ == "__main__": import argparse parser = argparse.ArgumentParser() parser.add_argument("--quarters", type=int, default=8) args = parser.parse_args() lock_fh = open(_LOCK_FILE, "w") try: fcntl.flock(lock_fh, fcntl.LOCK_EX | fcntl.LOCK_NB) except OSError: logger.error("Bootstrap is already running (lock file held). Exiting.") sys.exit(1) try: asyncio.run(bootstrap(n_quarters=args.quarters)) finally: fcntl.flock(lock_fh, fcntl.LOCK_UN) lock_fh.close() try: os.unlink(_LOCK_FILE) except OSError: pass