diff --git a/.env.example b/.env.example index f74c548..31cb036 100644 --- a/.env.example +++ b/.env.example @@ -52,6 +52,10 @@ SEC_REQUEST_SPACING_SECONDS=0.2 SEC_MAX_RETRIES=4 SEC_REQUEST_TIMEOUT_SECONDS=30.0 +# A5 read-only parity report archive. In production keep this outside the +# rsync deployment tree, e.g. /var/lib/signal-platform/reports/fundamentals-parity. +FUNDAMENTALS_PARITY_REPORT_DIR=reports/fundamentals-parity + # Regime Monitor — FRED (VIX + HY credit spreads). Free key: https://fred.stlouisfed.org/docs/api/api_key.html # Optional: without it the volatility (V1) and credit (C1) pillars show as n/a. FRED_API_KEY= diff --git a/.gitignore b/.gitignore index ca228bc..9c49cc9 100644 --- a/.gitignore +++ b/.gitignore @@ -51,3 +51,6 @@ backtest_snapshots/ reports/*.pkl reports/*.pk1 reports/.cache/ +# Runtime A5 parity bundles are generated on the production server. Research +# conclusions belong in docs/research, not as an ever-growing artifact archive. +reports/fundamentals-parity/ diff --git a/app/config.py b/app/config.py index d07d69f..e43ef63 100644 --- a/app/config.py +++ b/app/config.py @@ -61,6 +61,10 @@ 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 = "" diff --git a/app/routers/admin.py b/app/routers/admin.py index 73582f9..8ff9ad7 100644 --- a/app/routers/admin.py +++ b/app/routers/admin.py @@ -453,6 +453,36 @@ 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) # --------------------------------------------------------------------------- diff --git a/app/scheduler.py b/app/scheduler.py index 76b6df5..24db4d2 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -40,6 +40,7 @@ from app.services import ( sentiment_service, settings_store, shadow_book_service, + fundamentals_parity_service, ) from app.services.data_import import STATUS_FAILED, SourceImporter, run_import from app.services.dolt_earnings_importer import DoltEarningsImporter @@ -98,6 +99,7 @@ _JOB_NAMES = [ "fundamental_collector", "dolt_earnings_import", "sec_fundamentals_import", + "fundamentals_parity_report", "rr_scanner", "ticker_universe_sync", "alerts", @@ -981,6 +983,49 @@ async def run_sec_fundamentals_import() -> None: await _run_shadow_import("sec_fundamentals_import", SecFundamentalsImporter()) +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 # --------------------------------------------------------------------------- @@ -1524,6 +1569,7 @@ SCHEDULE_DEFAULTS: dict[str, str] = { # Shadow source imports. They never write legacy fundamental_data before A5. "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). @@ -1539,6 +1585,7 @@ _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", @@ -1645,6 +1692,17 @@ def configure_scheduler(schedule_config: dict[str, str] | None = None) -> None: name="SEC Fundamentals Import (shadow)", 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( @@ -1720,6 +1778,9 @@ 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], diff --git a/app/schemas/admin.py b/app/schemas/admin.py index 4c46639..3127750 100644 --- a/app/schemas/admin.py +++ b/app/schemas/admin.py @@ -78,6 +78,9 @@ class ScheduleConfigUpdate(BaseModel): (min hour dom month dow); timezone is an IANA name (e.g. America/New_York).""" schedule_timezone: str | None = Field(default=None, max_length=64) 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) diff --git a/app/schemas/fundamental.py b/app/schemas/fundamental.py index b28158c..b4d97f5 100644 --- a/app/schemas/fundamental.py +++ b/app/schemas/fundamental.py @@ -26,6 +26,7 @@ class MetricItem(BaseModel): industry: MetricIndustry | None = None period_end: str | None = None filed_date: str | None = None + caveat: str | None = None source: str = "sec" diff --git a/app/services/admin_service.py b/app/services/admin_service.py index 2cf3f59..4cf2424 100644 --- a/app/services/admin_service.py +++ b/app/services/admin_service.py @@ -614,6 +614,7 @@ VALID_JOB_NAMES = { "fundamental_collector", "dolt_earnings_import", "sec_fundamentals_import", + "fundamentals_parity_report", "rr_scanner", "ticker_universe_sync", "outcome_evaluator", @@ -637,6 +638,7 @@ JOB_LABELS = { "fundamental_collector": "Fundamental Collector", "dolt_earnings_import": "Dolt Earnings Import (shadow)", "sec_fundamentals_import": "SEC Fundamentals Import (shadow)", + "fundamentals_parity_report": "Fundamentals Parity Report (read-only)", "rr_scanner": "R:R Scanner", "ticker_universe_sync": "Ticker Universe Sync", "outcome_evaluator": "Outcome Evaluator", @@ -775,3 +777,30 @@ 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) diff --git a/app/services/fundamentals_api_service.py b/app/services/fundamentals_api_service.py index 3620bd9..33541bb 100644 --- a/app/services/fundamentals_api_service.py +++ b/app/services/fundamentals_api_service.py @@ -121,6 +121,7 @@ def _build_metrics(derived, peer_derived, two: str | None) -> list[dict[str, Any "industry": industry, "period_end": _iso(series.period_end) if series else None, "filed_date": _iso(series.filed_date) if series else None, + "caveat": series.caveat if series else None, "source": "sec", }) return out @@ -303,7 +304,8 @@ async def _latest_close(db, ticker_id: int) -> tuple[float, date] | None: def _empty_metrics() -> list[dict[str, Any]]: return [{"key": k, "value": None, "history": [], "industry": None, - "period_end": None, "filed_date": None, "source": "sec"} for k in METRIC_KEYS] + "period_end": None, "filed_date": None, "caveat": None, + "source": "sec"} for k in METRIC_KEYS] def _empty_earnings() -> dict[str, Any]: diff --git a/app/services/fundamentals_derivation.py b/app/services/fundamentals_derivation.py index e5d203c..d9790ce 100644 --- a/app/services/fundamentals_derivation.py +++ b/app/services/fundamentals_derivation.py @@ -27,6 +27,11 @@ _FP_TO_Q = {"Q1": 1, "Q2": 2, "Q3": 3, "FY": 4} _Q_TO_FP = {1: "Q1", 2: "Q2", 3: "Q3", 4: "FY"} _PREV_FP = {"Q2": "Q1", "Q3": "Q2", "FY": "Q3"} TAPE_LEN = 4 # quarter-tape length +SPLIT_SUSPECT_SHARE_CHANGE_PCT = 25.0 +SPLIT_SENSITIVE_CAVEAT = ( + "Not comparable: share count changed at least 25%; possible split or " + "corporate action." +) # Duration (flow) fields differenced from YTD into discrete quarters + summed to TTM. _FLOW_FIELDS = ( @@ -47,6 +52,7 @@ class MetricSeries: history: list[MetricPoint] = field(default_factory=list) # oldest -> newest, <= TAPE_LEN period_end: date | None = None filed_date: date | None = None + caveat: str | None = None @dataclass @@ -96,6 +102,7 @@ def derive(snapshots: Iterable[Any]) -> DerivedFundamentals: "net_debt_to_ebitda": _leverage_series(selected, discrete, tape), "share_count_change_yoy": _share_change_series(selected, tape), } + _guard_split_sensitive_metrics(result.metrics) for series in result.metrics.values(): series.period_end = latest_row.period_end series.filed_date = latest_row.filed_date @@ -248,6 +255,40 @@ def _share_change_series(selected, tape) -> MetricSeries: return _series(pts) +def _guard_split_sensitive_metrics(metrics: dict[str, MetricSeries]) -> None: + """Suppress historical comparisons likely distorted by a corporate action. + + Company Facts has no point-in-time split factors. A large YoY share-count + move can therefore make both the point-in-time share comparison and + per-share EPS growth non-comparable. Keep the raw facts in snapshots, but + expose nulls plus an explicit caveat in the user-facing derived series. + """ + shares = metrics.get("share_count_change_yoy") + eps = metrics.get("eps_growth_yoy") + if shares is None or eps is None: + return + + suspect_periods = { + point.period_end + for point in shares.history + if point.value is not None + and abs(point.value) >= SPLIT_SUSPECT_SHARE_CHANGE_PCT + } + if not suspect_periods: + return + + for series in (shares, eps): + latest_guarded = bool( + series.history and series.history[-1].period_end in suspect_periods + ) + for point in series.history: + if point.period_end in suspect_periods: + point.value = None + series.value = series.history[-1].value if series.history else None + if latest_guarded: + series.caveat = SPLIT_SENSITIVE_CAVEAT + + def _net_debt(row: Any) -> float | None: if row is None: return None diff --git a/app/services/fundamentals_parity_service.py b/app/services/fundamentals_parity_service.py new file mode 100644 index 0000000..8b99f42 --- /dev/null +++ b/app/services/fundamentals_parity_service.py @@ -0,0 +1,598 @@ +"""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 collections import defaultdict +from datetime import date, datetime, timezone +from pathlib import Path +from typing import Any, Iterable +from zoneinfo import ZoneInfo + +from sqlalchemy import func, select, text +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.data_import_run import DataImportRun +from app.models.earnings_event import EarningsEvent +from app.models.fundamental import FundamentalData +from app.models.fundamental_snapshot import FundamentalSnapshot +from app.models.ohlcv import OHLCVRecord +from app.models.ticker import Ticker +from app.services import fundamentals_derivation as deriv + +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")) + + tickers = list((await db.execute(select(Ticker).order_by(Ticker.symbol))).scalars()) + ticker_ids = [ticker.id for ticker in tickers] + ciks = sorted({ticker.cik for ticker in tickers if ticker.cik}) + + legacy_by_ticker = await _legacy_values(db, ticker_ids) + derived_by_cik = await _derived_by_cik(db, ciks) + closes_by_ticker = await _latest_closes(db, ticker_ids) + surprise_by_ticker = await _latest_surprises(db, ticker_ids, today) + source_runs = await _source_runs(db) + + rows: list[dict[str, Any]] = [] + for ticker in tickers: + legacy = legacy_by_ticker.get(ticker.id) + derived = derived_by_cik.get(ticker.cik) if ticker.cik else None + close = closes_by_ticker.get(ticker.id) + candidate_pe = ( + _pe(close[0], derived.ttm_diluted_eps) + if close is not None and derived is not None + else None + ) + growth_series = ( + derived.metrics.get("revenue_growth_yoy") if derived is not None else None + ) + candidate = { + "pe_ratio": candidate_pe, + "revenue_growth": growth_series.value if growth_series else None, + "earnings_surprise": surprise_by_ticker.get(ticker.id), + } + 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[key]) + for key in FIELD_KEYS + } + legacy_score = fundamental_score(**legacy_values) + candidate_score = fundamental_score(**candidate) + rows.append( + { + "symbol": ticker.symbol, + "cik": ticker.cik, + "legacy_fetched_at": _iso(legacy.fetched_at) if legacy else None, + "price_date": _iso(close[1]) if close else None, + "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 _derived_by_cik( + db: AsyncSession, ciks: list[str] +) -> dict[str, deriv.DerivedFundamentals]: + if not ciks: + return {} + grouped: dict[str, list[FundamentalSnapshot]] = defaultdict(list) + rows = ( + await db.execute( + select(FundamentalSnapshot).where(FundamentalSnapshot.cik.in_(ciks)) + ) + ).scalars() + for row in rows: + grouped[row.cik].append(row) + return {cik: deriv.derive(grouped.get(cik, [])) for cik in ciks} + + +async def _latest_closes( + db: AsyncSession, ticker_ids: list[int] +) -> dict[int, tuple[float, date]]: + if not ticker_ids: + return {} + latest = ( + select(OHLCVRecord.ticker_id, func.max(OHLCVRecord.date).label("max_date")) + .where(OHLCVRecord.ticker_id.in_(ticker_ids)) + .group_by(OHLCVRecord.ticker_id) + .subquery() + ) + rows = ( + await db.execute( + select(OHLCVRecord.ticker_id, OHLCVRecord.close, OHLCVRecord.date).join( + latest, + (OHLCVRecord.ticker_id == latest.c.ticker_id) + & (OHLCVRecord.date == latest.c.max_date), + ) + ) + ).all() + return { + ticker_id: (float(close), close_date) + for ticker_id, close, close_date in rows + if _finite(close) + } + + +async def _latest_surprises( + db: AsyncSession, ticker_ids: list[int], today: date +) -> dict[int, float]: + if not ticker_ids: + return {} + rows = ( + await db.execute( + select(EarningsEvent) + .where( + EarningsEvent.ticker_id.in_(ticker_ids), + EarningsEvent.announce_date < today, + ) + .order_by(EarningsEvent.ticker_id, EarningsEvent.announce_date.desc()) + ) + ).scalars() + out: dict[int, float] = {} + for row in rows: + if row.ticker_id in out: + continue + surprise = _surprise(row.eps_estimate, row.eps_actual) + if surprise is not None: + out[row.ticker_id] = surprise + return out + + +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 _surprise(estimate: float | None, actual: float | None) -> float | None: + if not _finite(estimate) or not _finite(actual) or estimate == 0: + return None + return (actual - estimate) / abs(estimate) * 100.0 + + +def _pe(price: float | None, ttm_eps: float | None) -> float | None: + if not _finite(price) or price <= 0 or not _finite(ttm_eps) or ttm_eps <= 0: + return None + return price / ttm_eps + + +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 diff --git a/deploy/provision_fundamentals.sh b/deploy/provision_fundamentals.sh index dcd2ee1..f80d6c6 100755 --- a/deploy/provision_fundamentals.sh +++ b/deploy/provision_fundamentals.sh @@ -13,6 +13,9 @@ APP_GROUP="${APP_GROUP:-deploy}" ENV_FILE="${ENV_FILE:-/opt/signalplatform/.env}" MIN_FREE_GB="${DOLT_MIN_FREE_DISK_GB:-5}" EARNINGS_DIR="${DOLT_DATA_DIR}/${DOLT_EARNINGS_SUBDIR}" +DOLT_IDENTITY_NAME="${DOLT_IDENTITY_NAME:-Signal Platform}" +DOLT_IDENTITY_EMAIL="${DOLT_IDENTITY_EMAIL:-signal-platform@localhost}" +FUNDAMENTALS_PARITY_REPORT_DIR="${FUNDAMENTALS_PARITY_REPORT_DIR:-/var/lib/signal-platform/reports/fundamentals-parity}" fail() { echo "ERROR: $*" >&2 @@ -25,6 +28,38 @@ version_ok() { grep -Eq "(^|[[:space:]])v?${DOLT_VERSION}([[:space:]]|$)" <<<"$output" } +as_app_user() { + if [[ "$(id -un)" == "$APP_USER" ]]; then + "$@" + else + command -v runuser >/dev/null 2>&1 || fail "runuser is required" + runuser -u "$APP_USER" -- "$@" + fi +} + +repo_command() { + ( + cd "$EARNINGS_DIR" + as_app_user "$@" + ) +} + +repo_config_value() { + repo_command "$DOLT_BINARY" config --get "$1" +} + +configure_identity() { + local name email + name="$(repo_config_value user.name 2>/dev/null || true)" + email="$(repo_config_value user.email 2>/dev/null || true)" + if [[ -z "$name" ]]; then + repo_command "$DOLT_BINARY" config --local --add user.name "$DOLT_IDENTITY_NAME" + fi + if [[ -z "$email" ]]; then + repo_command "$DOLT_BINARY" config --local --add user.email "$DOLT_IDENTITY_EMAIL" + fi +} + check_free_space() { local available_kb available_kb="$(df -Pk "$DOLT_DATA_DIR" | awk 'NR == 2 {print $4}')" @@ -45,9 +80,12 @@ check_env() { || fail "set DOLT_EARNINGS_SUBDIR=$DOLT_EARNINGS_SUBDIR in $ENV_FILE" grep -Eq '^SEC_USER_AGENT=.*@.*' "$ENV_FILE" \ || fail "SEC_USER_AGENT in $ENV_FILE must contain a real contact email" + grep -Fqx "FUNDAMENTALS_PARITY_REPORT_DIR=$FUNDAMENTALS_PARITY_REPORT_DIR" "$ENV_FILE" \ + || fail "set FUNDAMENTALS_PARITY_REPORT_DIR=$FUNDAMENTALS_PARITY_REPORT_DIR in $ENV_FILE" } check_all() { + local identity_name identity_email id "$APP_USER" >/dev/null 2>&1 || fail "missing service user: $APP_USER" [[ -x "$DOLT_BINARY" ]] || fail "missing Dolt binary: $DOLT_BINARY" version_ok || fail "expected Dolt $DOLT_VERSION at $DOLT_BINARY" @@ -62,6 +100,19 @@ check_all() { else fail "run --check as $APP_USER (or install runuser)" fi + identity_name="$(repo_config_value user.name 2>/dev/null || true)" + identity_email="$(repo_config_value user.email 2>/dev/null || true)" + [[ -n "$identity_name" ]] || fail "missing Dolt user.name for $EARNINGS_DIR" + [[ -n "$identity_email" ]] || fail "missing Dolt user.email for $EARNINGS_DIR" + [[ -d "$FUNDAMENTALS_PARITY_REPORT_DIR" ]] \ + || fail "missing parity report directory: $FUNDAMENTALS_PARITY_REPORT_DIR" + if [[ "$(id -un)" == "$APP_USER" ]]; then + [[ -w "$FUNDAMENTALS_PARITY_REPORT_DIR" ]] \ + || fail "parity report directory is not writable by $APP_USER" + else + runuser -u "$APP_USER" -- test -w "$FUNDAMENTALS_PARITY_REPORT_DIR" \ + || fail "parity report directory is not writable by $APP_USER" + fi check_free_space check_env echo "OK: Dolt $DOLT_VERSION and earnings clone are provisioned" @@ -88,6 +139,7 @@ fi version_ok || fail "Dolt $DOLT_VERSION installation failed" install -d -o "$APP_USER" -g "$APP_GROUP" -m 0750 "$DOLT_DATA_DIR" +install -d -o "$APP_USER" -g "$APP_GROUP" -m 0750 "$FUNDAMENTALS_PARITY_REPORT_DIR" check_free_space if [[ ! -d "$EARNINGS_DIR/.dolt" ]]; then @@ -97,4 +149,5 @@ if [[ ! -d "$EARNINGS_DIR/.dolt" ]]; then "$DOLT_BINARY" clone post-no-preference/earnings "$EARNINGS_DIR" fi +configure_identity check_all diff --git a/docs/dolt-integration-plan.md b/docs/dolt-integration-plan.md index 71d6b2a..fa22ac5 100644 --- a/docs/dolt-integration-plan.md +++ b/docs/dolt-integration-plan.md @@ -426,6 +426,8 @@ workstream B — Alpaca remains the price source throughout. API values across the tracked universe, report per-field deltas and resulting fundamental-score/ranking changes, require explicit approval. Definition changes (e.g. TTM vs provider convention) called out, not averaged away. + The read-only report job and Admin summary/download are implemented; production + observation and explicit cutover approval remain pending. - A6. Remove FMP/Finnhub/Alpha Vantage; keep monitoring + manual fallback. **Workstream B (independent, start when wanted):** diff --git a/docs/fundamentals-deployment.md b/docs/fundamentals-deployment.md index fa67d34..5f25ca7 100644 --- a/docs/fundamentals-deployment.md +++ b/docs/fundamentals-deployment.md @@ -8,6 +8,7 @@ approval. Do not add OS cron entries: the application scheduler owns both jobs. - `Dolt Earnings Import (shadow)` runs daily at 02:30 America/New_York. - `SEC Fundamentals Import (shadow)` runs daily at 04:00 America/New_York. +- `Fundamentals Parity Report (read-only)` runs daily at 05:30 America/New_York. - Both jobs are visible, toggleable, and manually triggerable in Admin → Jobs. - Cron expressions are editable in Admin → Schedule. - Every attempt is recorded in `data_import_runs`; failures also create a system @@ -28,11 +29,14 @@ DOLT_EARNINGS_SUBDIR=earnings DOLT_MIN_FREE_DISK_GB=5.0 SEC_USER_AGENT=signal-platform/1.0 (contact: real-address@example.com) SEC_REQUEST_SPACING_SECONDS=0.2 +FUNDAMENTALS_PARITY_REPORT_DIR=/var/lib/signal-platform/reports/fundamentals-parity ``` Use a real monitored contact address. Keep at least 5 GB free at the Dolt data path; 8–10 GB gives comfortable growth headroom. The data directory must stay outside `/opt/signalplatform`, because deployments use `rsync --delete` there. +The parity-report directory is also persistent and owned by the service user; +its small timestamped JSON/CSV bundles form the temporary A5 review trail. ## One-time provisioning @@ -49,8 +53,18 @@ curl -fsS http://127.0.0.1:8998/api/v1/health The provisioner is idempotent. It installs the pinned Dolt version, creates the persistent directory as `deploy:deploy`, clones -`post-no-preference/earnings`, verifies free space and `.env`, and refuses an -unexpected Dolt version. It does not modify PostgreSQL or start an import. +`post-no-preference/earnings`, configures a repository-local author identity for +`dolt pull`, verifies free space and `.env`, and refuses an unexpected +Dolt version. It does not modify PostgreSQL or start an import. The public clone +does not require `dolt login`. + +For a server provisioned before the author-identity check was added, repair the +existing clone once with: + +```bash +sudo -u deploy -H /usr/local/bin/dolt config --global --add user.name "Signal Platform" +sudo -u deploy -H /usr/local/bin/dolt config --global --add user.email "signal-platform@localhost" +``` Do not replace the pinned version with `latest`. A future Dolt upgrade should be a reviewed change to `DOLT_VERSION`, followed by the same provision/check flow. @@ -69,6 +83,28 @@ In Admin → Jobs, wait until no other job is running, then: 5. Open several ticker pages and confirm the fundamentals panel has populated data and still handles partial/missing issuers cleanly. +## A5 parity observation window + +After both shadow imports are healthy, trigger **Fundamentals Parity Report +(read-only)** once in Admin → Jobs. The **A5 Fundamentals Parity** card above +the jobs shows the latest coverage/delta summary and provides authenticated JSON +and CSV downloads. The canonical server-side bundles are archived at: + +```text +/var/lib/signal-platform/reports/fundamentals-parity/ +``` + +The scheduler then generates one report daily at 05:30 New York time, after the +02:30 Dolt and 04:00 SEC jobs. Review 5–7 consecutive reports before making the +cutover decision. A report never writes `fundamental_data`, dimension/composite +scores, rankings, qualification state, or an approval flag. Materiality bands +only highlight rows for review; A5 still requires explicit approval. + +Each bundle contains legacy and candidate P/E, revenue growth, and earnings +surprise; definition notes; source revisions and price dates; recomputed legacy +and candidate fundamental scores; and per-universe fundamental-rank changes. +Definition changes remain explicit even when numeric deltas are small. + Optional database verification: ```sql @@ -131,3 +167,6 @@ locks. A second Admin trigger should independently report the job as busy. audit rows) must remain covered by the normal production database backup. - Do not proceed to A5 while either shadow feed is unhealthy or the parity gate has not received explicit approval. +- If report generation fails, inspect Admin → System Events and verify + `FUNDAMENTALS_PARITY_REPORT_DIR` exists and is writable by `deploy`. Existing + reports and all live data remain untouched. diff --git a/docs/research/fundamentals-weight-backtest.md b/docs/research/fundamentals-weight-backtest.md new file mode 100644 index 0000000..b2d626f --- /dev/null +++ b/docs/research/fundamentals-weight-backtest.md @@ -0,0 +1,74 @@ +# Fundamentals ranking-overlay research + +Status: completed 2026-07-23. Decision: keep production scoring and qualification unchanged. + +## Question + +Does using point-in-time SEC fundamentals to reorder already-qualified long setups improve the production portfolio's risk-adjusted return? The experiments changed ranking only; qualification, execution, sizing, capacity, costs, ATR exits, and post-stop re-entry remained unchanged. + +## Method + +- The control was the production 80/20 residual-momentum / volatility rank. +- SEC facts became visible only after `accepted_at`, using the newest visible accession per fiscal period. +- Portfolio simulations used daily entry opportunities, close fills, the production gate-reset re-entry policy, and a 30-session horizon. +- Train contained entries before 2024-01-01, validation covered 2024, and test began 2025-01-01. +- Missing composite scores were neutral at 50. +- Deflated Sharpe used the complete registered arm count for each experiment. + +The snapshot contained 511 tracked tickers, 507 unique CIKs, 30,494 SEC snapshot rows, and prices from 2021-06-24 through 2026-07-22. + +## Initial experiment + +The first registered matrix tested quality, growth, and balanced composites at 10%, 20%, 30%, and 40% weights: 13 trials including control. + +No overlay passed the train and validation requirements. The most attractive full-period result, balanced at 10%, failed validation and improved test Sharpe by only 0.06. + +Review also found that filing-time diluted EPS and shares are not reliably comparable across stock splits. A snapshot audit found share-count changes above 25% for 90 of 461 issuers with comparable 2021+ periods, including recognizable split ratios for AMZN, GOOG, NVDA, CMG, and GE plus some obvious unit anomalies. Consequently, EPS growth and share-count change cannot be trusted for historical ranking without point-in-time split factors. + +The complete initial result is recoverable from Git commit `7f944d7`. + +## Split-safe follow-up + +The follow-up excluded diluted-EPS growth and share-count change completely. It tested: + +- Quality: operating margin, FCF margin, and low net-debt/EBITDA, requiring at least two inputs. +- Growth: revenue growth only. +- Balanced: equal quality and growth weights. +- Overlay weights: 5%, 10%, and 15%. + +This produced 10 registered trials including control. Growth coverage among qualified candidates was 88.96%; lack of data was not the limiting factor. + +| Window | Control Sharpe | Revenue-growth 5% | Delta | +|---|---:|---:|---:| +| Train | 1.26 | 1.31 | +0.05 | +| Validation | 2.52 | 2.64 | +0.12 | +| Test | 1.99 | 1.99 | 0.00 | +| Full | 1.86 | 1.87 | +0.01 | + +Revenue growth at 5% mechanically passed the deliberately permissive "not worse" gate, but did not demonstrate an economically meaningful edge: + +- Test CAGR rose from 54.4% to 56.1%, while full-period CAGR fell from 52.1% to 51.5%. +- Full-period trade overlap was 68.53%, so roughly one-third of selections changed for essentially unchanged Sharpe. +- Revenue-growth IC was 0.0006 in train, 0.0053 in test, and 0.0116 full-period with a full-period t-stat of 0.55. +- Growth weights of 10% and 15% deteriorated; quality and balanced composites failed. +- The test window had already been inspected, so this follow-up was sensitivity evidence rather than a fresh out-of-sample result. + +The complete split-safe result is recoverable from Git commit `dba7ea7`. + +## Decision + +- Do not add fundamental weight to production ranking or the automated qualification gate. +- Do not run another historical weight sweep on the same sample; it would add data-mining rather than new evidence. +- Keep fundamentals informational and user-facing in the UI. +- A5 source-parity and cutover work can proceed independently without changing scoring behavior. +- Treat historical EPS growth and share-count change as non-comparable across corporate actions until a split-aware solution or a conservative UI guard exists. + +Revisit automated weighting only with materially better data, such as point-in-time split factors and historical constituent/delisting coverage, followed by genuinely new forward paper evidence. + +## Limitations + +The snapshot uses today's tracked universe rather than historical membership and delisted securities, creating survivorship bias. Absolute CAGR and Sharpe must not be interpreted as unbiased live expectations. The relative comparison is useful, but the observed test window and short number of independent factor windows limit statistical power. + +## Repository cleanup + +The experiment-only scorer, runner, Mac launcher, caches, tests, and expanded report bundles were removed after this decision. They remain recoverable from commits `eae4d34`, `34d6dda`, `7f944d7`, and `dba7ea7`. Production fundamentals derivation and ingestion remain unchanged. diff --git a/frontend/src/api/admin.ts b/frontend/src/api/admin.ts index 20993a4..bf2fb28 100644 --- a/frontend/src/api/admin.ts +++ b/frontend/src/api/admin.ts @@ -233,6 +233,40 @@ export interface TriggerJobResponse { cadence?: BacktestCadence; } +export interface ParityFieldStats { + legacy_available: number; + candidate_available: number; + both_available: number; + material_differences: number; + median_absolute_delta: number | null; + p95_absolute_delta: number | null; + max_absolute_delta: number | null; +} + +export interface FundamentalsParityReport { + report_version: number; + generated_at: string; + as_of_date: string; + approval_status: string; + read_only: boolean; + summary: { + universe_count: number; + legacy_fundamental_score_available: number; + candidate_fundamental_score_available: number; + fundamental_scores_compared: number; + fundamental_score_material_changes: number; + fundamental_rank_changes: number; + field_stats: Record; + }; + source_runs: Record; +} + export type BacktestTargetModel = 'production_gtl' | 'structural_sr'; export type BacktestCadence = 'weekly' | 'daily'; @@ -259,6 +293,24 @@ export function triggerJob( .then((r) => r.data); } +export function getFundamentalsParityReport() { + return apiClient + .get('admin/fundamentals-parity') + .then((r) => r.data); +} + +export function getFundamentalsParityCsv() { + return apiClient + .get<{ filename: string; content: string } | null>('admin/fundamentals-parity/csv') + .then((r) => r.data); +} + +export function getFundamentalsParityJson() { + return apiClient + .get<{ filename: string; content: string } | null>('admin/fundamentals-parity/json') + .then((r) => r.data); +} + // System events (operational warnings / errors) export interface SystemEvent { id: number; diff --git a/frontend/src/components/admin/FundamentalsParityPanel.tsx b/frontend/src/components/admin/FundamentalsParityPanel.tsx new file mode 100644 index 0000000..36be387 --- /dev/null +++ b/frontend/src/components/admin/FundamentalsParityPanel.tsx @@ -0,0 +1,156 @@ +import { useState } from 'react'; +import { + getFundamentalsParityCsv, + getFundamentalsParityJson, +} from '../../api/admin'; +import { useFundamentalsParityReport } from '../../hooks/useAdmin'; +import { SkeletonTable } from '../ui/Skeleton'; + +const FIELD_LABELS: Record = { + pe_ratio: 'P/E', + revenue_growth: 'Revenue growth', + earnings_surprise: 'Earnings surprise', +}; + +function downloadText(filename: string, content: string, type: string) { + const blob = new Blob([content], { type }); + const url = URL.createObjectURL(blob); + const anchor = document.createElement('a'); + anchor.href = url; + anchor.download = filename; + anchor.click(); + URL.revokeObjectURL(url); +} + +export function FundamentalsParityPanel() { + const { data: report, isLoading, isError, error } = useFundamentalsParityReport(); + const [downloading, setDownloading] = useState(false); + + if (isLoading) return ; + if (isError) { + return

