chore: decommission FMP, Finnhub and Alpha Vantage (A6)

The A5 cutover has been on and observed in production, so SEC Company Facts +
DoltHub earnings are already the live source for `fundamental_data`. This
removes everything the legacy path still occupied.

Gone: the three providers and their config/env keys; the weekly
`fundamental_collector` job; the cutover toggle (SEC + Dolt is now the
unconditional path, so `off` can no longer silently freeze scoring inputs); the
A5 parity report, whose deltas became structurally zero once the candidate
builder started writing the table it compared against; and the FMP tier of
universe bootstrap.

Two behavioral notes:

- Disabling **SEC Fundamentals Import** now stops the SEC network fetch only.
  The local cache refresh moved outside the job-enable check, because candidates
  also derive from daily closes and earnings events — freezing those on an
  ingestion pause would stale scoring with no fallback left to recover from.
- `/ingestion/fetch?sources=fundamentals` still accepts the key and reports
  `skipped`; there is no per-ticker fetch any more.

Migration 029 does not blanket-delete the leftover settings rows. Migrations run
before the service restart, and pre-A6 code reads an absent `job_*_enabled` row
as *enabled* — so the two behavior-bearing keys become tombstones pinned to safe
values (hidden in Admin) and only the inert three are deleted. Removing the
provider keys from the production `.env` is the matching rollout step.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-07 11:19:28 +02:00
co-authored by Claude Opus 5
parent f5d4b516ab
commit 3e83d63b05
51 changed files with 368 additions and 3787 deletions
-20
View File
@@ -28,15 +28,6 @@ class Settings(BaseSettings):
deepseek_api_key: str = ""
xai_api_key: str = ""
# Fundamentals Provider — Financial Modeling Prep
fmp_api_key: str = ""
# Fundamentals Provider — Finnhub (optional fallback)
finnhub_api_key: str = ""
# Fundamentals Provider — Alpha Vantage (optional fallback)
alpha_vantage_api_key: str = ""
# Dolt bulk-data — local clone of post-no-preference/earnings (workstream A).
# dolt_binary: full path when not on PATH (dev/Windows install). dolt_data_dir
# holds the clones; in production it MUST be outside the deploy tree (deploy is
@@ -61,10 +52,6 @@ class Settings(BaseSettings):
sec_max_retries: int = 4
sec_request_timeout_seconds: float = 30.0
# A5 read-only comparison artifacts. Production must keep this outside the
# rsync deployment tree so the 5-7 day review window survives deploys.
fundamentals_parity_report_dir: str = "reports/fundamentals-parity"
# Regime Monitor — FRED (VIX level + HY credit spreads). Optional: without it
# the volatility (P5) and credit-spread (F2) signals are reported as n/a.
fred_api_key: str = ""
@@ -86,15 +73,8 @@ class Settings(BaseSettings):
# the score window is 7 days).
sentiment_fresh_hours: int = 120
sentiment_top_composite: int = 30
fundamental_fetch_frequency: str = "weekly" # quarterly-ish data; weekly conserves API quota
rr_scan_frequency: str = "daily" # legacy label; qualifying scan is cron near-close
# alerts_frequency removed: alerts fire only via morning + near-close pipelines
fundamental_rate_limit_retries: int = 3
fundamental_rate_limit_backoff_seconds: int = 15
# Pause between tickers in the bulk fundamentals job. Free tiers throttle
# hard (Finnhub ~60 calls/min, ~3 calls/ticker → ~3s/ticker); without
# spacing the job bursts straight into 429s. 0 disables.
fundamental_request_spacing_seconds: float = 3.0
# Scoring Defaults
default_watchlist_auto_size: int = 10
-174
View File
@@ -1,174 +0,0 @@
"""Financial Modeling Prep (FMP) fundamentals provider using httpx.
Uses the stable API endpoints (https://financialmodelingprep.com/stable/)
which replaced the legacy /api/v3/ endpoints deprecated in Aug 2025.
"""
from __future__ import annotations
import logging
import os
from datetime import datetime, timezone
from pathlib import Path
import httpx
from app.exceptions import ProviderError, RateLimitError
from app.providers.protocol import FundamentalData
logger = logging.getLogger(__name__)
_FMP_STABLE_URL = "https://financialmodelingprep.com/stable"
# Resolve CA bundle for explicit httpx verify
_CA_BUNDLE = os.environ.get("SSL_CERT_FILE", "")
if not _CA_BUNDLE or not Path(_CA_BUNDLE).exists():
_CA_BUNDLE_PATH: str | bool = True # use system default
else:
_CA_BUNDLE_PATH = _CA_BUNDLE
class FMPFundamentalProvider:
"""Fetches fundamental data from Financial Modeling Prep REST API."""
def __init__(self, api_key: str) -> None:
if not api_key:
raise ProviderError("FMP API key is required")
self._api_key = api_key
# Mapping from FMP endpoint name to the FundamentalData field it populates
_ENDPOINT_FIELD_MAP: dict[str, str] = {
"ratios-ttm": "pe_ratio",
"financial-growth": "revenue_growth",
"earnings": "earnings_surprise",
}
async def fetch_fundamentals(self, ticker: str) -> FundamentalData:
"""Fetch P/E, revenue growth, earnings surprise, and market cap.
Fetches from multiple stable endpoints. If a supplementary endpoint
(ratios, growth, earnings) returns 402 (paid tier), we gracefully
degrade and return partial data rather than failing entirely, and
record the affected field in ``unavailable_fields``.
"""
try:
endpoints_402: set[str] = set()
async with httpx.AsyncClient(timeout=30.0, verify=_CA_BUNDLE_PATH) as client:
params = {"symbol": ticker, "apikey": self._api_key}
# Profile is the primary source — must succeed
profile = await self._fetch_json(client, "profile", params, ticker)
# Supplementary sources — degrade gracefully on 402
ratios, was_402 = await self._fetch_json_optional(client, "ratios-ttm", params, ticker)
if was_402:
endpoints_402.add("ratios-ttm")
growth, was_402 = await self._fetch_json_optional(client, "financial-growth", params, ticker)
if was_402:
endpoints_402.add("financial-growth")
earnings, was_402 = await self._fetch_json_optional(client, "earnings", params, ticker)
if was_402:
endpoints_402.add("earnings")
pe_ratio = self._safe_float(ratios.get("priceToEarningsRatioTTM"))
revenue_growth = self._safe_float(growth.get("revenueGrowth"))
market_cap = self._safe_float(profile.get("marketCap"))
earnings_surprise = self._compute_earnings_surprise(earnings)
# Build unavailable_fields from 402 endpoints
unavailable_fields: dict[str, str] = {
self._ENDPOINT_FIELD_MAP[ep]: "requires paid plan"
for ep in endpoints_402
if ep in self._ENDPOINT_FIELD_MAP
}
return FundamentalData(
ticker=ticker,
pe_ratio=pe_ratio,
revenue_growth=revenue_growth,
earnings_surprise=earnings_surprise,
market_cap=market_cap,
fetched_at=datetime.now(timezone.utc),
unavailable_fields=unavailable_fields,
)
except (ProviderError, RateLimitError):
raise
except Exception as exc:
logger.error("FMP provider error for %s: %s", ticker, exc)
raise ProviderError(f"FMP provider error for {ticker}: {exc}") from exc
async def _fetch_json(
self,
client: httpx.AsyncClient,
endpoint: str,
params: dict,
ticker: str,
) -> dict:
"""Fetch a stable endpoint and return the first item (or empty dict)."""
url = f"{_FMP_STABLE_URL}/{endpoint}"
resp = await client.get(url, params=params)
self._check_response(resp, ticker, endpoint)
data = resp.json()
if isinstance(data, list):
return data[0] if data else {}
return data if isinstance(data, dict) else {}
async def _fetch_json_optional(
self,
client: httpx.AsyncClient,
endpoint: str,
params: dict,
ticker: str,
) -> tuple[dict, bool]:
"""Fetch a stable endpoint, returning ``({}, True)`` on 402 (paid tier).
Returns a tuple of (data_dict, was_402) so callers can track which
endpoints required a paid plan.
"""
url = f"{_FMP_STABLE_URL}/{endpoint}"
resp = await client.get(url, params=params)
if resp.status_code == 402:
logger.warning("FMP %s requires paid plan — skipping for %s", endpoint, ticker)
return {}, True
self._check_response(resp, ticker, endpoint)
data = resp.json()
if isinstance(data, list):
return (data[0] if data else {}, False)
return (data if isinstance(data, dict) else {}, False)
def _compute_earnings_surprise(self, earnings_data: dict) -> float | None:
"""Compute earnings surprise % from the most recent actual vs estimated EPS."""
actual = self._safe_float(earnings_data.get("epsActual"))
estimated = self._safe_float(earnings_data.get("epsEstimated"))
if actual is None or estimated is None or estimated == 0:
return None
return ((actual - estimated) / abs(estimated)) * 100
def _check_response(
self, resp: httpx.Response, ticker: str, endpoint: str
) -> None:
"""Raise appropriate errors for non-200 responses."""
if resp.status_code == 429:
raise RateLimitError(f"FMP rate limit hit for {ticker} ({endpoint})")
if resp.status_code == 403:
raise ProviderError(
f"FMP {endpoint} access denied for {ticker}: HTTP 403 — check API key validity and plan tier"
)
if resp.status_code != 200:
raise ProviderError(
f"FMP {endpoint} error for {ticker}: HTTP {resp.status_code}"
)
@staticmethod
def _safe_float(value: object) -> float | None:
"""Convert a value to float, returning None on failure."""
if value is None:
return None
try:
return float(value)
except (TypeError, ValueError):
return None
-354
View File
@@ -1,354 +0,0 @@
"""Chained fundamentals provider with fallback adapters.
Order:
1) FMP (if configured)
2) Finnhub (if configured)
3) Alpha Vantage (if configured)
"""
from __future__ import annotations
import logging
import os
from datetime import date, datetime, timedelta, timezone
from pathlib import Path
import httpx
from app.config import settings
from app.exceptions import ProviderError, RateLimitError
from app.providers.fmp import FMPFundamentalProvider
from app.providers.protocol import FundamentalData, FundamentalProvider
logger = logging.getLogger(__name__)
_CA_BUNDLE = os.environ.get("SSL_CERT_FILE", "")
if not _CA_BUNDLE or not Path(_CA_BUNDLE).exists():
_CA_BUNDLE_PATH: str | bool = True
else:
_CA_BUNDLE_PATH = _CA_BUNDLE
def _safe_float(value: object) -> float | None:
if value is None:
return None
try:
return float(value)
except (TypeError, ValueError):
return None
def _to_api_symbol(symbol: str) -> str:
"""Convert internal symbol format (BRK-B) to API format (BRK.B).
Finnhub and Alpha Vantage use dot-separated share class notation.
"""
return symbol.replace("-", ".")
class FinnhubFundamentalProvider:
"""Fundamentals provider backed by Finnhub free endpoints."""
def __init__(self, api_key: str) -> None:
if not api_key:
raise ProviderError("Finnhub API key is required")
self._api_key = api_key
self._base_url = "https://finnhub.io/api/v1"
async def fetch_fundamentals(self, ticker: str) -> FundamentalData:
unavailable: dict[str, str] = {}
api_symbol = _to_api_symbol(ticker)
today = date.today()
async with httpx.AsyncClient(timeout=30.0, verify=_CA_BUNDLE_PATH) as client:
profile_resp = await client.get(
f"{self._base_url}/stock/profile2",
params={"symbol": api_symbol, "token": self._api_key},
)
metric_resp = await client.get(
f"{self._base_url}/stock/metric",
params={"symbol": api_symbol, "metric": "all", "token": self._api_key},
)
earnings_resp = await client.get(
f"{self._base_url}/stock/earnings",
params={"symbol": api_symbol, "limit": 1, "token": self._api_key},
)
calendar_resp = await client.get(
f"{self._base_url}/calendar/earnings",
params={
"symbol": api_symbol,
"from": today.isoformat(),
"to": (today + timedelta(days=120)).isoformat(),
"token": self._api_key,
},
)
for resp, endpoint in (
(profile_resp, "profile2"),
(metric_resp, "stock/metric"),
(earnings_resp, "stock/earnings"),
(calendar_resp, "calendar/earnings"),
):
if resp.status_code == 429:
raise RateLimitError(f"Finnhub rate limit hit for {ticker} ({endpoint})")
if resp.status_code in (401, 403):
raise ProviderError(f"Finnhub access denied for {ticker} ({endpoint}): HTTP {resp.status_code}")
if resp.status_code != 200:
raise ProviderError(f"Finnhub error for {ticker} ({endpoint}): HTTP {resp.status_code}")
profile_payload = profile_resp.json() if profile_resp.text else {}
metric_payload = metric_resp.json() if metric_resp.text else {}
earnings_payload = earnings_resp.json() if earnings_resp.text else []
metrics = metric_payload.get("metric", {}) if isinstance(metric_payload, dict) else {}
# Finnhub profile2 marketCapitalization is in millions of USD.
# Normalize to absolute dollars so cap bands / formatters match FMP & Alpha Vantage.
market_cap_millions = _safe_float((profile_payload or {}).get("marketCapitalization"))
market_cap = market_cap_millions * 1_000_000.0 if market_cap_millions is not None else None
pe_ratio = _safe_float(metrics.get("peTTM") or metrics.get("peNormalizedAnnual"))
revenue_growth = _safe_float(metrics.get("revenueGrowthTTMYoy") or metrics.get("revenueGrowth5Y"))
earnings_surprise = None
if isinstance(earnings_payload, list) and earnings_payload:
first = earnings_payload[0] if isinstance(earnings_payload[0], dict) else {}
earnings_surprise = _safe_float(first.get("surprisePercent"))
next_earnings_date = self._next_earnings(calendar_resp)
if pe_ratio is None:
unavailable["pe_ratio"] = "not available from provider payload"
if revenue_growth is None:
unavailable["revenue_growth"] = "not available from provider payload"
if earnings_surprise is None:
unavailable["earnings_surprise"] = "not available from provider payload"
if market_cap is None:
unavailable["market_cap"] = "not available from provider payload"
return FundamentalData(
ticker=ticker,
pe_ratio=pe_ratio,
revenue_growth=revenue_growth,
earnings_surprise=earnings_surprise,
market_cap=market_cap,
fetched_at=datetime.now(timezone.utc),
next_earnings_date=next_earnings_date,
unavailable_fields=unavailable,
)
@staticmethod
def _next_earnings(resp: httpx.Response) -> date | None:
"""Earliest upcoming earnings date from Finnhub's calendar payload."""
try:
payload = resp.json() if resp.text else {}
except ValueError:
return None
entries = payload.get("earningsCalendar", []) if isinstance(payload, dict) else []
dates: list[date] = []
today = date.today()
for entry in entries if isinstance(entries, list) else []:
raw = entry.get("date") if isinstance(entry, dict) else None
if not raw:
continue
try:
parsed = date.fromisoformat(raw)
except ValueError:
continue
if parsed >= today:
dates.append(parsed)
return min(dates) if dates else None
class AlphaVantageFundamentalProvider:
"""Fundamentals provider backed by Alpha Vantage free endpoints."""
def __init__(self, api_key: str) -> None:
if not api_key:
raise ProviderError("Alpha Vantage API key is required")
self._api_key = api_key
self._base_url = "https://www.alphavantage.co/query"
async def fetch_fundamentals(self, ticker: str) -> FundamentalData:
unavailable: dict[str, str] = {}
api_symbol = _to_api_symbol(ticker)
async with httpx.AsyncClient(timeout=30.0, verify=_CA_BUNDLE_PATH) as client:
overview_resp = await client.get(
self._base_url,
params={"function": "OVERVIEW", "symbol": api_symbol, "apikey": self._api_key},
)
earnings_resp = await client.get(
self._base_url,
params={"function": "EARNINGS", "symbol": api_symbol, "apikey": self._api_key},
)
income_resp = await client.get(
self._base_url,
params={"function": "INCOME_STATEMENT", "symbol": api_symbol, "apikey": self._api_key},
)
for resp, endpoint in (
(overview_resp, "OVERVIEW"),
(earnings_resp, "EARNINGS"),
(income_resp, "INCOME_STATEMENT"),
):
if resp.status_code == 429:
raise RateLimitError(f"Alpha Vantage rate limit hit for {ticker} ({endpoint})")
if resp.status_code != 200:
raise ProviderError(f"Alpha Vantage error for {ticker} ({endpoint}): HTTP {resp.status_code}")
overview = overview_resp.json() if overview_resp.text else {}
earnings = earnings_resp.json() if earnings_resp.text else {}
income = income_resp.json() if income_resp.text else {}
if isinstance(overview, dict) and overview.get("Information"):
raise ProviderError(f"Alpha Vantage unavailable for {ticker}: {overview.get('Information')}")
if isinstance(overview, dict) and overview.get("Note"):
raise RateLimitError(f"Alpha Vantage rate limit for {ticker}: {overview.get('Note')}")
pe_ratio = _safe_float((overview or {}).get("PERatio"))
market_cap = _safe_float((overview or {}).get("MarketCapitalization"))
earnings_surprise = None
quarterly = earnings.get("quarterlyEarnings", []) if isinstance(earnings, dict) else []
if isinstance(quarterly, list) and quarterly:
first = quarterly[0] if isinstance(quarterly[0], dict) else {}
earnings_surprise = _safe_float(first.get("surprisePercentage"))
revenue_growth = None
annual = income.get("annualReports", []) if isinstance(income, dict) else []
if isinstance(annual, list) and len(annual) >= 2:
curr = _safe_float((annual[0] or {}).get("totalRevenue"))
prev = _safe_float((annual[1] or {}).get("totalRevenue"))
if curr is not None and prev not in (None, 0):
revenue_growth = ((curr - prev) / abs(prev)) * 100.0
if pe_ratio is None:
unavailable["pe_ratio"] = "not available from provider payload"
if revenue_growth is None:
unavailable["revenue_growth"] = "not available from provider payload"
if earnings_surprise is None:
unavailable["earnings_surprise"] = "not available from provider payload"
if market_cap is None:
unavailable["market_cap"] = "not available from provider payload"
return FundamentalData(
ticker=ticker,
pe_ratio=pe_ratio,
revenue_growth=revenue_growth,
earnings_surprise=earnings_surprise,
market_cap=market_cap,
fetched_at=datetime.now(timezone.utc),
unavailable_fields=unavailable,
)
_FUNDAMENTAL_FIELDS = ("pe_ratio", "revenue_growth", "earnings_surprise", "market_cap")
class ChainedFundamentalProvider:
"""Merge fundamentals across providers, filling gaps from later sources.
A single provider rarely covers everything on free tiers — FMP's free plan,
for example, returns only market cap (the ratios/growth/earnings endpoints
402). Rather than stop at the first provider with *any* field, we take each
field from the first provider that supplies it, so FMP's market cap is
combined with Finnhub's P/E and earnings surprise.
"""
def __init__(self, providers: list[tuple[str, FundamentalProvider]]) -> None:
if not providers:
raise ProviderError("No fundamental providers configured")
self._providers = providers
async def fetch_fundamentals(self, ticker: str, allow_partial: bool = False) -> FundamentalData:
"""Merge fundamentals across providers.
``allow_partial`` controls behaviour when a fallback provider is *rate
limited* and we end up with missing fields. By default we raise
RateLimitError so the caller (the bulk collector) can back off and retry
the ticker once the window frees — otherwise a transient 429 on Finnhub
would be silently stored as market-cap-only. Pass ``allow_partial=True``
(manual single fetches, or the collector's final give-up attempt) to
accept whatever was gathered instead of raising.
"""
merged: dict[str, float | None] = {f: None for f in _FUNDAMENTAL_FIELDS}
field_source: dict[str, str] = {}
errors: list[str] = []
rate_limited = False
next_earnings_date = None
for provider_name, provider in self._providers:
if all(merged[f] is not None for f in _FUNDAMENTAL_FIELDS) and next_earnings_date:
break
try:
data = await provider.fetch_fundamentals(ticker)
except RateLimitError as exc:
rate_limited = True
errors.append(f"{provider_name}: RateLimitError: {exc}")
continue
except Exception as exc:
errors.append(f"{provider_name}: {type(exc).__name__}: {exc}")
continue
if next_earnings_date is None and data.next_earnings_date is not None:
next_earnings_date = data.next_earnings_date
for field in _FUNDAMENTAL_FIELDS:
if merged[field] is None:
value = getattr(data, field)
if value is not None:
merged[field] = value
field_source[field] = provider_name
missing = [f for f in _FUNDAMENTAL_FIELDS if merged[f] is None]
# A rate limit left data incomplete: signal it (unless partial is OK) so
# the collector backs off rather than persisting a degraded record.
if rate_limited and missing and not allow_partial:
attempts = "; ".join(errors[:6])
raise RateLimitError(
f"Fundamentals incomplete for {ticker} due to provider rate limits "
f"(missing {', '.join(missing)}). Attempts: {attempts}"
)
if all(merged[f] is None for f in _FUNDAMENTAL_FIELDS):
attempts = "; ".join(errors[:6]) if errors else "no usable metrics from any provider"
raise ProviderError(f"All fundamentals providers failed for {ticker}. Attempts: {attempts}")
unavailable: dict[str, str] = {
field: "not available from any configured provider"
for field in _FUNDAMENTAL_FIELDS
if merged[field] is None
}
# Record which provider supplied each field for transparency.
for field, src in field_source.items():
unavailable[f"source_{field}"] = src
return FundamentalData(
ticker=ticker,
pe_ratio=merged["pe_ratio"],
revenue_growth=merged["revenue_growth"],
earnings_surprise=merged["earnings_surprise"],
market_cap=merged["market_cap"],
fetched_at=datetime.now(timezone.utc),
next_earnings_date=next_earnings_date,
unavailable_fields=unavailable,
)
def build_fundamental_provider_chain() -> FundamentalProvider:
providers: list[tuple[str, FundamentalProvider]] = []
if settings.fmp_api_key:
providers.append(("fmp", FMPFundamentalProvider(settings.fmp_api_key)))
if settings.finnhub_api_key:
providers.append(("finnhub", FinnhubFundamentalProvider(settings.finnhub_api_key)))
if settings.alpha_vantage_api_key:
providers.append(("alpha_vantage", AlphaVantageFundamentalProvider(settings.alpha_vantage_api_key)))
if not providers:
raise ProviderError(
"No fundamentals provider configured. Set one of FMP_API_KEY, FINNHUB_API_KEY, ALPHA_VANTAGE_API_KEY"
)
logger.info("Fundamentals provider chain configured: %s", [name for name, _ in providers])
return ChainedFundamentalProvider(providers)
+2 -20
View File
@@ -44,20 +44,6 @@ class SentimentData:
recommendation: str | None = None # "buy" | "hold" | "avoid" — actionable LLM view
@dataclass(frozen=True, slots=True)
class FundamentalData:
"""Fundamental metrics returned by fundamental providers."""
ticker: str
pe_ratio: float | None
revenue_growth: float | None
earnings_surprise: float | None
market_cap: float | None
fetched_at: datetime
next_earnings_date: date | None = None
unavailable_fields: dict[str, str] = field(default_factory=dict)
# ---------------------------------------------------------------------------
# Provider Protocols
# ---------------------------------------------------------------------------
@@ -81,9 +67,5 @@ class SentimentProvider(Protocol):
...
class FundamentalProvider(Protocol):
"""Protocol for fundamental data providers."""
async def fetch_fundamentals(self, ticker: str) -> FundamentalData:
"""Fetch fundamental data for a ticker."""
...
# No fundamentals provider protocol: since A6 fundamentals come only from the
# batch SEC/Dolt imports, never from a request-time provider call.
-52
View File
@@ -13,7 +13,6 @@ from app.schemas.admin import (
AlertConfigUpdate,
CreateUserRequest,
DataCleanupRequest,
FundamentalsCutoverConfigUpdate,
JobTriggerRequest,
JobToggle,
RecommendationConfigUpdate,
@@ -138,27 +137,6 @@ async def list_settings(
)
@router.get("/admin/settings/fundamentals-cutover", response_model=APIEnvelope)
async def get_fundamentals_cutover_settings(
_admin: User = Depends(require_admin),
db: AsyncSession = Depends(get_db),
):
config = await admin_service.get_fundamentals_cutover_config(db)
return APIEnvelope(status="success", data=config)
@router.put("/admin/settings/fundamentals-cutover", response_model=APIEnvelope)
async def update_fundamentals_cutover_settings(
body: FundamentalsCutoverConfigUpdate,
_admin: User = Depends(require_admin),
db: AsyncSession = Depends(get_db),
):
config = await admin_service.update_fundamentals_cutover_config(
db, body.enabled
)
return APIEnvelope(status="success", data=config)
@router.get("/admin/settings/recommendations", response_model=APIEnvelope)
async def get_recommendation_settings(
_admin: User = Depends(require_admin),
@@ -475,36 +453,6 @@ async def toggle_job(
)
@router.get("/admin/fundamentals-parity", response_model=APIEnvelope)
async def get_fundamentals_parity_report(
_admin: User = Depends(require_admin),
):
"""Latest read-only A5 source/score comparison, or null before first run."""
return APIEnvelope(
status="success", data=admin_service.get_fundamentals_parity_report()
)
@router.get("/admin/fundamentals-parity/csv", response_model=APIEnvelope)
async def get_fundamentals_parity_csv(
_admin: User = Depends(require_admin),
):
"""Latest flattened A5 report for an authenticated browser download."""
artifact = admin_service.get_fundamentals_parity_csv()
data = None if artifact is None else {"filename": artifact[0], "content": artifact[1]}
return APIEnvelope(status="success", data=data)
@router.get("/admin/fundamentals-parity/json", response_model=APIEnvelope)
async def get_fundamentals_parity_json(
_admin: User = Depends(require_admin),
):
"""Canonical A5 JSON artifact for an authenticated browser download."""
artifact = admin_service.get_fundamentals_parity_json()
data = None if artifact is None else {"filename": artifact[0], "content": artifact[1]}
return APIEnvelope(status="success", data=data)
# ---------------------------------------------------------------------------
# System events (operational warnings / errors)
# ---------------------------------------------------------------------------
+7 -29
View File
@@ -23,7 +23,6 @@ from app.models.sr_level import SRLevel
from app.models.ticker import Ticker
from app.models.user import User
from app.providers.alpaca import AlpacaOHLCVProvider
from app.providers.fundamentals_chain import build_fundamental_provider_chain
from app.services.rr_scanner_service import (
resolve_activation_ranks_for_symbol,
scan_ticker,
@@ -31,7 +30,6 @@ from app.services.rr_scanner_service import (
from app.services.sentiment_provider_service import build_sentiment_provider
from app.schemas.common import APIEnvelope
from app.services import (
fundamental_service,
ingestion_service,
scoring_service,
sentiment_service,
@@ -185,34 +183,14 @@ async def fetch_symbol(
sources_out["sentiment"] = {"status": "error", "message": str(exc)}
# --- Fundamentals ---
# No per-ticker fetch exists any more: fundamental_data is rebuilt for the
# whole universe by the nightly SEC + Dolt imports, from local PostgreSQL.
# The source key is still accepted so older clients get a truthful answer.
if "fundamentals" in requested:
if settings.fmp_api_key or settings.finnhub_api_key or settings.alpha_vantage_api_key:
try:
fundamentals_provider = build_fundamental_provider_chain()
# Manual single fetch: take whatever we can get (a lone 429 on a
# fallback shouldn't fail the whole refresh).
fdata = await fundamentals_provider.fetch_fundamentals(
symbol_upper, allow_partial=True
)
await fundamental_service.store_fundamental(
db,
symbol=symbol_upper,
pe_ratio=fdata.pe_ratio,
revenue_growth=fdata.revenue_growth,
earnings_surprise=fdata.earnings_surprise,
market_cap=fdata.market_cap,
next_earnings_date=fdata.next_earnings_date,
unavailable_fields=fdata.unavailable_fields,
)
sources_out["fundamentals"] = {"status": "ok", "message": None}
except Exception as exc:
logger.error("Fundamentals fetch failed for %s: %s", symbol_upper, exc)
sources_out["fundamentals"] = {"status": "error", "message": str(exc)}
else:
sources_out["fundamentals"] = {
"status": "skipped",
"message": "No fundamentals provider key configured",
}
sources_out["fundamentals"] = {
"status": "skipped",
"message": "Fundamentals refresh nightly from the SEC + Dolt imports",
}
# --- Derived pipeline: S/R levels (free, always) ---
try:
+32 -260
View File
@@ -1,9 +1,9 @@
"""APScheduler job definitions and FastAPI lifespan integration.
Defines four scheduled jobs:
Defines the scheduled jobs, among them:
- Data Collector (OHLCV fetch for all tickers)
- Sentiment Collector (sentiment for all tickers)
- Fundamental Collector (fundamentals for all tickers)
- Dolt Earnings / SEC Fundamentals imports (bulk fundamentals sources)
- R:R Scanner (trade setup scan for all tickers)
Each job processes tickers independently, logs errors as structured JSON,
@@ -25,22 +25,18 @@ from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
from app.database import async_session_factory
from app.models.fundamental import FundamentalData
from app.models.ohlcv import OHLCVRecord
from app.models.sentiment import SentimentScore
from app.models.ticker import Ticker
from app.exceptions import ProviderError
from app.providers.alpaca import AlpacaOHLCVProvider
from app.providers.fundamentals_chain import build_fundamental_provider_chain
from app.providers.protocol import SentimentData
from app.services import (
fundamental_service,
ingestion_service,
pipeline_run,
sentiment_service,
settings_store,
shadow_book_service,
fundamentals_parity_service,
fundamental_data_refresh_service,
)
from app.services.data_import import (
@@ -93,7 +89,6 @@ _last_successful: dict[str, str | None] = {
"data_collector": None,
"data_backfill": None,
"sentiment_collector": None,
"fundamental_collector": None,
}
# Jobs whose per-run progress is surfaced to Admin → Jobs. (outcome_evaluator is
@@ -102,10 +97,8 @@ _JOB_NAMES = [
"data_collector",
"data_backfill",
"sentiment_collector",
"fundamental_collector",
"dolt_earnings_import",
"sec_fundamentals_import",
"fundamentals_parity_report",
"rr_scanner",
"ticker_universe_sync",
"alerts",
@@ -466,23 +459,6 @@ async def _get_sentiment_priority_tickers(db: AsyncSession) -> list[str]:
return priority_syms + filler_syms
async def _get_fundamental_priority_tickers(db: AsyncSession) -> list[str]:
"""Return symbols prioritized for fundamentals refresh.
Priority:
1) Tickers with no fundamentals snapshot yet
2) Tickers with existing fundamentals, oldest fetched_at first
3) Alphabetical tiebreaker
"""
missing_first = case((FundamentalData.fetched_at.is_(None), 0), else_=1)
result = await db.execute(
select(Ticker.symbol)
.outerjoin(FundamentalData, FundamentalData.ticker_id == Ticker.id)
.order_by(missing_first.asc(), FundamentalData.fetched_at.asc(), Ticker.symbol.asc())
)
return list(result.scalars().all())
def _resume_tickers(symbols: list[str], job_name: str) -> list[str]:
"""Reorder tickers to resume after the last successful one (rate-limit resume).
@@ -815,139 +791,6 @@ async def collect_sentiment() -> None:
_runtime_finish(job_name, "error", processed=processed, total=total, message=str(exc))
# ---------------------------------------------------------------------------
# Job: Fundamental Collector
# ---------------------------------------------------------------------------
async def collect_fundamentals() -> None:
"""Fetch fundamentals for all tracked tickers via FMP.
Processes each ticker independently. On rate limit, records last
successful ticker for resume.
"""
job_name = "fundamental_collector"
_log_event(logging.INFO, "job_start", job=job_name)
_runtime_start(job_name)
processed = 0
total: int | None = None
try:
async with async_session_factory() as db:
if not await _is_job_enabled(db, job_name):
_log_event(logging.INFO, "job_skipped", job=job_name, reason="disabled")
_runtime_finish(job_name, "skipped", processed=0, total=0, message="Disabled")
return
if await fundamental_data_refresh_service.is_enabled(db):
message = "SEC + Dolt fundamentals cutover is active"
_log_event(
logging.INFO,
"job_skipped",
job=job_name,
reason="sec_dolt_cutover_active",
)
_runtime_finish(
job_name,
"skipped",
processed=0,
total=0,
message=message,
)
return
symbols = await _get_fundamental_priority_tickers(db)
if not symbols:
_log_event(logging.INFO, "job_complete", job=job_name, tickers=0)
_runtime_finish(job_name, "completed", processed=0, total=0, message="No tickers")
return
total = len(symbols)
_runtime_progress(job_name, processed=0, total=total)
if not (settings.fmp_api_key or settings.finnhub_api_key or settings.alpha_vantage_api_key):
_log_event(logging.WARNING, "job_skipped", job=job_name, reason="no fundamentals provider keys configured")
_runtime_finish(job_name, "skipped", processed=0, total=total, message="No fundamentals provider keys configured")
return
try:
provider = build_fundamental_provider_chain()
except Exception as exc:
_log_event(logging.ERROR, "job_error", job=job_name, error_type=type(exc).__name__, message=str(exc))
_runtime_finish(job_name, "error", processed=0, total=total, message=str(exc))
return
max_retries = max(0, settings.fundamental_rate_limit_retries)
base_backoff = max(1, settings.fundamental_rate_limit_backoff_seconds)
spacing = max(0.0, settings.fundamental_request_spacing_seconds)
async def _store(symbol: str, data) -> None:
async with async_session_factory() as db:
await fundamental_service.store_fundamental(
db,
symbol=symbol,
pe_ratio=data.pe_ratio,
revenue_growth=data.revenue_growth,
earnings_surprise=data.earnings_surprise,
market_cap=data.market_cap,
next_earnings_date=data.next_earnings_date,
unavailable_fields=data.unavailable_fields,
)
for symbol in symbols:
_runtime_progress(job_name, processed=processed, total=total, current_ticker=symbol)
attempt = 0
while True:
try:
data = await provider.fetch_fundamentals(symbol)
await _store(symbol, data)
_last_successful[job_name] = symbol
processed += 1
_runtime_progress(job_name, processed=processed, total=total, current_ticker=symbol)
_log_event(logging.INFO, "ticker_collected", job=job_name, ticker=symbol)
break
except Exception as exc:
msg = str(exc).lower()
if "rate" in msg or "429" in msg:
if attempt < max_retries:
wait_seconds = base_backoff * (2 ** attempt)
attempt += 1
_log_event(logging.WARNING, "rate_limited_retry", job=job_name, ticker=symbol, attempt=attempt, max_retries=max_retries, wait_seconds=wait_seconds, processed=processed)
_runtime_progress(
job_name,
processed=processed,
total=total,
current_ticker=symbol,
message=f"Rate-limited at {symbol}; retry {attempt}/{max_retries} in {wait_seconds}s",
)
await asyncio.sleep(wait_seconds)
continue
# Retries exhausted: store whatever partial data we can
# still get (e.g. FMP market cap) and move on, rather than
# aborting the whole run and leaving every later ticker
# untouched.
_log_event(logging.WARNING, "rate_limited_partial", job=job_name, ticker=symbol, processed=processed)
try:
data = await provider.fetch_fundamentals(symbol, allow_partial=True)
await _store(symbol, data)
processed += 1
except Exception as exc2:
_log_job_error(job_name, symbol, exc2)
break
_log_job_error(job_name, symbol, exc)
break
if spacing:
await asyncio.sleep(spacing)
_last_successful[job_name] = None
_log_event(logging.INFO, "job_complete", job=job_name, tickers=processed)
_runtime_finish(job_name, "completed", processed=processed, total=total, message=f"Processed {processed} tickers")
except Exception as exc:
_log_event(logging.ERROR, "job_error", job=job_name, error_type=type(exc).__name__, message=str(exc))
_runtime_finish(job_name, "error", processed=processed, total=total, message=str(exc))
# ---------------------------------------------------------------------------
# Jobs: shadow fundamentals sources
# ---------------------------------------------------------------------------
@@ -956,9 +799,9 @@ async def collect_fundamentals() -> None:
async def _run_shadow_import(job_name: str, importer: SourceImporter) -> bool:
"""Run an importer and return whether its scheduled job was enabled.
The SEC wrapper uses the return value to run its activated local cache step
after deferred, failed, no-op, promoted, or source-locked attempts while honoring
the job-level disable switch.
The SEC wrapper uses the return value only to word its runtime message: its
local cache step runs after deferred, failed, no-op, promoted, source-locked
and disabled attempts alike.
"""
_log_event(logging.INFO, "job_start", job=job_name)
_runtime_start(job_name, total=1)
@@ -968,7 +811,7 @@ async def _run_shadow_import(job_name: str, importer: SourceImporter) -> bool:
if not await _is_job_enabled(db, job_name):
_log_event(logging.INFO, "job_skipped", job=job_name, reason="disabled")
_runtime_finish(job_name, "skipped", processed=0, total=1, message="Disabled")
return
return False
run = await run_import(importer)
if run is None:
@@ -1015,25 +858,26 @@ async def _run_shadow_import(job_name: str, importer: SourceImporter) -> bool:
async def run_dolt_earnings_import() -> None:
"""Pull and import the Dolt earnings calendar/results feed in shadow."""
"""Pull and import the Dolt earnings calendar/results feed."""
await _run_shadow_import("dolt_earnings_import", DoltEarningsImporter())
async def run_sec_fundamentals_import() -> None:
"""Import SEC facts, then run the activated local compat-cache refresh.
"""Import SEC facts, then refresh the local compat cache.
The refresh is deliberately separate from the network import result. Once
activated it therefore still runs from stored snapshots/earnings/prices when
SEC is unavailable, unchanged, or another SEC import owns the source lock.
The refresh is deliberately independent of the network import: it reads only
stored snapshots, earnings events and closes, so it runs identically when SEC
is unavailable, unchanged, or owned by another import — and also when the
job's ingestion is switched off in Admin → Jobs. Disabling the job stops
SEC network access, not the cache; prices and earnings move daily even when
no filing does, and `fundamental_data` feeds scoring.
"""
job_name = "sec_fundamentals_import"
job_enabled = await _run_shadow_import(job_name, SecFundamentalsImporter())
if not job_enabled:
return
import_ran = await _run_shadow_import(job_name, SecFundamentalsImporter())
try:
async with async_session_factory() as db:
summary = await fundamental_data_refresh_service.refresh_if_enabled(db)
summary = await fundamental_data_refresh_service.refresh(db)
except asyncio.CancelledError:
_runtime_finish(
job_name, "error", processed=0, total=1, message="Cancelled"
@@ -1051,29 +895,27 @@ async def run_sec_fundamentals_import() -> None:
_runtime_finish(job_name, "error", processed=0, total=1, message=message)
return
if not summary["enabled"]:
_log_event(
logging.INFO,
"fundamental_data_refresh_skipped",
job=job_name,
reason="cutover_disabled",
setting=fundamental_data_refresh_service.ACTIVATION_KEY,
)
return
_log_event(
logging.INFO,
"fundamental_data_refresh_complete",
job=job_name,
**summary,
)
cache_message = (
f"cache {summary['refreshed']} · "
f"{summary['score_inputs_changed']} score inputs changed"
)
runtime = get_job_runtime_snapshot(job_name)
if runtime.get("status") == "completed":
import_message = runtime.get("message") or "import completed"
cache_message = (
f"cache {summary['refreshed']} · "
f"{summary['score_inputs_changed']} score inputs changed"
if not import_ran:
_runtime_finish(
job_name,
"completed",
processed=1,
total=1,
message=f"Import disabled · {cache_message}",
)
elif runtime.get("status") == "completed":
import_message = runtime.get("message") or "import completed"
_runtime_finish(
job_name,
"completed",
@@ -1083,49 +925,6 @@ async def run_sec_fundamentals_import() -> None:
)
async def run_fundamentals_parity_report() -> None:
"""Generate the A5 comparison bundle without mutating live fundamentals/scores."""
job_name = "fundamentals_parity_report"
_log_event(logging.INFO, "job_start", job=job_name)
_runtime_start(job_name, total=1)
try:
async with async_session_factory() as db:
if not await _is_job_enabled(db, job_name):
_runtime_finish(
job_name, "skipped", processed=0, total=1, message="Disabled"
)
return
report, artifacts = await fundamentals_parity_service.generate_and_store(
db, settings.fundamentals_parity_report_dir
)
summary = report["summary"]
message = (
f"{summary['universe_count']} tickers · "
f"{summary['fundamental_score_material_changes']} material score changes"
)
_runtime_finish(job_name, "completed", processed=1, total=1, message=message)
_log_event(
logging.INFO,
"job_complete",
job=job_name,
generated_at=report["generated_at"],
json_path=artifacts["json"],
csv_path=artifacts["csv"],
)
except asyncio.CancelledError:
_runtime_finish(job_name, "error", processed=0, total=1, message="Cancelled")
raise
except Exception as exc:
_runtime_finish(job_name, "error", processed=0, total=1, message=str(exc))
_log_event(
logging.ERROR,
"job_error",
job=job_name,
error_type=type(exc).__name__,
message=str(exc),
)
# ---------------------------------------------------------------------------
# Job: R:R Scanner
# ---------------------------------------------------------------------------
@@ -1666,19 +1465,16 @@ SCHEDULE_DEFAULTS: dict[str, str] = {
"schedule_timezone": "America/New_York",
# Morning data/display refresh (no qualifying R:R scan).
"schedule_daily_pipeline_cron": "0 2 * * *",
# Bulk source imports. The SEC job writes the legacy compat cache only after
# the explicit, default-off A5 cutover setting is enabled.
# Bulk source imports. The SEC job also refreshes the fundamental_data compat
# cache that scoring reads — locally, from stored snapshots/earnings/closes.
"schedule_dolt_earnings_cron": "30 2 * * *",
"schedule_sec_fundamentals_cron": "0 4 * * *",
"schedule_fundamentals_parity_cron": "30 5 * * *",
# Fetch in-progress bars → scan → Telegram (manual MOC window).
"schedule_near_close_pipeline_cron": "30 15 * * mon-fri",
# Fetch final bars → outcome eval (must not run on the partial near-close bar).
"schedule_after_close_pipeline_cron": "45 16 * * mon-fri",
# Hourly mid-session price + outcome (10:0015:00 ET MonFri).
"schedule_intraday_pipeline_cron": "0 10-15 * * mon-fri",
# Weekly fundamentals early Monday NY.
"schedule_fundamentals_cron": "0 1 * * mon",
}
# job id -> schedule setting key
@@ -1686,11 +1482,9 @@ _CRON_JOBS: dict[str, str] = {
"daily_pipeline": "schedule_daily_pipeline_cron",
"dolt_earnings_import": "schedule_dolt_earnings_cron",
"sec_fundamentals_import": "schedule_sec_fundamentals_cron",
"fundamentals_parity_report": "schedule_fundamentals_parity_cron",
"near_close_pipeline": "schedule_near_close_pipeline_cron",
"after_close_pipeline": "schedule_after_close_pipeline_cron",
"intraday_pipeline": "schedule_intraday_pipeline_cron",
"fundamental_collector": "schedule_fundamentals_cron",
}
@@ -1779,7 +1573,7 @@ def configure_scheduler(schedule_config: dict[str, str] | None = None) -> None:
"schedule_dolt_earnings_cron",
),
id="dolt_earnings_import",
name="Dolt Earnings Import (shadow)",
name="Dolt Earnings Import",
replace_existing=True,
)
scheduler.add_job(
@@ -1793,17 +1587,6 @@ def configure_scheduler(schedule_config: dict[str, str] | None = None) -> None:
name="SEC Fundamentals Import",
replace_existing=True,
)
scheduler.add_job(
run_fundamentals_parity_report,
_cron_trigger(
cfg["schedule_fundamentals_parity_cron"],
tz,
"schedule_fundamentals_parity_cron",
),
id="fundamentals_parity_report",
name="Fundamentals Parity Report (read-only)",
replace_existing=True,
)
scheduler.add_job(
run_near_close_pipeline,
_cron_trigger(
@@ -1831,13 +1614,6 @@ def configure_scheduler(schedule_config: dict[str, str] | None = None) -> None:
_cron_trigger(cfg["schedule_intraday_pipeline_cron"], tz, "schedule_intraday_pipeline_cron"),
id="intraday_pipeline", name="Intraday Pipeline", replace_existing=True,
)
# Fundamentals — quarterly-ish data; weekly by default (conserves API quota).
# Its own early cron so the slow, rate-limited fetch finishes before the day.
scheduler.add_job(
collect_fundamentals,
_cron_trigger(cfg["schedule_fundamentals_cron"], tz, "schedule_fundamentals_cron"),
id="fundamental_collector", name="Fundamental Collector", replace_existing=True,
)
# Independent interval jobs (own cadence, no ordering dependency)
scheduler.add_job(
@@ -1879,9 +1655,6 @@ def configure_scheduler(schedule_config: dict[str, str] | None = None) -> None:
},
dolt_earnings_import={"cron": cfg["schedule_dolt_earnings_cron"]},
sec_fundamentals_import={"cron": cfg["schedule_sec_fundamentals_cron"]},
fundamentals_parity_report={
"cron": cfg["schedule_fundamentals_parity_cron"]
},
near_close_pipeline={
"cron": cfg["schedule_near_close_pipeline_cron"],
"steps": [name for name, _ in _NEAR_CLOSE_PIPELINE_STEPS],
@@ -1894,7 +1667,6 @@ def configure_scheduler(schedule_config: dict[str, str] | None = None) -> None:
"cron": cfg["schedule_intraday_pipeline_cron"],
"steps": [name for name, _ in _INTRADAY_PIPELINE_STEPS],
},
fundamental_collector={"cron": cfg["schedule_fundamentals_cron"]},
independent=["ticker_universe_sync", "backtest"],
manual_only=["alerts", "data_backfill", "event_study"],
)
-7
View File
@@ -73,11 +73,6 @@ class ActivationConfigUpdate(BaseModel):
exclude_neutral: bool | None = None
class FundamentalsCutoverConfigUpdate(BaseModel):
"""Switch the legacy fundamentals cache from quota APIs to SEC/Dolt."""
enabled: bool
class ScheduleConfigUpdate(BaseModel):
"""Cron schedule for the pipelines + fundamentals. Crons are 5-field
(min hour dom month dow); timezone is an IANA name (e.g. America/New_York)."""
@@ -85,11 +80,9 @@ class ScheduleConfigUpdate(BaseModel):
schedule_daily_pipeline_cron: str | None = Field(default=None, max_length=120)
schedule_dolt_earnings_cron: str | None = Field(default=None, max_length=120)
schedule_sec_fundamentals_cron: str | None = Field(default=None, max_length=120)
schedule_fundamentals_parity_cron: str | None = Field(default=None, max_length=120)
schedule_near_close_pipeline_cron: str | None = Field(default=None, max_length=120)
schedule_after_close_pipeline_cron: str | None = Field(default=None, max_length=120)
schedule_intraday_pipeline_cron: str | None = Field(default=None, max_length=120)
schedule_fundamentals_cron: str | None = Field(default=None, max_length=120)
class PerformanceConfigUpdate(BaseModel):
+2 -55
View File
@@ -17,7 +17,7 @@ from app.models.settings import SystemSetting
from app.models.ticker import Ticker
from app.models.trade_setup import TradeSetup
from app.models.user import User
from app.services import fundamental_data_refresh_service, settings_store
from app.services import settings_store
logger = logging.getLogger(__name__)
@@ -159,28 +159,6 @@ async def update_setting(db: AsyncSession, key: str, value: str) -> SystemSettin
return setting
# ---------------------------------------------------------------------------
# Fundamentals source cutover
# ---------------------------------------------------------------------------
async def get_fundamentals_cutover_config(db: AsyncSession) -> dict[str, bool]:
"""Return the explicit A5 cache-cutover switch (default off)."""
return {"enabled": await fundamental_data_refresh_service.is_enabled(db)}
async def update_fundamentals_cutover_config(
db: AsyncSession, enabled: bool
) -> dict[str, bool]:
"""Activate or pause SEC/Dolt writes to the legacy fundamentals cache."""
await settings_store.upsert_setting(
db,
fundamental_data_refresh_service.ACTIVATION_KEY,
"true" if enabled else "false",
)
await db.commit()
return await get_fundamentals_cutover_config(db)
# ---------------------------------------------------------------------------
# Activation thresholds
# ---------------------------------------------------------------------------
@@ -633,10 +611,8 @@ VALID_JOB_NAMES = {
"data_backfill",
"benchmark_collector",
"sentiment_collector",
"fundamental_collector",
"dolt_earnings_import",
"sec_fundamentals_import",
"fundamentals_parity_report",
"rr_scanner",
"ticker_universe_sync",
"outcome_evaluator",
@@ -657,10 +633,8 @@ JOB_LABELS = {
"data_backfill": "Data Backfill (deep history)",
"benchmark_collector": "Benchmark Collector",
"sentiment_collector": "Sentiment Collector",
"fundamental_collector": "Fundamental Collector",
"dolt_earnings_import": "Dolt Earnings Import (shadow)",
"dolt_earnings_import": "Dolt Earnings Import",
"sec_fundamentals_import": "SEC Fundamentals Import",
"fundamentals_parity_report": "Fundamentals Parity Report (read-only)",
"rr_scanner": "R:R Scanner",
"ticker_universe_sync": "Ticker Universe Sync",
"outcome_evaluator": "Outcome Evaluator",
@@ -799,30 +773,3 @@ async def toggle_job(db: AsyncSession, job_name: str, enabled: bool) -> SystemSe
key = f"job_{job_name}_enabled"
return await update_setting(db, key, str(enabled).lower())
def get_fundamentals_parity_report() -> dict | None:
"""Return the latest compact A5 summary, if the job has run."""
from app.config import settings
from app.services.fundamentals_parity_service import load_latest
report = load_latest(settings.fundamentals_parity_report_dir)
if report is not None:
report.pop("rows", None) # full per-ticker data is download-only
return report
def get_fundamentals_parity_csv() -> tuple[str, str] | None:
"""Return the latest A5 CSV filename and content for authenticated download."""
from app.config import settings
from app.services.fundamentals_parity_service import load_latest_csv
return load_latest_csv(settings.fundamentals_parity_report_dir)
def get_fundamentals_parity_json() -> tuple[str, str] | None:
"""Return the canonical A5 JSON artifact for authenticated download."""
from app.config import settings
from app.services.fundamentals_parity_service import load_latest_json
return load_latest_json(settings.fundamentals_parity_report_dir)
@@ -1,4 +1,7 @@
"""A5 activation: refresh the legacy fundamentals cache from local bulk data."""
"""Refresh the fundamentals compat cache from local SEC/Dolt bulk data.
``fundamental_data`` is the table scoring reads. This is its only writer.
"""
from __future__ import annotations
@@ -12,38 +15,11 @@ from sqlalchemy.ext.asyncio import AsyncSession
from app.database import insert_for_session
from app.models.fundamental import FundamentalData
from app.models.score import CompositeScore, DimensionScore
from app.services import fundamentals_candidate_service, settings_store
from app.services import fundamentals_candidate_service
# Absence is deliberately false. Production activation therefore requires one
# explicit, durable SystemSetting change after the A5 evidence is approved.
ACTIVATION_KEY = "fundamental_data_sec_dolt_cutover_enabled"
_SCORE_FIELDS = ("pe_ratio", "revenue_growth", "earnings_surprise")
async def is_enabled(db: AsyncSession) -> bool:
raw = await settings_store.get_value(db, ACTIVATION_KEY, "false")
return str(raw).strip().lower() == "true"
async def refresh_if_enabled(
db: AsyncSession,
*,
now: datetime | None = None,
today: date | None = None,
) -> dict[str, Any]:
"""Refresh atomically when activated; otherwise perform no writes."""
if not await is_enabled(db):
return {
"enabled": False,
"refreshed": 0,
"score_inputs_changed": 0,
"dimension_scores_staled": 0,
"composite_scores_staled": 0,
}
return await refresh(db, now=now, today=today)
async def refresh(
db: AsyncSession,
*,
@@ -117,7 +93,6 @@ async def refresh(
await db.commit()
return {
"enabled": True,
"refreshed": len(candidates),
"score_inputs_changed": len(changed_ids),
"dimension_scores_staled": len(dimension_ids),
+5 -67
View File
@@ -1,22 +1,19 @@
"""Fundamental data service.
"""Fundamental data read access.
Stores fundamental data (P/E, revenue growth, earnings surprise, market cap)
and marks the fundamental dimension score as stale on new data.
``fundamental_data`` is the compat cache scoring reads. It is written solely by
``fundamental_data_refresh_service`` from SEC snapshots, Dolt earnings events and
stored closes; nothing fetches it per ticker.
"""
from __future__ import annotations
import json
import logging
from datetime import datetime, timezone
from sqlalchemy import select, update
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.database import insert_for_session
from app.exceptions import NotFoundError
from app.models.fundamental import FundamentalData
from app.models.score import DimensionScore
from app.models.ticker import Ticker
logger = logging.getLogger(__name__)
@@ -32,65 +29,6 @@ async def _get_ticker(db: AsyncSession, symbol: str) -> Ticker:
return ticker
async def store_fundamental(
db: AsyncSession,
symbol: str,
pe_ratio: float | None = None,
revenue_growth: float | None = None,
earnings_surprise: float | None = None,
market_cap: float | None = None,
next_earnings_date=None,
unavailable_fields: dict[str, str] | None = None,
) -> FundamentalData:
"""Store or update fundamental data for a ticker.
Keeps a single latest snapshot per ticker. On new data, marks the
fundamental dimension score as stale (if one exists).
"""
ticker = await _get_ticker(db, symbol)
now = datetime.now(timezone.utc)
unavailable_fields_json = json.dumps(unavailable_fields or {})
stmt = insert_for_session(db, FundamentalData).values(
ticker_id=ticker.id,
pe_ratio=pe_ratio,
revenue_growth=revenue_growth,
earnings_surprise=earnings_surprise,
market_cap=market_cap,
next_earnings_date=next_earnings_date,
fetched_at=now,
unavailable_fields_json=unavailable_fields_json,
)
stmt = stmt.on_conflict_do_update(
index_elements=["ticker_id"],
set_={
"pe_ratio": stmt.excluded.pe_ratio,
"revenue_growth": stmt.excluded.revenue_growth,
"earnings_surprise": stmt.excluded.earnings_surprise,
"market_cap": stmt.excluded.market_cap,
"next_earnings_date": stmt.excluded.next_earnings_date,
"fetched_at": stmt.excluded.fetched_at,
"unavailable_fields_json": stmt.excluded.unavailable_fields_json,
},
).returning(FundamentalData)
record = (await db.execute(stmt)).scalar_one()
# Mark fundamental dimension score as stale if it exists
# TODO: Use DimensionScore service when built
await db.execute(
update(DimensionScore)
.where(
DimensionScore.ticker_id == ticker.id,
DimensionScore.dimension == "fundamental",
)
.values(is_stale=True)
)
await db.commit()
return record
async def get_fundamental(
db: AsyncSession,
symbol: str,
@@ -1,9 +1,8 @@
"""Local SEC/Dolt candidate values for the legacy fundamentals cache.
"""Local SEC/Dolt candidate values for the fundamentals compat cache.
This is the single read path shared by the A5 parity report and the activated
``fundamental_data`` refresh. It never contacts SEC or Dolt: every input comes
from PostgreSQL, so price- and earnings-driven values can still refresh when an
upstream import is unchanged or unavailable.
This is the read path behind the ``fundamental_data`` refresh. It never contacts
SEC or Dolt: every input comes from PostgreSQL, so price- and earnings-driven
values can still refresh when an upstream import is unchanged or unavailable.
"""
from __future__ import annotations
-498
View File
@@ -1,498 +0,0 @@
"""Read-only A5 comparison of legacy and SEC/Dolt fundamental inputs.
The report deliberately does not write ``fundamental_data`` or score tables.
It reconstructs the current legacy and candidate fundamental scores, projects
their composite-score/rank effect with the active weights, and archives a
timestamped JSON + CSV bundle for explicit human approval.
"""
from __future__ import annotations
import csv
import io
import json
import math
import os
import statistics
from datetime import date, datetime, timezone
from pathlib import Path
from typing import Any, Iterable
from zoneinfo import ZoneInfo
from sqlalchemy import select, text
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.data_import_run import DataImportRun
from app.models.fundamental import FundamentalData
from app.services import fundamentals_candidate_service as candidate_service
REPORT_VERSION = 1
APPROVAL_STATUS = "pending_explicit_approval"
FIELD_KEYS = ("pe_ratio", "revenue_growth", "earnings_surprise")
MIN_SCORE_METRICS = 2
# Materiality is a review aid, never an automatic cutover verdict. Definition
# changes remain visible even when a delta falls inside these bands.
FIELD_TOLERANCES = {
"pe_ratio": {"absolute": 1.0, "relative_pct": 10.0},
"revenue_growth": {"absolute": 2.0, "relative_pct": None},
"earnings_surprise": {"absolute": 2.0, "relative_pct": None},
}
DEFINITION_NOTES = {
"pe_ratio": (
"Legacy provider P/E convention versus latest close divided by "
"SEC-derived TTM diluted EPS."
),
"revenue_growth": (
"Legacy provider growth convention versus SEC-derived TTM revenue YoY."
),
"earnings_surprise": (
"Legacy provider latest surprise versus latest completed Dolt earnings "
"event with actual and estimate."
),
}
def fundamental_score(
pe_ratio: float | None,
revenue_growth: float | None,
earnings_surprise: float | None,
) -> float | None:
"""Match the production fundamental-dimension formula without persistence."""
scores: list[float] = []
if _finite(pe_ratio) and pe_ratio > 0:
scores.append(max(0.0, min(100.0, 100.0 - (pe_ratio - 15.0) * (100.0 / 30.0))))
if _finite(revenue_growth):
scores.append(max(0.0, min(100.0, 50.0 + revenue_growth * 2.5)))
if _finite(earnings_surprise):
scores.append(max(0.0, min(100.0, 50.0 + earnings_surprise * 5.0)))
return sum(scores) / len(scores) if len(scores) >= MIN_SCORE_METRICS else None
async def build_report(
db: AsyncSession,
*,
generated_at: datetime | None = None,
today: date | None = None,
) -> dict[str, Any]:
"""Build a point-in-time parity report from one database session."""
generated_at = generated_at or datetime.now(timezone.utc)
today = today or datetime.now(ZoneInfo("America/New_York")).date()
# A report must not mix rows from before and after a concurrent import
# promotion. The scheduled job provides a fresh session, so establish the
# production snapshot before its first query and have Postgres enforce the
# no-write contract as well. SQLite tests retain their normal transaction.
if db.get_bind().dialect.name == "postgresql":
connection = await db.connection(
execution_options={"isolation_level": "REPEATABLE READ"}
)
await connection.execute(text("SET TRANSACTION READ ONLY"))
candidates = await candidate_service.build_candidates(db, today=today)
ticker_ids = [candidate.ticker_id for candidate in candidates]
legacy_by_ticker = await _legacy_values(db, ticker_ids)
source_runs = await _source_runs(db)
rows: list[dict[str, Any]] = []
for candidate in candidates:
legacy = legacy_by_ticker.get(candidate.ticker_id)
candidate_values = {
"pe_ratio": candidate.pe_ratio,
"revenue_growth": candidate.revenue_growth,
"earnings_surprise": candidate.earnings_surprise,
}
legacy_values = {
"pe_ratio": legacy.pe_ratio if legacy else None,
"revenue_growth": legacy.revenue_growth if legacy else None,
"earnings_surprise": legacy.earnings_surprise if legacy else None,
}
fields = {
key: _field_comparison(key, legacy_values[key], candidate_values[key])
for key in FIELD_KEYS
}
legacy_score = fundamental_score(**legacy_values)
candidate_score = fundamental_score(**candidate_values)
rows.append(
{
"symbol": candidate.symbol,
"cik": candidate.cik,
"legacy_fetched_at": _iso(legacy.fetched_at) if legacy else None,
"price_date": _iso(candidate.price_date),
"fields": fields,
"scores": {
"legacy_fundamental": _round(legacy_score),
"candidate_fundamental": _round(candidate_score),
"fundamental_delta": _delta(legacy_score, candidate_score),
"legacy_fundamental_rank": None,
"candidate_fundamental_rank": None,
"fundamental_rank_change": None,
},
}
)
_attach_ranks(rows, "legacy_fundamental", "legacy_fundamental_rank")
_attach_ranks(rows, "candidate_fundamental", "candidate_fundamental_rank")
for row in rows:
scores = row["scores"]
scores["fundamental_rank_change"] = _rank_change(
scores["legacy_fundamental_rank"], scores["candidate_fundamental_rank"]
)
return {
"report_version": REPORT_VERSION,
"generated_at": generated_at.isoformat(),
"as_of_date": today.isoformat(),
"approval_status": APPROVAL_STATUS,
"read_only": True,
"fundamental_score_formula": (
"Equal-weighted mean of 2+ available sub-scores: P/E = "
"clamp(100-(pe-15)*(100/30)); revenue growth = "
"clamp(50+growth*2.5); earnings surprise = "
"clamp(50+surprise*5)."
),
"source_runs": source_runs,
"definition_notes": DEFINITION_NOTES,
"materiality_notes": {
"fields": FIELD_TOLERANCES,
"fundamental_score_absolute": 5.0,
"automatic_cutover": False,
},
"summary": _summary(rows),
"rows": rows,
}
def store_report(report: dict[str, Any], report_dir: str | Path) -> dict[str, str]:
"""Atomically archive JSON/CSV artifacts and update the latest manifest."""
directory = Path(report_dir).expanduser().resolve()
directory.mkdir(parents=True, exist_ok=True)
stamp = _artifact_stamp(report["generated_at"])
json_name = f"fundamentals-parity-{stamp}.json"
csv_name = f"fundamentals-parity-{stamp}.csv"
json_path = directory / json_name
csv_path = directory / csv_name
_atomic_write(json_path, json.dumps(report, indent=2, sort_keys=True) + "\n")
_atomic_write(csv_path, report_csv(report))
manifest = {
"generated_at": report["generated_at"],
"json_file": json_name,
"csv_file": csv_name,
}
_atomic_write(
directory / "latest.json",
json.dumps(manifest, indent=2, sort_keys=True) + "\n",
)
return {
"json": str(json_path),
"csv": str(csv_path),
"manifest": str(directory / "latest.json"),
}
async def generate_and_store(
db: AsyncSession,
report_dir: str | Path,
*,
generated_at: datetime | None = None,
today: date | None = None,
) -> tuple[dict[str, Any], dict[str, str]]:
report = await build_report(db, generated_at=generated_at, today=today)
return report, store_report(report, report_dir)
def load_latest(report_dir: str | Path) -> dict[str, Any] | None:
manifest = _load_manifest(report_dir)
if manifest is None:
return None
try:
path = _manifest_artifact(report_dir, manifest, "json_file")
loaded = json.loads(path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError, TypeError, ValueError):
return None
return loaded if isinstance(loaded, dict) else None
def load_latest_csv(report_dir: str | Path) -> tuple[str, str] | None:
return _load_latest_text_artifact(report_dir, "csv_file")
def load_latest_json(report_dir: str | Path) -> tuple[str, str] | None:
return _load_latest_text_artifact(report_dir, "json_file")
def _load_latest_text_artifact(
report_dir: str | Path, manifest_key: str
) -> tuple[str, str] | None:
manifest = _load_manifest(report_dir)
if manifest is None:
return None
try:
path = _manifest_artifact(report_dir, manifest, manifest_key)
return path.name, path.read_text(encoding="utf-8")
except (OSError, TypeError, ValueError):
return None
def report_csv(report: dict[str, Any]) -> str:
output = io.StringIO(newline="")
columns = [
"symbol",
"cik",
"legacy_fetched_at",
"price_date",
*(
f"{field}_{suffix}"
for field in FIELD_KEYS
for suffix in ("legacy", "candidate", "absolute_delta", "relative_delta_pct", "material")
),
"legacy_fundamental",
"candidate_fundamental",
"fundamental_delta",
"legacy_fundamental_rank",
"candidate_fundamental_rank",
"fundamental_rank_change",
]
writer = csv.DictWriter(output, fieldnames=columns)
writer.writeheader()
for row in report.get("rows", []):
flat = {
"symbol": row["symbol"],
"cik": row.get("cik"),
"legacy_fetched_at": row.get("legacy_fetched_at"),
"price_date": row.get("price_date"),
**row["scores"],
}
for field in FIELD_KEYS:
comparison = row["fields"][field]
for suffix in (
"legacy",
"candidate",
"absolute_delta",
"relative_delta_pct",
"material",
):
flat[f"{field}_{suffix}"] = comparison.get(suffix)
writer.writerow(flat)
return output.getvalue()
async def _legacy_values(
db: AsyncSession, ticker_ids: list[int]
) -> dict[int, FundamentalData]:
if not ticker_ids:
return {}
rows = (
await db.execute(
select(FundamentalData).where(FundamentalData.ticker_id.in_(ticker_ids))
)
).scalars()
return {row.ticker_id: row for row in rows}
async def _source_runs(db: AsyncSession) -> dict[str, dict[str, Any] | None]:
sources = ("sec_facts", "dolt_earnings")
rows = (
await db.execute(
select(DataImportRun)
.where(
DataImportRun.source.in_(sources),
DataImportRun.status.in_(("promoted", "no_op")),
)
.order_by(DataImportRun.id.desc())
)
).scalars()
latest: dict[str, dict[str, Any] | None] = {source: None for source in sources}
for row in rows:
if latest[row.source] is None:
latest[row.source] = {
"run_id": row.id,
"status": row.status,
"revision": row.revision,
"source_max_date": _iso(row.source_max_date),
"completed_at": _iso(row.completed_at),
}
return latest
def _field_comparison(
key: str, legacy: float | None, candidate: float | None
) -> dict[str, Any]:
legacy = float(legacy) if _finite(legacy) else None
candidate = float(candidate) if _finite(candidate) else None
absolute = _delta(legacy, candidate)
relative = (
None
if absolute is None or legacy in (None, 0)
else round(absolute / abs(legacy) * 100.0, 4)
)
tolerance = FIELD_TOLERANCES[key]
material = False
if absolute is not None:
material = abs(absolute) > tolerance["absolute"]
relative_limit = tolerance["relative_pct"]
if relative_limit is not None:
material = material and relative is not None and abs(relative) > relative_limit
return {
"legacy": _round(legacy),
"candidate": _round(candidate),
"absolute_delta": absolute,
"relative_delta_pct": relative,
"material": material,
"definition_changed": True,
}
def _attach_ranks(rows: list[dict[str, Any]], value_key: str, rank_key: str) -> None:
values = [
row["scores"][value_key]
for row in rows
if _finite(row["scores"][value_key])
]
for row in rows:
value = row["scores"][value_key]
row["scores"][rank_key] = (
1 + sum(other > value for other in values) if _finite(value) else None
)
def _summary(rows: list[dict[str, Any]]) -> dict[str, Any]:
field_stats = {}
for key in FIELD_KEYS:
comparisons = [row["fields"][key] for row in rows]
deltas = [
abs(item["absolute_delta"])
for item in comparisons
if item["absolute_delta"] is not None
]
field_stats[key] = {
"legacy_available": sum(item["legacy"] is not None for item in comparisons),
"candidate_available": sum(
item["candidate"] is not None for item in comparisons
),
"both_available": len(deltas),
"material_differences": sum(item["material"] for item in comparisons),
"median_absolute_delta": _round(statistics.median(deltas) if deltas else None),
"p95_absolute_delta": _round(_percentile(deltas, 0.95)),
"max_absolute_delta": _round(max(deltas) if deltas else None),
}
fundamental_deltas = _score_deltas(rows, "fundamental_delta")
changed_rows = sorted(
(
{
"symbol": row["symbol"],
"fundamental_delta": row["scores"]["fundamental_delta"],
"fundamental_rank_change": row["scores"]["fundamental_rank_change"],
}
for row in rows
if row["scores"]["fundamental_delta"] is not None
),
key=lambda item: (
abs(item["fundamental_delta"] or 0),
),
reverse=True,
)[:20]
return {
"universe_count": len(rows),
"legacy_fundamental_score_available": _count_score(
rows, "legacy_fundamental"
),
"candidate_fundamental_score_available": _count_score(
rows, "candidate_fundamental"
),
"fundamental_scores_compared": len(fundamental_deltas),
"fundamental_score_material_changes": sum(
abs(delta) > 5.0 for delta in fundamental_deltas
),
"fundamental_rank_changes": _rank_change_count(
rows, "fundamental_rank_change"
),
"field_stats": field_stats,
"largest_changes": changed_rows,
}
def _score_deltas(rows: Iterable[dict[str, Any]], key: str) -> list[float]:
return [
row["scores"][key]
for row in rows
if row["scores"][key] is not None
]
def _count_score(rows: Iterable[dict[str, Any]], key: str) -> int:
return sum(row["scores"][key] is not None for row in rows)
def _rank_change_count(rows: Iterable[dict[str, Any]], key: str) -> int:
return sum(
row["scores"][key] not in (None, 0)
for row in rows
)
def _rank_change(legacy: int | None, candidate: int | None) -> int | None:
# Positive means the candidate improved its rank.
return legacy - candidate if legacy is not None and candidate is not None else None
def _delta(legacy: float | None, candidate: float | None) -> float | None:
if not _finite(legacy) or not _finite(candidate):
return None
return round(candidate - legacy, 4)
def _round(value: float | None, digits: int = 4) -> float | None:
return round(float(value), digits) if _finite(value) else None
def _percentile(values: list[float], quantile: float) -> float | None:
if not values:
return None
ordered = sorted(values)
index = max(0, math.ceil(quantile * len(ordered)) - 1)
return ordered[index]
def _finite(value: Any) -> bool:
return (
isinstance(value, (int, float))
and not isinstance(value, bool)
and math.isfinite(value)
)
def _iso(value: Any) -> str | None:
return value.isoformat() if value is not None else None
def _artifact_stamp(raw: str) -> str:
parsed = datetime.fromisoformat(raw.replace("Z", "+00:00"))
return parsed.astimezone(timezone.utc).strftime("%Y%m%dT%H%M%S%fZ")
def _atomic_write(path: Path, content: str) -> None:
temp = path.with_name(f".{path.name}.{os.getpid()}.tmp")
temp.write_text(content, encoding="utf-8", newline="")
os.replace(temp, path)
def _load_manifest(report_dir: str | Path) -> dict[str, Any] | None:
path = Path(report_dir).expanduser().resolve() / "latest.json"
try:
loaded = json.loads(path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError, TypeError, ValueError):
return None
return loaded if isinstance(loaded, dict) else None
def _manifest_artifact(
report_dir: str | Path, manifest: dict[str, Any], key: str
) -> Path:
directory = Path(report_dir).expanduser().resolve()
name = Path(str(manifest.get(key, ""))).name
if not name:
raise ValueError(f"Latest parity manifest has no {key}")
return directory / name
@@ -12,7 +12,6 @@ from app.models.data_import_run import DataImportRun
from app.models.fundamental_snapshot import FundamentalSnapshot
from app.models.sec_filing_gap import SecFilingGap
from app.models.ticker import Ticker
from app.services import fundamental_data_refresh_service
_SEC_FORMS = ("10-K", "10-Q", "10-K/A", "10-Q/A")
@@ -78,8 +77,6 @@ async def blocked_reasons_by_cik(
ciks: set[str] | None = None,
) -> dict[str, str]:
"""Current SEC blocker code by CIK; no historical audit scan."""
if not await fundamental_data_refresh_service.is_enabled(db):
return {}
if ciks is not None and not ciks:
return {}
+2 -2
View File
@@ -497,8 +497,8 @@ async def _compute_fundamental_score(
"reason": "Earnings surprise data not available",
})
# Require at least two real metrics — a single available metric (e.g. only
# market cap is free on FMP) does not make a meaningful fundamental score.
# Require at least two real metrics — a single available metric (e.g. an
# issuer with only a market cap) does not make a meaningful fundamental score.
MIN_METRICS = 2
if len(scores) < MIN_METRICS:
unavailable.append({
+1 -1
View File
@@ -39,7 +39,7 @@ logger = logging.getLogger(__name__)
_WWW = "https://www.sec.gov"
_DATA = "https://data.sec.gov"
# Resolve CA bundle for explicit httpx verify (matches app/providers/fmp.py).
# Resolve CA bundle for explicit httpx verify (matches app/providers/alpaca.py).
_CA = os.environ.get("SSL_CERT_FILE", "")
_CA_VERIFY: str | bool = _CA if _CA and Path(_CA).exists() else True
+8 -124
View File
@@ -113,116 +113,6 @@ def _normalise_symbols(symbols: Iterable[str]) -> list[str]:
return sorted(deduped)
def _extract_symbols_from_fmp_payload(payload: object) -> list[str]:
if not isinstance(payload, list):
return []
symbols: list[str] = []
for item in payload:
if not isinstance(item, dict):
continue
candidate = item.get("symbol") or item.get("ticker")
if isinstance(candidate, str):
symbols.append(candidate)
return symbols
async def _try_fmp_urls(
client: httpx.AsyncClient,
urls: list[str],
) -> tuple[list[str], list[str]]:
failures: list[str] = []
for url in urls:
endpoint = url.split("?")[0]
try:
response = await client.get(url)
except httpx.HTTPError as exc:
failures.append(f"{endpoint}: network error ({type(exc).__name__}: {exc})")
continue
if response.status_code != 200:
failures.append(f"{endpoint}: HTTP {response.status_code}")
continue
try:
payload = response.json()
except ValueError:
failures.append(f"{endpoint}: invalid JSON payload")
continue
symbols = _extract_symbols_from_fmp_payload(payload)
if symbols:
return symbols, failures
failures.append(f"{endpoint}: empty/unsupported payload")
return [], failures
async def _fetch_universe_symbols_from_fmp(universe: str) -> list[str]:
if not settings.fmp_api_key:
raise ValidationError(
"FMP API key is required for universe bootstrap (set FMP_API_KEY)"
)
api_key = settings.fmp_api_key
stable_base = "https://financialmodelingprep.com/stable"
legacy_base = "https://financialmodelingprep.com/api/v3"
stable_candidates: dict[str, list[str]] = {
"sp500": [
f"{stable_base}/sp500-constituent?apikey={api_key}",
f"{stable_base}/sp500-constituents?apikey={api_key}",
],
"nasdaq100": [
f"{stable_base}/nasdaq-100-constituent?apikey={api_key}",
f"{stable_base}/nasdaq100-constituent?apikey={api_key}",
f"{stable_base}/nasdaq-100-constituents?apikey={api_key}",
],
"nasdaq_all": [
f"{stable_base}/stock-screener?exchange=NASDAQ&isEtf=false&limit=10000&apikey={api_key}",
f"{stable_base}/available-traded/list?apikey={api_key}",
],
}
legacy_candidates: dict[str, list[str]] = {
"sp500": [
f"{legacy_base}/sp500_constituent?apikey={api_key}",
f"{legacy_base}/sp500_constituent",
],
"nasdaq100": [
f"{legacy_base}/nasdaq_constituent?apikey={api_key}",
f"{legacy_base}/nasdaq_constituent",
],
"nasdaq_all": [
f"{legacy_base}/stock-screener?exchange=NASDAQ&isEtf=false&limit=10000&apikey={api_key}",
],
}
failures: list[str] = []
async with httpx.AsyncClient(timeout=30.0, verify=_CA_BUNDLE_PATH) as client:
stable_symbols, stable_failures = await _try_fmp_urls(client, stable_candidates[universe])
failures.extend(stable_failures)
if stable_symbols:
return stable_symbols
legacy_symbols, legacy_failures = await _try_fmp_urls(client, legacy_candidates[universe])
failures.extend(legacy_failures)
if legacy_symbols:
return legacy_symbols
if failures:
reason = "; ".join(failures[:6])
logger.warning("FMP universe fetch failed for %s: %s", universe, reason)
raise ProviderError(
f"Failed to fetch universe symbols from FMP for '{universe}'. Attempts: {reason}"
)
raise ProviderError(f"Failed to fetch universe symbols from FMP for '{universe}'")
async def _fetch_wiki_constituent_symbols(
client: httpx.AsyncClient,
url: str,
@@ -351,13 +241,16 @@ async def fetch_universe_symbols(
Fallback order:
1) Free public sources (Wikipedia/NASDAQ trader)
2) FMP endpoints (if available)
3) Cached snapshot in SystemSetting
4) Built-in seed symbols
2) Cached snapshot in SystemSetting
3) Built-in seed symbols
Returns ``(symbols, source_label)`` so bootstrap UI can show where the
list came from (important when Wikipedia/FMP fail and a stale cache still
lists BK instead of BNY).
list came from (important when the public source fails and a stale cache
still lists BK instead of BNY).
The seeds are representative, not complete, so a *fresh* install whose
public source is down bootstraps a partial universe. A warm instance is
unaffected it falls through to its cached snapshot.
"""
normalised_universe = _validate_universe(universe)
failures: list[str] = []
@@ -369,15 +262,6 @@ async def fetch_universe_symbols(
await _write_cached_symbols(db, normalised_universe, cleaned_public, public_source or "public")
return cleaned_public, public_source or "public"
try:
fmp_symbols = await _fetch_universe_symbols_from_fmp(normalised_universe)
cleaned_fmp = _normalise_symbols(fmp_symbols)
if cleaned_fmp:
await _write_cached_symbols(db, normalised_universe, cleaned_fmp, "fmp")
return cleaned_fmp, "fmp"
except (ProviderError, ValidationError) as exc:
failures.append(str(exc))
cached_symbols = await _read_cached_symbols(db, normalised_universe)
if cached_symbols:
logger.warning(