feat: /form4 auto-fetch fallback + targeted backfill 지원

- /form4/{ticker}, /form4/aggregate/{ticker}: DB에 ticker 데이터 없으면 SEC에서 800일치 자동 fetch (index_form4s)
- force_refresh 파라미터 추가로 on-demand 재인덱싱 지원
- Form4Response, Form4AggregateResponse에 metadata 필드 추가 (auto_fetched, fetched_count)
- bootstrap_form4_by_ticker.py: --tickers 옵션 추가로 특정 ticker만 targeted backfill 가능
- audit_form4_coverage.py: universe 커버리지 검증 스크립트 신규

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
main
I Luk Kim 4 months ago
parent 93a81fd35a
commit f19d31f125

@ -8,9 +8,11 @@ from typing import Optional
from fastapi import APIRouter, Depends, HTTPException, Query from fastapi import APIRouter, Depends, HTTPException, Query
from fastapi.responses import Response from fastapi.responses import Response
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.core.database import get_db from app.core.database import get_db
from app.models.insider_transaction import InsiderTransaction
from app.schemas.insider import ( from app.schemas.insider import (
Form4AggregateResponse, Form4AggregateResponse,
Form4ByDateResponse, Form4ByDateResponse,
@ -126,11 +128,12 @@ async def get_insider_summary(
description=( description=(
"Returns Form 4 transactions for a ticker where **filing_date ≤ as_of** (point-in-time safe).\n\n" "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" "`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 13 min)."
), ),
) )
@with_cache(namespace="insider:form4", ttl=300, @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( async def get_form4(
ticker: str, ticker: str,
response: Response, response: Response,
@ -139,10 +142,28 @@ async def get_form4(
end: Optional[date] = Query(None, description="Window end (filing_date ≤ end)"), 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."), 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"), 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), db: AsyncSession = Depends(get_db),
): ):
svc = InsiderTransactionService() svc = InsiderTransactionService()
try: 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( rows, total = await svc.get_form4_pit(
db, ticker=ticker, as_of=as_of, start=start, end=end, db, ticker=ticker, as_of=as_of, start=start, end=end,
buy_only=buy_only, csuite_only=csuite_only, 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}, window={"start": start.isoformat() if start else None, "end": end.isoformat() if end else None},
transactions=entries, transactions=entries,
total_count=total, total_count=total,
metadata={"auto_fetched": auto_fetched, "fetched_count": fetched_count},
) )
except Exception as e: except Exception as e:
logger.error(f"Form4 error for {ticker}: {e}") logger.error(f"Form4 error for {ticker}: {e}")
@ -195,21 +217,39 @@ async def get_form4_by_date(
description=( description=(
"Aggregated insider buy metrics within [as_of - window_days, as_of].\n\n" "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, " "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( async def get_form4_aggregate(
ticker: str, ticker: str,
response: Response, response: Response,
as_of: date = Query(..., description="Point-in-time cutoff. Required."), as_of: date = Query(..., description="Point-in-time cutoff. Required."),
window_days: int = Query(30, ge=1, le=365, description="Lookback window in days"), 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), db: AsyncSession = Depends(get_db),
): ):
svc = InsiderTransactionService() svc = InsiderTransactionService()
try: 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) 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: except Exception as e:
logger.error(f"Form4 aggregate error for {ticker}: {e}") logger.error(f"Form4 aggregate error for {ticker}: {e}")
raise HTTPException(status_code=502, detail=str(e)) raise HTTPException(status_code=502, detail=str(e))

@ -151,6 +151,7 @@ class Form4Response(BaseModel):
window: Dict[str, Any] = Field(default_factory=dict) window: Dict[str, Any] = Field(default_factory=dict)
transactions: List[Form4Entry] transactions: List[Form4Entry]
total_count: int total_count: int
metadata: Dict[str, Any] = Field(default_factory=dict)
class Form4ByDateResponse(BaseModel): class Form4ByDateResponse(BaseModel):
@ -175,3 +176,4 @@ class Form4AggregateResponse(BaseModel):
csuite_count: int csuite_count: int
avg_pct_of_holding: Optional[float] = None avg_pct_of_holding: Optional[float] = None
recency_days: int recency_days: int
metadata: Dict[str, Any] = Field(default_factory=dict)

@ -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))

@ -39,6 +39,7 @@ async def bootstrap(
resume_from: Optional[str] = None, resume_from: Optional[str] = None,
days: int = 730, days: int = 730,
delay_ms: int = 300, delay_ms: int = 300,
tickers_override: Optional[list] = None,
) -> None: ) -> None:
from app.core.database import AsyncSessionLocal from app.core.database import AsyncSessionLocal
from app.services.insider_transaction_service import InsiderTransactionService from app.services.insider_transaction_service import InsiderTransactionService
@ -46,14 +47,18 @@ async def bootstrap(
txn_svc = InsiderTransactionService() txn_svc = InsiderTransactionService()
async with AsyncSessionLocal() as db: if tickers_override:
rows = await db.execute( tickers = [t.upper() for t in tickers_override]
text("SELECT ticker FROM universe_ticker_registry ORDER BY ticker") logger.info(f"Using explicit ticker list: {tickers}")
) else:
tickers = [r[0] for r in rows.fetchall()] 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) 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: if resume_from:
resume_from = resume_from.upper() resume_from = resume_from.upper()
@ -114,8 +119,15 @@ if __name__ == "__main__":
parser.add_argument("--resume-from", default=None) parser.add_argument("--resume-from", default=None)
parser.add_argument("--days", type=int, default=730) parser.add_argument("--days", type=int, default=730)
parser.add_argument("--delay-ms", type=int, default=300) 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() 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") lock_fh = open(_LOCK_FILE, "w")
try: try:
fcntl.flock(lock_fh, fcntl.LOCK_EX | fcntl.LOCK_NB) fcntl.flock(lock_fh, fcntl.LOCK_EX | fcntl.LOCK_NB)
@ -128,6 +140,7 @@ if __name__ == "__main__":
resume_from=args.resume_from, resume_from=args.resume_from,
days=args.days, days=args.days,
delay_ms=args.delay_ms, delay_ms=args.delay_ms,
tickers_override=tickers_override,
)) ))
finally: finally:
fcntl.flock(lock_fh, fcntl.LOCK_UN) fcntl.flock(lock_fh, fcntl.LOCK_UN)

Loading…
Cancel
Save