{(error as Error).message}

; + } + + if (!report) { + return ( +
+

A5 Fundamentals Parity

+

+ No report yet. Trigger “Fundamentals Parity Report (read-only)” below. +

+
+ ); + } + + const summary = report.summary; + const generated = new Date(report.generated_at).toLocaleString(); + + async function downloadCsv() { + setDownloading(true); + try { + const artifact = await getFundamentalsParityCsv(); + if (artifact) downloadText(artifact.filename, artifact.content, 'text/csv;charset=utf-8'); + } finally { + setDownloading(false); + } + } + + async function downloadJson() { + setDownloading(true); + try { + const artifact = await getFundamentalsParityJson(); + if (artifact) downloadText(artifact.filename, artifact.content, 'application/json;charset=utf-8'); + } finally { + setDownloading(false); + } + } + + return ( +
+
+
+
+

A5 Fundamentals Parity

+ + approval pending + + + read-only + +
+

+ Generated {generated} · as of {report.as_of_date} · {summary.universe_count} tracked tickers +

+
+
+ + +
+
+ +
+ + + + +
+ +
+ + + + + + + + + + + + + {Object.entries(summary.field_stats).map(([key, stats]) => ( + + + + + + + + + ))} + +
FieldLegacyCandidateComparedMaterialMedian |Δ|
{FIELD_LABELS[key] ?? key}{stats.legacy_available}{stats.candidate_available}{stats.both_available}{stats.material_differences} + {stats.median_absolute_delta == null ? 'n/a' : stats.median_absolute_delta.toFixed(2)} +
+
+ +

+ Materiality bands highlight review candidates only. They do not approve a cutover or write fundamentals, + scores, rankings, or qualification state. +

+
+ ); +} + +function Summary({ label, value }: { label: string; value: string | number }) { + return ( +
+
{label}
+
{value}
+
+ ); +} diff --git a/frontend/src/components/admin/ScheduleSettings.tsx b/frontend/src/components/admin/ScheduleSettings.tsx index 0dd2472..ccb44b5 100644 --- a/frontend/src/components/admin/ScheduleSettings.tsx +++ b/frontend/src/components/admin/ScheduleSettings.tsx @@ -8,6 +8,7 @@ const DEFAULTS: ScheduleConfig = { schedule_daily_pipeline_cron: '0 2 * * *', schedule_dolt_earnings_cron: '30 2 * * *', schedule_sec_fundamentals_cron: '0 4 * * *', + schedule_fundamentals_parity_cron: '30 5 * * *', schedule_near_close_pipeline_cron: '30 15 * * mon-fri', schedule_after_close_pipeline_cron: '45 16 * * mon-fri', schedule_intraday_pipeline_cron: '0 10-15 * * mon-fri', @@ -38,6 +39,12 @@ const FIELDS: { key: keyof ScheduleConfig; label: string; hint: string; mono?: b hint: 'Import tracked-universe SEC facts daily at 04:00 ET. Unchanged revisions become no-op runs.', mono: true, }, + { + key: 'schedule_fundamentals_parity_cron', + label: 'Fundamentals parity report', + hint: 'Read-only legacy vs SEC/Dolt comparison daily at 05:30 ET, after the shadow imports.', + mono: true, + }, { key: 'schedule_near_close_pipeline_cron', label: 'Near-close pipeline (scan + alert)', diff --git a/frontend/src/components/ticker/FundamentalsPanel.tsx b/frontend/src/components/ticker/FundamentalsPanel.tsx index 712e65d..05390e2 100644 --- a/frontend/src/components/ticker/FundamentalsPanel.tsx +++ b/frontend/src/components/ticker/FundamentalsPanel.tsx @@ -138,7 +138,8 @@ export function FundamentalsPanel({ data }: FundamentalsPanelProps) {
{trendRows.map((r) => ( + metric={metrics[r.key]} read={reads[r.key]} + caveat={metrics[r.key]?.caveat} /> ))}
@@ -188,9 +189,10 @@ function Bullet({ label, value, rail, comparison }: { // ---- operating-trend row (delta vs reference, favorable = right) ------------ -function TrendRow({ label, kind, metric, read }: { +function TrendRow({ label, kind, metric, read, caveat }: { label: string; kind: 'growth' | 'margin' | 'share'; metric: MetricItem | undefined; read: string | null | undefined; + caveat: string | null | undefined; }) { const tone = readTone(read); const value = finiteOrNull(metric?.value); @@ -213,7 +215,9 @@ function TrendRow({ label, kind, metric, read }: { } const delta = value != null && ref != null ? value - ref : null; - const comparison = value == null ? ( + const comparison = caveat ? ( + {caveat} + ) : value == null ? ( n/a ) : delta == null ? ( history n/a diff --git a/frontend/src/dev/harness.tsx b/frontend/src/dev/harness.tsx index 0b60783..4d9010e 100644 --- a/frontend/src/dev/harness.tsx +++ b/frontend/src/dev/harness.tsx @@ -23,11 +23,13 @@ function dateFromToday(days: number): string { } function metric(key: string, value: number | null, hist: (number | null)[], - industry: MetricItem['industry'] = null): MetricItem { + industry: MetricItem['industry'] = null, + caveat: string | null = null): MetricItem { return { key: key as MetricItem['key'], value, history: hist.map((v, i) => h(P[i], v)), - industry, period_end: '2026-03-28', filed_date: '2026-05-01', source: 'sec', + industry, period_end: '2026-03-28', filed_date: '2026-05-01', caveat, + source: 'sec', }; } const ind = (median: number, favorable_percentile: number) => @@ -78,12 +80,24 @@ const partial: FundamentalResponse = { earnings: { next: { date: dateFromToday(0), session: 'unknown', days_until: 0 }, recent: [] }, metrics: [ metric('revenue_growth_yoy', 12, [null, 8, 10, 12], null), - metric('eps_growth_yoy', null, [null, null, null, null], null), + metric( + 'eps_growth_yoy', + null, + [null, null, null, null], + null, + 'Not comparable: share count changed at least 25%; possible split or corporate action.', + ), metric('operating_margin', 25, [24, 24, 25, 25], null), metric('fcf_margin', null, [null, null, null, null], null), metric('net_debt', null, [], null), metric('net_debt_to_ebitda', 1.9, [1.7, 1.8, 1.9, 1.9], null), - metric('share_count_change_yoy', 2.1, [1.8, 2.0, 2.0, 2.1], null), + metric( + 'share_count_change_yoy', + null, + [1.8, 2.0, 2.0, null], + null, + 'Not comparable: share count changed at least 25%; possible split or corporate action.', + ), ], valuation: { pe: 15.2, fcf_yield: null, market_cap_est: 5.4e8, diff --git a/frontend/src/hooks/useAdmin.ts b/frontend/src/hooks/useAdmin.ts index 7d2d6b2..3b5d31a 100644 --- a/frontend/src/hooks/useAdmin.ts +++ b/frontend/src/hooks/useAdmin.ts @@ -316,6 +316,14 @@ export function useJobs() { }); } +export function useFundamentalsParityReport() { + return useQuery({ + queryKey: ['admin', 'fundamentals-parity'], + queryFn: () => adminApi.getFundamentalsParityReport(), + refetchInterval: 15_000, + }); +} + export function usePipelineReadiness() { return useQuery({ queryKey: ['admin', 'pipeline-readiness'], diff --git a/frontend/src/lib/types.ts b/frontend/src/lib/types.ts index bd15b43..aef0736 100644 --- a/frontend/src/lib/types.ts +++ b/frontend/src/lib/types.ts @@ -193,6 +193,7 @@ export interface ScheduleConfig { schedule_daily_pipeline_cron: string; schedule_dolt_earnings_cron: string; schedule_sec_fundamentals_cron: string; + schedule_fundamentals_parity_cron: string; schedule_near_close_pipeline_cron: string; schedule_after_close_pipeline_cron: string; schedule_intraday_pipeline_cron: string; @@ -733,6 +734,7 @@ export interface MetricItem { industry: MetricIndustry | null; period_end: string | null; filed_date: string | null; + caveat: string | null; source: string; // 'sec' | 'legacy_api' } diff --git a/frontend/src/pages/AdminPage.tsx b/frontend/src/pages/AdminPage.tsx index 57f11e2..520faee 100644 --- a/frontend/src/pages/AdminPage.tsx +++ b/frontend/src/pages/AdminPage.tsx @@ -5,6 +5,7 @@ import { AlertSettings } from '../components/admin/AlertSettings'; import { SentimentProviderSettings } from '../components/admin/SentimentProviderSettings'; import { DataCleanup } from '../components/admin/DataCleanup'; import { JobControls } from '../components/admin/JobControls'; +import { FundamentalsParityPanel } from '../components/admin/FundamentalsParityPanel'; import { PerformanceSettings } from '../components/admin/PerformanceSettings'; import { PipelineReadinessPanel } from '../components/admin/PipelineReadinessPanel'; import { SystemEventsPanel } from '../components/admin/SystemEventsPanel'; @@ -48,6 +49,7 @@ export default function AdminPage() { {activeTab === 'Jobs' && (
+
diff --git a/tests/unit/test_fundamentals_derivation.py b/tests/unit/test_fundamentals_derivation.py index ca249f4..d77a0c2 100644 --- a/tests/unit/test_fundamentals_derivation.py +++ b/tests/unit/test_fundamentals_derivation.py @@ -104,6 +104,22 @@ def test_net_debt_leverage_and_share_dilution(): assert d.metrics["share_count_change_yoy"].value == pytest.approx(-10.0, abs=1e-6) +def test_split_suspect_share_move_suppresses_share_and_eps_comparisons(): + rows = _two_years() + for row in rows: + if row.fiscal_year == 2026: + row.shares_outstanding = 2000 # +100% resembles an unadjusted 2-for-1 split + + d = fd.derive(rows) + + for key in ("share_count_change_yoy", "eps_growth_yoy"): + series = d.metrics[key] + assert series.value is None + assert series.history[-1].value is None + assert "possible split" in series.caveat + assert d.metrics["revenue_growth_yoy"].value == pytest.approx(10.0) + + def test_valuation_inputs(): d = fd.derive(_two_years()) # TTM diluted EPS FY2026 = 1.1+1.21+1.32+1.43 = 5.06 diff --git a/tests/unit/test_fundamentals_parity.py b/tests/unit/test_fundamentals_parity.py new file mode 100644 index 0000000..a0495ae --- /dev/null +++ b/tests/unit/test_fundamentals_parity.py @@ -0,0 +1,209 @@ +"""A5 fundamentals parity report: read-only comparison + artifact archive.""" + +from __future__ import annotations + +from datetime import date, datetime, timezone + +import pytest +from sqlalchemy import func, select + +from app.models.data_import_run import DataImportRun +from app.models.earnings_event import EarningsEvent +from app.models.fundamental import FundamentalData +from app.models.fundamental_snapshot import FundamentalSnapshot +from app.models.ohlcv import OHLCVRecord +from app.models.ticker import Ticker +from app.services.fundamentals_parity_service import ( + build_report, + fundamental_score, + load_latest, + load_latest_csv, + load_latest_json, + store_report, +) + +UTC = timezone.utc +GENERATED = datetime(2026, 7, 23, 10, 30, tzinfo=UTC) + + +def _snapshot_rows(cik: str) -> list[FundamentalSnapshot]: + rows = [] + periods = ("Q1", "Q2", "Q3", "FY") + months = (3, 6, 9, 12) + for fy, multiplier in ((2025, 1.0), (2026, 1.1)): + revenues = [100 * multiplier, 110 * multiplier, 120 * multiplier, 130 * multiplier] + eps = [1.0 * multiplier, 1.1 * multiplier, 1.2 * multiplier, 1.3 * multiplier] + for index, period in enumerate(periods): + period_end = date(fy, months[index], 28) + rows.append( + FundamentalSnapshot( + cik=cik, + accession=f"{cik}-{fy}-{period}", + form="10-K" if period == "FY" else "10-Q", + filed_date=period_end, + accepted_at=datetime(fy, months[index], 28, tzinfo=UTC), + period_end=period_end, + fiscal_year=fy, + fiscal_period=period, + revenue=sum(revenues[: index + 1]), + operating_income=sum(revenues[: index + 1]) * 0.2, + diluted_eps=sum(eps[: index + 1]), + cfo=sum(revenues[: index + 1]) * 0.25, + capex=sum(revenues[: index + 1]) * 0.05, + depreciation_amortization=sum(revenues[: index + 1]) * 0.05, + cash_and_st_investments=40, + total_debt=100, + shares_outstanding=1000, + ) + ) + return rows + + +async def _seed(db_session): + first = Ticker(symbol="AAA", cik="0000000001", sic="3571") + second = Ticker(symbol="BBB", cik=None, sic=None) + db_session.add_all([first, second]) + await db_session.flush() + db_session.add_all(_snapshot_rows(first.cik)) + db_session.add_all( + [ + FundamentalData( + ticker_id=first.id, + pe_ratio=25, + revenue_growth=5, + earnings_surprise=0, + fetched_at=GENERATED, + ), + FundamentalData( + ticker_id=second.id, + pe_ratio=12, + revenue_growth=3, + earnings_surprise=None, + fetched_at=GENERATED, + ), + OHLCVRecord( + ticker_id=first.id, + date=date(2026, 7, 22), + open=100, + high=100, + low=100, + close=100, + volume=100, + ), + EarningsEvent( + ticker_id=first.id, + announce_date=date(2026, 7, 1), + session="amc", + eps_estimate=2, + eps_actual=2.2, + source="dolt_earnings", + ), + DataImportRun( + source="sec_facts", + revision="sec-rev", + status="promoted", + source_max_date=date(2026, 7, 22), + started_at=GENERATED, + completed_at=GENERATED, + ), + DataImportRun( + source="dolt_earnings", + revision="dolt-rev", + status="no_op", + source_max_date=date(2026, 7, 22), + started_at=GENERATED, + completed_at=GENERATED, + ), + ] + ) + await db_session.flush() + + +def test_score_formula_matches_production_rules(): + score = fundamental_score(pe_ratio=15, revenue_growth=0, earnings_surprise=0) + assert score == pytest.approx((100 + 50 + 50) / 3) + assert fundamental_score(pe_ratio=15, revenue_growth=None, earnings_surprise=None) is None + +async def test_report_compares_sources_and_leaves_database_untouched(db_session): + await _seed(db_session) + before = await db_session.scalar(select(func.count()).select_from(FundamentalData)) + + report = await build_report( + db_session, + generated_at=GENERATED, + today=date(2026, 7, 23), + ) + + after = await db_session.scalar(select(func.count()).select_from(FundamentalData)) + assert before == after == 2 + assert not db_session.new and not db_session.dirty and not db_session.deleted + assert report["read_only"] is True + assert report["approval_status"] == "pending_explicit_approval" + assert report["source_runs"]["sec_facts"]["revision"] == "sec-rev" + assert report["source_runs"]["dolt_earnings"]["revision"] == "dolt-rev" + + first = next(row for row in report["rows"] if row["symbol"] == "AAA") + assert first["fields"]["pe_ratio"]["candidate"] == pytest.approx( + 100 / 5.06, abs=1e-4 + ) + assert first["fields"]["revenue_growth"]["candidate"] == pytest.approx(10) + assert first["fields"]["earnings_surprise"]["candidate"] == pytest.approx(10) + assert first["scores"]["candidate_fundamental"] is not None + assert report["summary"]["universe_count"] == 2 + assert report["summary"]["field_stats"]["pe_ratio"]["both_available"] == 1 + + +async def test_artifacts_archive_and_latest_manifest(db_session, tmp_path): + await _seed(db_session) + report = await build_report( + db_session, + generated_at=GENERATED, + today=date(2026, 7, 23), + ) + + paths = store_report(report, tmp_path) + + assert tmp_path.joinpath("latest.json").exists() + assert paths["json"].endswith(".json") and paths["csv"].endswith(".csv") + assert load_latest(tmp_path)["generated_at"] == GENERATED.isoformat() + csv_artifact = load_latest_csv(tmp_path) + assert csv_artifact is not None + assert csv_artifact[0].endswith(".csv") + assert "legacy_fundamental,candidate_fundamental" in csv_artifact[1] + assert "AAA" in csv_artifact[1] + json_artifact = load_latest_json(tmp_path) + assert json_artifact is not None and '"rows"' in json_artifact[1] + + +async def test_admin_endpoints_return_compact_summary_and_downloads( + client, db_session, tmp_path, monkeypatch +): + from app.config import settings + from app.dependencies import require_admin + from app.main import app + + await _seed(db_session) + report = await build_report( + db_session, + generated_at=GENERATED, + today=date(2026, 7, 23), + ) + store_report(report, tmp_path) + monkeypatch.setattr(settings, "fundamentals_parity_report_dir", str(tmp_path)) + app.dependency_overrides[require_admin] = lambda: None + try: + summary_response = await client.get("/api/v1/admin/fundamentals-parity") + assert summary_response.status_code == 200 + summary = summary_response.json()["data"] + assert summary["summary"]["universe_count"] == 2 + assert "rows" not in summary + + csv_response = await client.get("/api/v1/admin/fundamentals-parity/csv") + assert csv_response.status_code == 200 + assert "AAA" in csv_response.json()["data"]["content"] + + json_response = await client.get("/api/v1/admin/fundamentals-parity/json") + assert json_response.status_code == 200 + assert '"rows"' in json_response.json()["data"]["content"] + finally: + app.dependency_overrides.pop(require_admin, None) diff --git a/tests/unit/test_schedule_config.py b/tests/unit/test_schedule_config.py index bc745a2..211edd3 100644 --- a/tests/unit/test_schedule_config.py +++ b/tests/unit/test_schedule_config.py @@ -85,6 +85,7 @@ class TestTradingDayCrons: ( ("schedule_dolt_earnings_cron", 2, 30), ("schedule_sec_fundamentals_cron", 4, 0), + ("schedule_fundamentals_parity_cron", 5, 30), ), ) def test_shadow_imports_run_daily_at_expected_et_time( diff --git a/tests/unit/test_scheduler.py b/tests/unit/test_scheduler.py index bb69087..07fb8dd 100644 --- a/tests/unit/test_scheduler.py +++ b/tests/unit/test_scheduler.py @@ -11,6 +11,7 @@ from app.scheduler import ( _resume_tickers, _last_successful, _run_shadow_import, + run_fundamentals_parity_report, configure_scheduler, get_job_runtime_snapshot, queue_backtest_options, @@ -112,6 +113,7 @@ class TestConfigureScheduler: "fundamental_collector", "dolt_earnings_import", "sec_fundamentals_import", + "fundamentals_parity_report", "rr_scanner", "shadow_book", "ticker_universe_sync", @@ -145,6 +147,7 @@ class TestConfigureScheduler: "fundamental_collector", "dolt_earnings_import", "sec_fundamentals_import", + "fundamentals_parity_report", "market_regime", "near_close_pipeline", "regime_monitor", @@ -243,3 +246,32 @@ class TestShadowImportJobs: runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "skipped" assert runtime["message"] == "Disabled" + + +async def test_fundamentals_parity_job_surfaces_report_summary(monkeypatch): + async def enabled(db, job_name): + return True + + async def generated(db, report_dir): + return ( + { + "generated_at": "2026-07-23T10:30:00+00:00", + "summary": { + "universe_count": 511, + "fundamental_score_material_changes": 12, + }, + }, + {"json": "report.json", "csv": "report.csv"}, + ) + + monkeypatch.setattr("app.scheduler.async_session_factory", TestShadowImportJobs._session_factory) + monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) + monkeypatch.setattr( + "app.scheduler.fundamentals_parity_service.generate_and_store", generated + ) + + await run_fundamentals_parity_report() + + runtime = get_job_runtime_snapshot("fundamentals_parity_report") + assert runtime["status"] == "completed" + assert runtime["message"] == "511 tickers · 12 material score changes"