diff --git a/app/api/v1/endpoints/insider.py b/app/api/v1/endpoints/insider.py index c1a72e2..37bea7d 100644 --- a/app/api/v1/endpoints/insider.py +++ b/app/api/v1/endpoints/insider.py @@ -8,9 +8,11 @@ from typing import Optional from fastapi import APIRouter, Depends, HTTPException, Query from fastapi.responses import Response +from sqlalchemy import func, select from sqlalchemy.ext.asyncio import AsyncSession from app.core.database import get_db +from app.models.insider_transaction import InsiderTransaction from app.schemas.insider import ( Form4AggregateResponse, Form4ByDateResponse, @@ -126,11 +128,12 @@ async def get_insider_summary( description=( "Returns Form 4 transactions for a ticker where **filing_date ≤ as_of** (point-in-time safe).\n\n" "`as_of` is required to prevent lookahead in backtests.\n\n" - "`start`/`end` also filter by `filing_date` (not transaction_date)." + "`start`/`end` also filter by `filing_date` (not transaction_date).\n\n" + "If no data exists for the ticker, auto-fetches ~2 years of history from SEC EDGAR (first call may take 1–3 min)." ), ) @with_cache(namespace="insider:form4", ttl=300, - key_params=["ticker", "start", "end", "as_of", "buy_only", "csuite_only"]) + key_params=["ticker", "start", "end", "as_of", "buy_only", "csuite_only", "force_refresh"]) async def get_form4( ticker: str, response: Response, @@ -139,10 +142,28 @@ async def get_form4( end: Optional[date] = Query(None, description="Window end (filing_date ≤ end)"), buy_only: bool = Query(False, description="Only return open-market purchases (transaction_code=P, shares > 0). Excludes awards/grants."), csuite_only: bool = Query(False, description="Only return C-suite insider transactions"), + force_refresh: bool = Query(False, description="Re-fetch ~2yr of Form 4 history from SEC EDGAR before querying. Slow on first call."), db: AsyncSession = Depends(get_db), ): svc = InsiderTransactionService() try: + auto_fetched = False + fetched_count = 0 + + # Auto-fetch when ticker has no data at all or force_refresh requested + if force_refresh: + fetched_count = await svc.index_form4s(db, ticker.upper(), days=800, force_refresh=True) + auto_fetched = True + else: + count_q = await db.execute( + select(func.count(InsiderTransaction.id)).where( + InsiderTransaction.ticker == ticker.upper() + ) + ) + if (count_q.scalar() or 0) == 0: + fetched_count = await svc.index_form4s(db, ticker.upper(), days=800, force_refresh=False) + auto_fetched = True + rows, total = await svc.get_form4_pit( db, ticker=ticker, as_of=as_of, start=start, end=end, buy_only=buy_only, csuite_only=csuite_only, @@ -154,6 +175,7 @@ async def get_form4( window={"start": start.isoformat() if start else None, "end": end.isoformat() if end else None}, transactions=entries, total_count=total, + metadata={"auto_fetched": auto_fetched, "fetched_count": fetched_count}, ) except Exception as e: logger.error(f"Form4 error for {ticker}: {e}") @@ -195,21 +217,39 @@ async def get_form4_by_date( description=( "Aggregated insider buy metrics within [as_of - window_days, as_of].\n\n" "All based on `filing_date` (PIT-safe). Returns buy_count, buy_dollar_total, " - "cluster_size (unique insiders), csuite_count, avg_pct_of_holding, recency_days." + "cluster_size (unique insiders), csuite_count, avg_pct_of_holding, recency_days.\n\n" + "If no data exists for the ticker, auto-fetches ~2 years of history from SEC EDGAR." ), ) -@with_cache(namespace="insider:form4_agg", ttl=300, key_params=["ticker", "as_of", "window_days"]) +@with_cache(namespace="insider:form4_agg", ttl=300, key_params=["ticker", "as_of", "window_days", "force_refresh"]) async def get_form4_aggregate( ticker: str, response: Response, as_of: date = Query(..., description="Point-in-time cutoff. Required."), window_days: int = Query(30, ge=1, le=365, description="Lookback window in days"), + force_refresh: bool = Query(False, description="Re-fetch ~2yr of Form 4 history from SEC EDGAR before querying."), db: AsyncSession = Depends(get_db), ): svc = InsiderTransactionService() try: + auto_fetched = False + fetched_count = 0 + + if force_refresh: + fetched_count = await svc.index_form4s(db, ticker.upper(), days=800, force_refresh=True) + auto_fetched = True + else: + count_q = await db.execute( + select(func.count(InsiderTransaction.id)).where( + InsiderTransaction.ticker == ticker.upper() + ) + ) + if (count_q.scalar() or 0) == 0: + fetched_count = await svc.index_form4s(db, ticker.upper(), days=800, force_refresh=False) + auto_fetched = True + agg = await svc.get_form4_aggregate(db, ticker=ticker, as_of=as_of, window_days=window_days) - return Form4AggregateResponse(**agg) + return Form4AggregateResponse(**agg, metadata={"auto_fetched": auto_fetched, "fetched_count": fetched_count}) except Exception as e: logger.error(f"Form4 aggregate error for {ticker}: {e}") raise HTTPException(status_code=502, detail=str(e)) diff --git a/app/schemas/insider.py b/app/schemas/insider.py index 958e503..e7c6f42 100644 --- a/app/schemas/insider.py +++ b/app/schemas/insider.py @@ -151,6 +151,7 @@ class Form4Response(BaseModel): window: Dict[str, Any] = Field(default_factory=dict) transactions: List[Form4Entry] total_count: int + metadata: Dict[str, Any] = Field(default_factory=dict) class Form4ByDateResponse(BaseModel): @@ -175,3 +176,4 @@ class Form4AggregateResponse(BaseModel): csuite_count: int avg_pct_of_holding: Optional[float] = None recency_days: int + metadata: Dict[str, Any] = Field(default_factory=dict) diff --git a/scripts/audit_form4_coverage.py b/scripts/audit_form4_coverage.py new file mode 100644 index 0000000..b5d9743 --- /dev/null +++ b/scripts/audit_form4_coverage.py @@ -0,0 +1,70 @@ +""" +Audit Form 4 coverage for all universe_ticker_registry tickers. + +Prints per-ticker earliest filing_date and P-code buy count. +Flags tickers with no data or earliest date after the baseline. + +Usage (inside container): + python scripts/audit_form4_coverage.py [--baseline YYYY-MM-DD] +""" + +import asyncio +import os +import sys +from datetime import date + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + + +async def audit(baseline: date) -> None: + from app.core.database import AsyncSessionLocal + from sqlalchemy import text + + async with AsyncSessionLocal() as db: + rows = await db.execute(text(""" + WITH u AS (SELECT ticker FROM universe_ticker_registry ORDER BY ticker), + e AS ( + SELECT ticker, + MIN(filing_date)::date AS earliest, + COUNT(*) AS total_rows, + SUM(CASE WHEN transaction_code = 'P' AND shares > 0 THEN 1 ELSE 0 END) AS p_buys + FROM insider_transactions + GROUP BY ticker + ) + SELECT u.ticker, e.earliest, e.total_rows, e.p_buys + FROM u LEFT JOIN e USING (ticker) + ORDER BY e.earliest NULLS FIRST, u.ticker + """)) + results = rows.fetchall() + + no_data, late, ok = [], [], [] + for ticker, earliest, total_rows, p_buys in results: + if earliest is None: + no_data.append(ticker) + elif earliest > baseline: + late.append((ticker, earliest, total_rows, p_buys)) + else: + ok.append((ticker, earliest, total_rows, p_buys)) + + print(f"\n=== Form 4 Coverage Audit (baseline: {baseline}) ===") + print(f"Universe: {len(results)} tickers | OK: {len(ok)} | Late: {len(late)} | No data: {len(no_data)}\n") + + if no_data: + print(f"NO DATA ({len(no_data)}): {', '.join(no_data)}\n") + + if late: + print(f"LATE (earliest > {baseline}):") + for ticker, earliest, total_rows, p_buys in late: + print(f" {ticker:<8} earliest={earliest} rows={total_rows or 0} P-buys={p_buys or 0}") + print() + + print(f"OK: {len(ok)} tickers meet {baseline} baseline") + + +if __name__ == "__main__": + import argparse + parser = argparse.ArgumentParser() + parser.add_argument("--baseline", default="2024-04-01") + args = parser.parse_args() + baseline = date.fromisoformat(args.baseline) + asyncio.run(audit(baseline)) diff --git a/scripts/bootstrap_form4_by_ticker.py b/scripts/bootstrap_form4_by_ticker.py index 5279545..ab18b8e 100644 --- a/scripts/bootstrap_form4_by_ticker.py +++ b/scripts/bootstrap_form4_by_ticker.py @@ -39,6 +39,7 @@ async def bootstrap( resume_from: Optional[str] = None, days: int = 730, delay_ms: int = 300, + tickers_override: Optional[list] = None, ) -> None: from app.core.database import AsyncSessionLocal from app.services.insider_transaction_service import InsiderTransactionService @@ -46,14 +47,18 @@ async def bootstrap( txn_svc = InsiderTransactionService() - async with AsyncSessionLocal() as db: - rows = await db.execute( - text("SELECT ticker FROM universe_ticker_registry ORDER BY ticker") - ) - tickers = [r[0] for r in rows.fetchall()] + if tickers_override: + tickers = [t.upper() for t in tickers_override] + logger.info(f"Using explicit ticker list: {tickers}") + else: + async with AsyncSessionLocal() as db: + rows = await db.execute( + text("SELECT ticker FROM universe_ticker_registry ORDER BY ticker") + ) + tickers = [r[0] for r in rows.fetchall()] total = len(tickers) - logger.info(f"Universe: {total} tickers, fetching {days} days of Form 4 history each") + logger.info(f"Target: {total} tickers, fetching {days} days of Form 4 history each") if resume_from: resume_from = resume_from.upper() @@ -114,8 +119,15 @@ if __name__ == "__main__": parser.add_argument("--resume-from", default=None) parser.add_argument("--days", type=int, default=730) parser.add_argument("--delay-ms", type=int, default=300) + parser.add_argument( + "--tickers", + default=None, + help="Comma-separated ticker list to process instead of full universe_ticker_registry. E.g. NVDA,MSTR,TSLA", + ) args = parser.parse_args() + tickers_override = [t.strip().upper() for t in args.tickers.split(",") if t.strip()] if args.tickers else None + lock_fh = open(_LOCK_FILE, "w") try: fcntl.flock(lock_fh, fcntl.LOCK_EX | fcntl.LOCK_NB) @@ -128,6 +140,7 @@ if __name__ == "__main__": resume_from=args.resume_from, days=args.days, delay_ms=args.delay_ms, + tickers_override=tickers_override, )) finally: fcntl.flock(lock_fh, fcntl.LOCK_UN)