============================================================ SL-004 — EXACT INGEST CONTRACT INSPECTION READ-ONLY / NO DB / NO PRODUCTION MUTATION ============================================================ UTC=2026-09-10T16:18:56Z HOST=Azeroth USER=root UID=0 RC=/data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453 ============================================================ 1. ingest_security — EXACT SOURCE ============================================================ """ Ingestion pipeline (spec §45/§52, docs/ARCHITECTURE.md §4). One security's worth of raw provider data -> normalized point-in-time rows in Postgres, source-hierarchy conflict resolution applied, then a `financials.updated` signal that triggers recompute (app/workers/ recompute.py). Demo and live data are never blended for the same security (docs/DATA_SOURCES.md §9) — enforced in `_assert_no_demo_live_mix` below, not just documented. AUDIT FIX (StockLab overhaul, final engineering pass, Part A6, docs/AUDIT_PROVIDER_CONFLICT_A6.md): `resolve_conflict` used to be imported here but never called — this file ingested exactly one provider's data per call, so there was nothing for it to resolve BETWEEN (docs/ AUDIT_PROVIDER_RESILIENCE.md's finding). `ingest_security()` now optionally accepts a `secondary_adapter`; when given, both providers' line items are resolved field-by-field via `resolve_line_items()` (app/adapters/resolver.py) and real disagreements are persisted as `DataQuality` rows. When `secondary_adapter` is None — the default, and what every existing caller still passes — behavior is byte-identical to before this pass: nothing about the single-source path changed. """ from __future__ import annotations from datetime import date, timedelta from sqlalchemy.orm import Session from app.adapters.base import ( ProviderAdapter, ProviderAuthError, ProviderNotFoundError, ProviderRateLimitError, ) from app.adapters.resolver import merge_statement_line_items, resolve_line_items from app.adapters.validation import ValidationReport, validate_line_items from app.core.config import get_settings from app.core.db import SessionLocal from app.core.logging import get_logger from app.engines.identity import resolve_company_identity from app.engines.reference_data import canonical_mic, exchange_for_mic from app.models import ( BalanceSheet, CashFlow, Company, Country, DataQuality, Estimate, Exchange, FinancialPeriod, Industry, IncomeStatement, Price, Sector, Security, Shares, Source, ) from app.workers.celery_app import celery_app logger = get_logger(__name__) class ReferenceDataError(RuntimeError): """A NOT NULL reference row could not be resolved, so ingestion must not proceed. AUDIT FIX (DEMO ingestion defect). `get_or_create_company()` used to build `Company` and `Exchange` rows whose `country_id` was a conditional expression falling back to a null. Both columns are **NOT NULL** (`alembic/versions/0001_initial_schema.py`, `companies.country_id` and `exchanges.country_id`), so that `else None` branch could never succeed: it handed Postgres a NULL for a NOT NULL column and the whole transaction died with psycopg.errors.NotNullViolation: null value in column "country_id" of relation "exchanges" several statements later, with a traceback pointing at `db.flush()` rather than at the missing reference row. The code was written as if the column were nullable; it is not. Failing here instead is the correct behaviour and not merely a nicer message: * an exchange row without a country is **not representable** in this schema, so there is no value to fall back to — inventing one (e.g. the company's country) is exactly the shared-reference-table corruption §18 forbids and that a previous pass removed; * the error names the ticker, the field and the specific remedy, so the fix is a one-line addition to a reference table rather than a database forensics session; * it is raised **before** any row is added to the session, so the caller's transaction is still clean and one bad ticker cannot poison the rest of a seed run. """ # Field tuples pulled out to module level (StockLab overhaul, Part A6) -- previously inline in the # single setattr loop below; now shared between the single-source path and resolve_line_items() # calls so the two can never drift out of sync with each other. _INCOME_FIELDS = ( "revenue", "cogs", "gross_profit", "operating_income", "ebit", "ebitda", "net_income", "eps_diluted", "tax_expense", "pretax_income", "interest_expense", "depreciation_and_amortization", "stock_based_compensation", ) _BALANCE_FIELDS = ( "cash_and_equivalents", "short_term_investments", "total_debt", "shareholders_equity", "minority_interest", "preferred_equity", "goodwill", "intangible_assets", "current_assets", "current_liabilities", "receivables", "inventory", ) _CASH_FLOW_FIELDS = ("operating_cash_flow", "capital_expenditure", "dividends_paid", "buybacks", "stock_issuance") _SHARES_FIELDS = ("shares_outstanding", "diluted_shares") def _get_or_create_source(db: Session, provider: str, tier: str, is_demo: bool, endpoint: str) -> Source: src = Source(provider=provider, provider_tier=tier, is_demo=is_demo, endpoint=endpoint) db.add(src) db.flush() return src def _assert_no_demo_live_mix(db: Session, security_id: str, is_demo: bool) -> None: """Refuse to write demo rows for a security that already has live rows, or vice versa (docs/DATA_SOURCES.md §9 — hard separation, enforced here, not just at read time).""" existing = ( db.query(FinancialPeriod.id) .join(Source, FinancialPeriod.source_id == Source.id) .filter(FinancialPeriod.security_id == security_id, Source.is_demo == (not is_demo)) .first() ) if existing: raise RuntimeError( f"Refusing to ingest {'demo' if is_demo else 'live'} data for security {security_id}: " f"it already has {'live' if is_demo else 'demo'} financial_periods rows." ) def _resolve_required_reference_data(db: Session, ticker: str, profile) -> tuple: """Resolve every NOT NULL reference row this company needs, or raise `ReferenceDataError`. Done in one place, up front, and **before any `db.add()`**, so that a missing reference row leaves the caller's session completely clean. That matters more than it looks: a seed loop that catches an exception per ticker cannot continue if the failing ticker already left pending rows behind — SQLAlchemy raises `PendingRollbackError` on every subsequent statement and the run ends with an empty database rather than 19 of 20 companies. Returns `(country, identity, mic, exchange, exchange_country)`, where `mic` is the venue's canonical MIC. `exchange` is the existing row when the venue is already known, otherwise None and the caller creates it with `mic` and `exchange_country`. """ country = ( db.query(Country).filter_by(iso2=profile.country_iso2).one_or_none() if profile.country_iso2 else None ) if country is None: # `companies.country_id` is NOT NULL. See ReferenceDataError above. raise ReferenceDataError( f"{ticker}: cannot resolve a company country. The provider profile reported " f"country_iso2={profile.country_iso2!r} and no `countries` row matches it. " f"`companies.country_id` is NOT NULL, so the company cannot be written without it. " f"Seed the country (scripts/seed_demo.py::seed_reference_data) or correct the " f"provider's country code — do not default it to the exchange's country (spec 18)." ) if not profile.exchange_mic: # `securities.exchange_id` is NOT NULL, and a security is by definition a listed line on a # venue. Writing the security with a NULL exchange fails the same way `exchanges.country_id` # did; writing it against an arbitrary exchange would be worse. raise ReferenceDataError( f"{ticker}: the provider profile carries no exchange MIC. `securities.exchange_id` is " f"NOT NULL — a security cannot be recorded without the venue it trades on. Fix the " f"provider profile rather than attaching the security to an arbitrary exchange." ) identity = resolve_company_identity(profile.country_iso2, profile.exchange_mic) # Look the venue up by its CANONICAL MIC. FMP returns `exchangeShortName` ("NASDAQ") where # EODHD returns a MIC ("XNAS"); without this, the same venue becomes two `exchanges` rows and # peer groups and country rollups silently split down the middle. mic = canonical_mic(profile.exchange_mic) exchange = db.query(Exchange).filter_by(mic=mic).one_or_none() exchange_country = None if exchange is None: exchange_country = ( db.query(Country).filter_by(iso2=identity.exchange_country.iso2).one_or_none() if identity.exchange_country.known else None ) if exchange_country is None: # `exchanges.country_id` is NOT NULL. This used to log a warning and then write a NULL # anyway, which turned a one-line reference-data gap into a NotNullViolation whose # traceback pointed at db.flush(). See ReferenceDataError above. logger.warning( "ingest.exchange.country_unknown", ticker=ticker, mic=profile.exchange_mic, exchange_country_iso2=identity.exchange_country.iso2, reason=identity.exchange_country.reason, ) if not identity.exchange_country.known: raise ReferenceDataError( f"{ticker}: exchange MIC {profile.exchange_mic!r} is not in " f"EXCHANGE_COUNTRY_BY_MIC (app/engines/identity.py), so the venue's country " f"is unknown. `exchanges.country_id` is NOT NULL. Add the MIC to that table " f"with its ISO 3166-1 alpha-2 country — deliberately as a stated fact, not " f"inferred from the ticker suffix or from the company's own country." ) raise ReferenceDataError( f"{ticker}: exchange MIC {profile.exchange_mic!r} resolves to country " f"{identity.exchange_country.iso2!r}, but no `countries` row has that iso2. " f"`exchanges.country_id` is NOT NULL. Seed the country " f"(scripts/seed_demo.py::seed_reference_data) before ingesting this security." ) return country, identity, mic, exchange, exchange_country def get_or_create_company(db: Session, adapter: ProviderAdapter, ticker: str) -> Security: profile = adapter.get_company_profile(ticker) country, identity, mic, exchange, exchange_country = _resolve_required_reference_data( db, ticker, profile ) sector = db.query(Sector).filter_by(code=(profile.sector or "DEFAULT").upper()).one_or_none() if sector is None: sector = Sector(code=(profile.sector or "DEFAULT").upper(), name=profile.sector or "Unclassified") db.add(sector) db.flush() industry_code = (profile.industry or profile.sector or "DEFAULT").upper() industry = db.query(Industry).filter_by(code=industry_code).one_or_none() if industry is None: industry = Industry(code=industry_code, name=profile.industry or "Unclassified", sector_id=sector.id) db.add(industry) db.flush() # AUDIT FIX (final master pass, §18 — shared-reference-table corruption). # This used to create the Exchange row with `country_id = country.id`, where `country` is the # COMPANY's country. Exchanges are shared across every company listed on them, so the first # company ingested set that venue's country for everyone after it: ingest one Chinese ADR # first and NYSE is permanently recorded as being in China. §18 states the opposite explicitly # — "Company Country = China, Incorporation Country = Cayman Islands, Exchange Country = USA" # is a valid and common combination, and company country must never be derived from the venue # (nor, as here, the venue's country from the company). # `exchange_country_for_mic()` resolves the venue's country from an explicit MIC table: a fact # about the venue, independent of any company on it. An unrecognised MIC yields None rather # than a guess, and the reason is logged so the table can be extended deliberately. # The venue's country was resolved and validated in `_resolve_required_reference_data()` # above, before anything was added to the session. if exchange is None: # Normally unreachable: `seed_reference_data()` writes every venue in # `app/engines/reference_data.py::EXCHANGES` up front, with its real name and IANA # timezone. This branch is the fallback for a venue that table does not list — a provider # code with no ISO-assigned MIC ("OTC", "PNK", "AIM"), or a venue added by a provider # before the table catches up. It records the MIC as the name and UTC as the timezone # rather than inventing either; the country still comes from the MIC table, never from the # company (spec 18). ref = exchange_for_mic(mic) exchange = Exchange( mic=mic, name=ref.name if ref else mic, country_id=exchange_country.id, timezone=ref.timezone if ref else "UTC", ) db.add(exchange) db.flush() if identity.is_cross_border_listing: # Not an error — an ADR or a foreign listing. Logged because it is the case where a single # `country` column would have been misleading, and because it is worth being able to count # how much of the universe is cross-border. logger.info("ingest.company.cross_border_listing", ticker=ticker, company_country=identity.company_country.iso2, exchange_country=identity.exchange_country.iso2) security = db.query(Security).filter_by(ticker=ticker, exchange_id=exchange.id).one_or_none() if security is not None: return security company = Company( legal_name=profile.legal_name, display_name=profile.display_name, country_id=country.id, sector_id=sector.id, industry_id=industry.id, website=profile.website, description=profile.description, ) db.add(company) db.flush() security = Security( company_id=company.id, exchange_id=exchange.id, ticker=ticker, isin=profile.isin, currency=profile.currency, discovery_status="QUALIFIED", ) db.add(security) db.flush() return security def _record_conflicts(db: Session, entity_table: str, entity_id: str, conflicts: list, source_a: Source, source_b: Source) -> None: """Persist resolve_line_items()'s flagged disagreements as DataQuality rows (StockLab overhaul, Part A6). entity_id is the owning FinancialPeriod's id for every statement type -- IncomeStatement/BalanceSheet/CashFlow/Shares are each a 1:1 extension of FinancialPeriod with no independently meaningful id of their own, so `entity_table` (e.g. "income_statements") is what disambiguates which statement a given row's conflicting field belongs to, not entity_id itself. A deliberate simplification, not an oversight -- see docs/AUDIT_PROVIDER_CONFLICT_A6.md for the alternative (a dedicated per-statement-row id) and why it wasn't needed for this.""" for c in conflicts: db.add(DataQuality( entity_table=entity_table, entity_id=entity_id, field=c.field, status="CONFLICTING", source_a_id=source_a.id, value_a=c.value_a, source_b_id=source_b.id if source_b is not None else None, value_b=c.value_b, selected_value=c.selected_value, selected_source_id=source_a.id if c.selected_value == c.value_a else (source_b.id if source_b is not None else None), selection_reason=c.selection_reason, )) def _fetch_statements_by_period(adapter, method, ticker: str, label: str) -> dict: """Fetch one statement type and index it by (period_end, period_type). Returns an empty dict -- never raises -- when the adapter has not implemented this statement (EODHDAdapter's balance-sheet/cash-flow methods are `NotImplementedError` today). An adapter with a partial implementation must degrade to "that statement's columns stay empty for this run", exactly as before this fix, rather than crashing every ingestion task. """ try: rows = method(ticker) except NotImplementedError: logger.warning( "ingest.security.statement_not_implemented", ticker=ticker, provider=adapter.name, statement=label, ) return {} return {(r.period_end, r.period_type): r.line_items for r in rows} def _ingest_prices( db: Session, adapter: ProviderAdapter, security_id: str, ticker: str, source: Source, ) -> int: """Ingest daily OHLCV bars for one security, idempotently by security/date.""" end = date.today() start = end - timedelta(days=365) bars = adapter.get_prices(ticker, start, end) written = 0 for bar in bars: row = ( db.query(Price) .filter_by(security_id=security_id, date=bar.date) .one_or_none() ) if row is None: row = Price( security_id=security_id, date=bar.date, open=bar.open, high=bar.high, low=bar.low, close=bar.close, adjusted_close=bar.adjusted_close, volume=bar.volume, currency=bar.currency, source_id=source.id, ) db.add(row) written += 1 else: row.open = bar.open row.high = bar.high row.low = bar.low row.close = bar.close row.adjusted_close = bar.adjusted_close row.volume = bar.volume row.currency = bar.currency row.source_id = source.id logger.info( "ingest.prices.complete", ticker=ticker, provider=adapter.name, bars=len(bars), inserted=written, ) return written def ingest_security(db: Session, adapter: ProviderAdapter, ticker: str, secondary_adapter: ProviderAdapter | None = None, include_quarterly: bool = False) -> str: """Ingest one security's financials from `adapter` (the primary/only source for this call). AUDIT FIX (StockLab overhaul, final engineering pass, Part A6): `secondary_adapter` is new and optional. When None -- the default, and what every caller in this codebase passed before this pass -- this function's behavior is completely unchanged: every field comes straight from `adapter`, exactly as before. When given a second adapter, each period's line items are resolved field-by-field against the primary's via `resolve_line_items()` (source-hierarchy tie-break, real disagreements flagged) instead of blindly overwritten by whichever ran last. See docs/AUDIT_PROVIDER_CONFLICT_A6.md for the full design, the period-alignment limitation, and honest TESTED/NOT TESTED status -- none of the dual-source path has been executed against real provider data or a real database in this build environment. """ is_demo = adapter.tier == "DEMO" security = get_or_create_company(db, adapter, ticker) _assert_no_demo_live_mix(db, security.id, is_demo) source = _get_or_create_source(db, adapter.name, adapter.tier, is_demo, f"ingest:{ticker}") # AUDIT FIX (StockLab final engineering pass, Part A9 -- CRITICAL, see # docs/AUDIT_INGEST_STATEMENTS_A9.md). Before this fix, this function fetched ONLY the income # statement and then wrote the BalanceSheet, CashFlow and Shares rows out of that same # response's line items. adapter.get_balance_sheets() and adapter.get_cash_flows() were never # called anywhere in the application. That is invisible in DEMO mode (DemoDataAdapter returns # one fully-merged dict from all three methods) and catastrophic against FMP, whose income # response carries only _INCOME_MAP's fields -- every balance-sheet, cash-flow and share-count # column would have been written NULL for real provider data, taking most metrics, all five # pillar scores and every valuation down with it. income_periods = adapter.get_income_statements(ticker) balance_by_period = _fetch_statements_by_period(adapter, adapter.get_balance_sheets, ticker, "balance_sheets") cash_by_period = _fetch_statements_by_period(adapter, adapter.get_cash_flows, ticker, "cash_flows") # Part A9: optional quarterly ingestion. Off by default (Settings.INGEST_QUARTERLY_PERIODS). # Quarterly rows are stored as ordinary FinancialPeriod rows with period_type Q1..Q4 -- the # (security_id, period_end, period_type, filing_date) unique constraint keeps them distinct # from the annual rows, and every existing query that filters period_type == "FY" is # unaffected. app/engines/ttm.py turns four of them into a TTM basis at recompute time. if include_quarterly: try: quarterly_income = adapter.get_income_statements(ticker, period="quarter") except NotImplementedError: logger.warning("ingest.security.quarterly_not_implemented", ticker=ticker, provider=adapter.name) quarterly_income = [] if quarterly_income: income_periods = list(income_periods) + list(quarterly_income) balance_by_period.update(_fetch_statements_by_period( adapter, lambda t: adapter.get_balance_sheets(t, period="quarter"), ticker, "balance_sheets_quarterly")) cash_by_period.update(_fetch_statements_by_period( adapter, lambda t: adapter.get_cash_flows(t, period="quarter"), ticker, "cash_flows_quarterly")) secondary_source = None secondary_by_period: dict = {} secondary_balance: dict = {} secondary_cash: dict = {} if secondary_adapter is not None: if secondary_adapter.tier == "DEMO" or is_demo: # Mirrors _assert_no_demo_live_mix's own rule at the single-call level: a demo source # can never be one half of a dual-source resolution, since demo data is synthetic and # was never meant to be cross-checked against (or corrupt) real provider data. raise RuntimeError("Refusing dual-source ingestion: DEMO cannot be paired with a second provider.") secondary_source = _get_or_create_source(db, secondary_adapter.name, secondary_adapter.tier, is_demo, f"ingest:{ticker}") # AUDIT (Part A6): period alignment across two providers is keyed on (period_end, # period_type) only -- NOT filing_date, which two providers can legitimately report a day # or more apart for the "same" filing. This is a real, documented limitation: a provider # that reports a fiscal period boundary a few days off from the other (a genuine, if rare, # real-world occurrence) will not be matched and that period is silently treated as # primary-only for this run, not flagged as a mismatch. See docs/AUDIT_PROVIDER_CONFLICT_A6.md. try: secondary_by_period = { (p.period_end, p.period_type): p for p in secondary_adapter.get_income_statements(ticker) } # Part A9: the secondary provider's balance-sheet and cash-flow responses are fetched # too, so cross-source conflict resolution covers all four statement groups rather # than only the income statement. _fetch_statements_by_period returns {} for an # adapter that has not implemented a statement, so this adds no new failure mode. secondary_balance = _fetch_statements_by_period( secondary_adapter, secondary_adapter.get_balance_sheets, ticker, "balance_sheets") secondary_cash = _fetch_statements_by_period( secondary_adapter, secondary_adapter.get_cash_flows, ticker, "cash_flows") except NotImplementedError: # AUDIT FIX (Part A6, real bug caught before this shipped): EODHDAdapter.get_income_ # statements/get_balance_sheets/get_cash_flows are NotImplementedError today # (app/adapters/eodhd.py's own module docstring -- statement-level field mapping was # deliberately deferred until verified against a live response). EODHD is the only # other real (non-DEMO) adapter besides FMP, so PROVIDER_SECONDARY="EODHD" is the one # realistic dual-source configuration this codebase can actually be set to today -- and # without this except clause, setting it would crash every single ingestion task the # instant this function tried to call the secondary adapter's statement method. Treated # as "no secondary statement data available this run", not a crash: falls back to # exactly the primary-only behavior secondary_adapter=None already has for every period. logger.warning("ingest.security.secondary_statements_not_implemented", ticker=ticker, secondary_provider=secondary_adapter.name) for p in income_periods: fp = db.query(FinancialPeriod).filter_by( security_id=security.id, period_end=p.period_end, period_type=p.period_type, filing_date=p.filing_date, ).one_or_none() if fp is None: fp = FinancialPeriod( security_id=security.id, period_end=p.period_end, period_type=p.period_type, filing_date=p.filing_date, currency=p.currency, source_id=source.id, ) db.add(fp) db.flush() key = (p.period_end, p.period_type) li = merge_statement_line_items( p.line_items, balance_by_period.get(key), cash_by_period.get(key), ) secondary_p = secondary_by_period.get(key) li_b = ( merge_statement_line_items( secondary_p.line_items, secondary_balance.get(key), secondary_cash.get(key) ) if secondary_p is not None else None ) secondary_tier = secondary_adapter.tier if (secondary_adapter is not None and li_b is not None) else None income = fp.income_statement or IncomeStatement(financial_period_id=fp.id) resolved, conflicts = resolve_line_items( li, li_b, _INCOME_FIELDS, adapter.tier, secondary_tier ) merged_report = ValidationReport() resolved = validate_line_items( resolved, merged_report, row_key=f"{ticker}:{p.period_end}:{p.period_type}:income", ) for field, value in resolved.items(): setattr(income, field, value) db.add(income) if conflicts: _record_conflicts(db, "income_statements", fp.id, conflicts, source, secondary_source) balance = fp.balance_sheet or BalanceSheet(financial_period_id=fp.id) resolved, conflicts = resolve_line_items( li, li_b, _BALANCE_FIELDS, adapter.tier, secondary_tier ) merged_report = ValidationReport() resolved = validate_line_items( resolved, merged_report, row_key=f"{ticker}:{p.period_end}:{p.period_type}:balance", ) for field, value in resolved.items(): setattr(balance, field, value) db.add(balance) if conflicts: _record_conflicts(db, "balance_sheets", fp.id, conflicts, source, secondary_source) cash_flow = fp.cash_flow or CashFlow(financial_period_id=fp.id) resolved, conflicts = resolve_line_items( li, li_b, _CASH_FLOW_FIELDS, adapter.tier, secondary_tier ) merged_report = ValidationReport() resolved = validate_line_items( resolved, merged_report, row_key=f"{ticker}:{p.period_end}:{p.period_type}:cash_flow", ) for field, value in resolved.items(): setattr(cash_flow, field, value) db.add(cash_flow) if conflicts: _record_conflicts(db, "cash_flows", fp.id, conflicts, source, secondary_source) shares = fp.shares or Shares(financial_period_id=fp.id) resolved, conflicts = resolve_line_items( li, li_b, _SHARES_FIELDS, adapter.tier, secondary_tier ) merged_report = ValidationReport() resolved = validate_line_items( resolved, merged_report, row_key=f"{ticker}:{p.period_end}:{p.period_type}:shares", ) for field, value in resolved.items(): setattr(shares, field, value) db.add(shares) if conflicts: _record_conflicts(db, "shares", fp.id, conflicts, source, secondary_source) # AUDIT FIX (final master pass, §21) — the fourth "engine built, never wired" defect. # `ProviderAdapter.get_estimates()` is implemented by FMP and by the demo adapter, # `ProviderEstimateRow` exists, and the `estimates` table exists with a unique constraint on # (security_id, period_end, metric, as_of_date). Nothing ever called it, nothing ever wrote a # row, and `build_snapshot_from_db()` never passed `forward_eps_estimate` — so it was None on # every real run and `forward_pe` and `eps_growth_forward` were PERMANENTLY NULL. `forward_pe` # is a member of the Valuation pillar, so every Valuation score was computed from six of its # seven metrics with nothing saying so. # Writes are idempotent against that unique constraint: an estimate already stored for the # same (period_end, metric, as_of_date) is updated in place, never duplicated on re-ingestion. estimates_written = _ingest_estimates(db, adapter, security.id, ticker, source) prices_written = _ingest_prices(db, adapter, security.id, ticker, source) db.commit() logger.info("ingest.security.complete", ticker=ticker, provider=adapter.name, is_demo=is_demo, periods=len(income_periods), estimates=estimates_written, prices=prices_written, secondary_provider=(secondary_adapter.name if secondary_adapter else None)) return security.id def _ingest_estimates(db: Session, adapter: ProviderAdapter, security_id: str, ticker: str, source: Source) -> int: """Fetch and persist forward analyst estimates. Returns how many rows were written or updated. Never raises for an adapter that has not implemented estimates (EODHD's is `NotImplementedError`) — that degrades to "no estimates this run", which is exactly what `select_forward_estimate()` is built to report honestly, rather than crashing an ingestion that otherwise succeeded. """ try: rows = adapter.get_estimates(ticker) except NotImplementedError: logger.warning("ingest.estimates.not_implemented", ticker=ticker, provider=adapter.name) return 0 written = 0 for r in rows: if r.consensus_value is None or r.period_end is None: continue existing = db.query(Estimate).filter_by( security_id=security_id, period_end=r.period_end, metric=r.metric, as_of_date=r.as_of_date, ).one_or_none() if existing is None: existing = Estimate( security_id=security_id, period_end=r.period_end, metric=r.metric, as_of_date=r.as_of_date, ) db.add(existing) existing.consensus_value = r.consensus_value existing.num_analysts = r.num_analysts existing.source_id = source.id written += 1 return written # AUDIT FIX (StockLab overhaul, Part 23): this task previously had NO retry policy at all — a # transient rate-limit (`ProviderRateLimitError`, e.g. HTTP 429) or network blip crashed the task # outright with no retry, despite `Settings.PROVIDER_MAX_RETRIES` existing in config and looking # like it was already wired in (it was not referenced anywhere in this file before this fix). # `ProviderAuthError` (bad/expired API key) and `ProviderNotFoundError` (unknown ticker) are # deliberately NOT in `autoretry_for` — retrying a bad API key or a ticker that doesn't exist just # wastes the retry budget on an error that will never resolve itself. @celery_app.task( name="app.workers.ingest.ingest_security_task", bind=True, autoretry_for=(ProviderRateLimitError, ConnectionError, TimeoutError), retry_backoff=True, # exponential backoff between attempts retry_backoff_max=300, # cap backoff at 5 minutes retry_jitter=True, # avoid a thundering herd of simultaneous retries max_retries=get_settings().PROVIDER_MAX_RETRIES, ) def ingest_security_task(self, ticker: str, provider_name: str) -> str: from app.adapters import build_adapter, get_secondary_adapter settings = get_settings() adapter = build_adapter(provider_name, settings) # AUDIT FIX (StockLab overhaul, final engineering pass, Part A6): dual-source conflict # resolution activates automatically, with no separate flag to remember to set, exactly when a # deployment has actually configured a second provider (Settings.PROVIDER_SECONDARY != "NONE") # AND this task is running the primary provider's own scheduled ingestion -- an explicit # re-ingest of one named provider (e.g. a manual backfill using only EODHD) does not also pull # in a second source, since that call already named the one provider it wants. When # PROVIDER_SECONDARY is "NONE" (the default), get_secondary_adapter() returns None and this is # a no-op -- ingest_security() below then runs exactly the single-source path it always has. secondary_adapter = ( get_secondary_adapter(settings) if provider_name == settings.PROVIDER_PRIMARY else None ) db = SessionLocal() try: security_id = ingest_security( db, adapter, ticker, secondary_adapter=secondary_adapter, include_quarterly=settings.INGEST_QUARTERLY_PERIODS, ) except ProviderAuthError: logger.error("ingest.security.auth_failed", ticker=ticker, provider=provider_name) raise # not retried — see autoretry_for note above except ProviderNotFoundError: logger.warning("ingest.security.not_found", ticker=ticker, provider=provider_name) raise # not retried finally: db.close() from app.workers.recompute import recompute_security_task recompute_security_task.delay(security_id) return security_id @celery_app.task(name="app.workers.ingest.ingest_universe_task") def ingest_universe_task(provider_name: str | None = None) -> int: from app.adapters import build_adapter settings = get_settings() adapter = build_adapter(provider_name or settings.PROVIDER_PRIMARY, settings) tickers = adapter.list_universe() for ticker in tickers: ingest_security_task.delay(ticker, provider_name or settings.PROVIDER_PRIMARY) logger.info("ingest.universe.dispatched", count=len(tickers)) return len(tickers) ============================================================ 2. ProviderAdapter — EXACT INTERFACE ============================================================ """ ProviderAdapter interface (spec §7). Every adapter — FMP, EODHD, a future filing-based adapter, or the DemoDataAdapter — implements this and only this; the ingestion worker and every caller depend on this interface, never on a concrete provider class, so swapping providers is a config change (`PROVIDER_PRIMARY` / `PROVIDER_SECONDARY`). """ from __future__ import annotations from abc import ABC, abstractmethod from datetime import date from typing import Optional from app.adapters.schemas import ( ProviderCompanyProfile, ProviderDividendRow, ProviderEstimateRow, ProviderFinancialPeriod, ProviderPriceBar, ) class ProviderRateLimitError(Exception): pass class ProviderAuthError(Exception): pass class ProviderNotFoundError(Exception): pass class ProviderAdapter(ABC): """All methods are synchronous-return dataclass lists — concrete adapters may be async internally (FMPAdapter uses httpx under the hood) but this ABC keeps the calling code simple for v1; see SPEC_COVERAGE.md for the note on making this fully async in a later pass.""" name: str tier: str # "PRIMARY" | "SECONDARY" | "OFFICIAL_FILING" | "DEMO" @abstractmethod def get_company_profile(self, ticker: str, exchange_mic: Optional[str] = None) -> ProviderCompanyProfile: ... @abstractmethod def get_income_statements(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: ... @abstractmethod def get_balance_sheets(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: ... @abstractmethod def get_cash_flows(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: ... @abstractmethod def get_prices(self, ticker: str, start: date, end: date) -> list[ProviderPriceBar]: ... @abstractmethod def get_estimates(self, ticker: str) -> list[ProviderEstimateRow]: ... @abstractmethod def get_dividends(self, ticker: str) -> list[ProviderDividendRow]: ... @abstractmethod def list_universe(self, exchange_mic: Optional[str] = None, country_iso2: Optional[str] = None) -> list[str]: """Return a batch of tickers to ingest for a given exchange/country (spec §5/§6 universe + discovery scanning).""" ... ============================================================ 3. Adapter schemas — EXACT DATA CONTRACT ============================================================ """ Canonical, provider-agnostic shapes that every ProviderAdapter must return (docs/DATA_SOURCES.md §7). The calculation engines never see a provider-specific field name — normalization from these shapes into `app.engines.types.LineItems` / DB rows happens once, in `app/workers/normalize.py`, not scattered across the codebase. """ from __future__ import annotations from dataclasses import dataclass from datetime import date from typing import Optional @dataclass(frozen=True) class ProviderCompanyProfile: ticker: str exchange_mic: Optional[str] legal_name: str display_name: str country_iso2: Optional[str] sector: Optional[str] industry: Optional[str] currency: str isin: Optional[str] = None website: Optional[str] = None description: Optional[str] = None beta: Optional[float] = None @dataclass(frozen=True) class ProviderFinancialPeriod: """One reporting period's worth of raw line items, provider-normalized field names already mapped to our canonical keys (this dataclass's field names == LineItems' field names so the mapping in app/engines/types.LineItems(**asdict(provider_period_minus_meta)) is direct).""" period_end: date period_type: str filing_date: Optional[date] currency: str line_items: dict # canonical field name -> value, matches LineItems fields (see types.py) @dataclass(frozen=True) class ProviderPriceBar: date: date open: Optional[float] high: Optional[float] low: Optional[float] close: float adjusted_close: Optional[float] volume: Optional[float] currency: str @dataclass(frozen=True) class ProviderEstimateRow: period_end: date metric: str consensus_value: Optional[float] num_analysts: Optional[int] as_of_date: date @dataclass(frozen=True) class ProviderDividendRow: ex_date: date pay_date: Optional[date] amount_per_share: float currency: str is_special: bool = False ============================================================ 4. DEMO ADAPTER — CONCRETE IMPLEMENTATION ============================================================ """ DemoDataAdapter — the ONLY adapter that runs without a real API key (docs/DATA_SOURCES.md §9). Produces 11 years of internally-consistent, deterministic SYNTHETIC financials for a set of real, well-known tickers spanning multiple regions/sectors, so the full pipeline (ingestion -> metrics -> scoring -> valuation -> recommendation -> screener -> UI) is demonstrable end to end without live provider access. Every row this adapter returns must be persisted with `sources.is_demo = True` by the caller — see `app/workers/ingest.py` — and the numbers here are NOT real reported financials for these companies; they are a plausible synthetic trajectory generated from a small seed profile per company, with a fixed random seed per ticker so the same "demo run" is reproducible. """ from __future__ import annotations import random from dataclasses import dataclass from datetime import date, timedelta from typing import Optional from app.adapters.base import ProviderAdapter, ProviderNotFoundError from app.adapters.schemas import ( ProviderCompanyProfile, ProviderDividendRow, ProviderEstimateRow, ProviderFinancialPeriod, ProviderPriceBar, ) YEARS_OF_HISTORY = 11 # current + 10 prior, enough for every CAGR window in FINANCIAL_FORMULAS.md @dataclass(frozen=True) class DemoSeedProfile: ticker: str exchange_mic: str legal_name: str country_iso2: str sector_code: str industry_code: str currency: str starting_revenue_musd: float # 10 years ago revenue_cagr: float gross_margin: float operating_margin: float net_margin_adj: float # adjustment vs. a naive EBIT*(1-tax) estimate, for realism tax_rate: float capex_pct_revenue: float debt_to_ebitda_target: float dividend_payout_pct: float # of net income; 0 for non-payers buyback_pct_of_fcf: float shares_outstanding_musd_equiv: float # millions of shares, current beta: float pe_assumption: float # used only to derive a plausible demo price # A representative global set spanning the spec's minimum region list (§5). Real, well-known # companies; ALL financial figures below are illustrative seed assumptions, not reported data. DEMO_SEED_PROFILES: list[DemoSeedProfile] = [ DemoSeedProfile("AAPL", "XNAS", "Apple Inc. (DEMO)", "US", "TECHNOLOGY", "CONSUMER_ELECTRONICS", "USD", 180000, 0.08, 0.42, 0.30, 0.0, 0.15, 0.02, 0.5, 0.15, 0.75, 15500, 1.25, 28), DemoSeedProfile("MSFT", "XNAS", "Microsoft Corp. (DEMO)", "US", "TECHNOLOGY", "SOFTWARE", "USD", 90000, 0.13, 0.68, 0.42, 0.0, 0.17, 0.10, 1.0, 0.28, 0.25, 7400, 0.90, 32), DemoSeedProfile("NVDA", "XNAS", "NVIDIA Corp. (DEMO)", "US", "TECHNOLOGY", "SEMICONDUCTORS", "USD", 11000, 0.35, 0.72, 0.45, 0.0, 0.16, 0.05, 0.1, 0.02, 0.10, 24500, 1.7, 45), DemoSeedProfile("JNJ", "XNYS", "Johnson & Johnson (DEMO)", "US", "HEALTHCARE", "PHARMACEUTICALS", "USD", 70000, 0.04, 0.66, 0.24, 0.0, 0.15, 0.05, 1.2, 0.55, 0.20, 2600, 0.55, 17), DemoSeedProfile("JPM", "XNYS", "JPMorgan Chase (DEMO)", "US", "FINANCIALS", "BANKS", "USD", 100000, 0.06, None, 0.35, 0.0, 0.22, None, None, 0.28, 0.20, 2900, 1.1, 12), DemoSeedProfile("PG", "XNYS", "Procter & Gamble (DEMO)", "US", "CONSUMER_STAPLES", "HOUSEHOLD_PRODUCTS", "USD", 65000, 0.03, 0.49, 0.22, 0.0, 0.19, 0.04, 1.0, 0.58, 0.20, 2350, 0.45, 24), DemoSeedProfile("XOM", "XNYS", "ExxonMobil (DEMO)", "US", "ENERGY", "OIL_GAS", "USD", 240000, 0.02, 0.22, 0.10, 0.0, 0.20, 0.07, 1.0, 0.40, 0.15, 4000, 1.15, 13), DemoSeedProfile("SAP", "XETR", "SAP SE (DEMO)", "DE", "TECHNOLOGY", "SOFTWARE", "EUR", 24000, 0.07, 0.72, 0.24, 0.0, 0.24, 0.03, 1.5, 0.35, 0.10, 1230, 0.95, 26), DemoSeedProfile("ASML", "XAMS", "ASML Holding (DEMO)", "NL", "TECHNOLOGY", "SEMICONDUCTOR_EQUIPMENT", "EUR", 11000, 0.15, 0.50, 0.32, 0.0, 0.16, 0.06, 0.3, 0.20, 0.05, 393, 1.15, 33), DemoSeedProfile("NESN", "XSWX", "Nestle S.A. (DEMO)", "CH", "CONSUMER_STAPLES", "PACKAGED_FOODS", "CHF", 88000, 0.02, 0.47, 0.16, 0.0, 0.20, 0.05, 1.3, 0.65, 0.15, 2650, 0.55, 20), DemoSeedProfile("MC", "XPAR", "LVMH (DEMO)", "FR", "CONSUMER_DISCRETIONARY", "LUXURY_GOODS", "EUR", 35000, 0.10, 0.68, 0.26, 0.0, 0.23, 0.05, 0.9, 0.45, 0.05, 500, 1.05, 24), DemoSeedProfile("NOVO-B", "XCSE", "Novo Nordisk (DEMO)", "DK", "HEALTHCARE", "PHARMACEUTICALS", "DKK", 90000, 0.14, 0.83, 0.44, 0.0, 0.20, 0.06, 0.6, 0.45, 0.05, 4480, 0.35, 34), DemoSeedProfile("SHEL", "XLON", "Shell plc (DEMO)", "GB", "ENERGY", "OIL_GAS", "USD", 260000, 0.02, 0.20, 0.09, 0.0, 0.30, 0.06, 1.0, 0.40, 0.20, 6900, 1.05, 11), DemoSeedProfile("7203", "XTKS", "Toyota Motor (DEMO)", "JP", "CONSUMER_DISCRETIONARY", "AUTOMOBILES", "JPY", 25000000, 0.03, 0.19, 0.08, 0.0, 0.24, 0.05, 1.8, 0.30, 0.05, 14300, 0.75, 10), DemoSeedProfile("005930", "XKRX", "Samsung Electronics (DEMO)", "KR", "TECHNOLOGY", "SEMICONDUCTORS", "KRW", 240000000, 0.05, 0.38, 0.16, 0.0, 0.20, 0.14, 0.4, 0.25, 0.02, 5970000, 1.0, 14), DemoSeedProfile("0700", "XHKG", "Tencent Holdings (DEMO)", "HK", "TECHNOLOGY", "INTERNET_SERVICES", "HKD", 400000, 0.14, 0.48, 0.32, 0.0, 0.15, 0.06, 0.5, 0.0, 0.10, 9400, 1.1, 22), DemoSeedProfile("005380", "XKRX", "Hyundai Motor (DEMO)", "KR", "CONSUMER_DISCRETIONARY", "AUTOMOBILES", "KRW", 100000000, 0.04, 0.20, 0.08, 0.0, 0.23, 0.04, 1.2, 0.25, 0.05, 208000, 1.1, 6), DemoSeedProfile("RELIANCE", "XBOM", "Reliance Industries (DEMO)", "IN", "ENERGY", "CONGLOMERATE", "INR", 4500000, 0.12, 0.28, 0.14, 0.0, 0.24, 0.09, 1.5, 0.10, 0.02, 6770, 0.95, 25), DemoSeedProfile("CBA", "XASX", "Commonwealth Bank (DEMO)", "AU", "FINANCIALS", "BANKS", "AUD", 20000, 0.04, None, 0.42, 0.0, 0.28, None, None, 0.70, 0.05, 1700, 0.85, 17), DemoSeedProfile("4SBK", "XBUL", "First Investment Bank (DEMO)", "BG", "FINANCIALS", "BANKS", "EUR", 500, 0.05, None, 0.30, 0.0, 0.10, None, None, 0.15, 0.0, 176, 0.90, 8), ] _INDEX_BY_TICKER = {p.ticker: p for p in DEMO_SEED_PROFILES} def _seeded_random(ticker: str) -> random.Random: return random.Random(f"STOCKLAB-DEMO-{ticker}") def _jitter(rng: random.Random, base: float, pct: float = 0.06) -> float: return base * (1 + rng.uniform(-pct, pct)) def _generate_annual_series(profile: DemoSeedProfile) -> list[dict]: """Returns YEARS_OF_HISTORY dicts, oldest first, each a full LineItems-shaped field set.""" rng = _seeded_random(profile.ticker) is_bank = profile.sector_code == "FINANCIALS" revenue = profile.starting_revenue_musd years = [] today_year = date.today().year for i in range(YEARS_OF_HISTORY): year = today_year - (YEARS_OF_HISTORY - 1 - i) if i > 0: revenue = revenue * (1 + _jitter(rng, profile.revenue_cagr)) gross_margin = profile.gross_margin if profile.gross_margin is not None else None cogs = revenue * (1 - gross_margin) if gross_margin is not None else None gross_profit = revenue - cogs if cogs is not None else None operating_income = revenue * _jitter(rng, profile.operating_margin, 0.04) ebit = operating_income d_and_a = revenue * 0.03 ebitda = ebit + d_and_a interest_expense = revenue * 0.01 if not is_bank else revenue * 0.005 pretax_income = ebit - interest_expense tax_expense = max(pretax_income, 0) * profile.tax_rate net_income = pretax_income - tax_expense shares = profile.shares_outstanding_musd_equiv * ((1 - profile.buyback_pct_of_fcf * 0.01) ** i) eps_diluted = net_income / shares if shares else None capex = revenue * profile.capex_pct_revenue if profile.capex_pct_revenue else revenue * 0.03 operating_cash_flow = net_income + d_and_a * 1.1 fcf = operating_cash_flow - capex dividends_paid = max(net_income, 0) * profile.dividend_payout_pct * 0.5 if profile.dividend_payout_pct else 0.0 buybacks = max(fcf, 0) * profile.buyback_pct_of_fcf if profile.buyback_pct_of_fcf else 0.0 ebitda_for_debt = ebitda if ebitda else revenue * 0.15 total_debt = ebitda_for_debt * profile.debt_to_ebitda_target if profile.debt_to_ebitda_target is not None else revenue * 0.6 cash = revenue * 0.15 shareholders_equity = revenue * 0.35 if not is_bank else revenue * 1.2 current_assets = revenue * 0.30 current_liabilities = revenue * 0.20 # AUDIT FIX (StockLab final engineering pass, Part A9): the demo generator never emitted # total_assets, short_term_debt, long_term_debt or lease_liabilities. total_assets in # particular is an input to ROA, asset turnover, the goodwill-concentration red flag and # the acquisition-driven-growth red flag -- all of which were therefore permanently # SKIPPED_INSUFFICIENT_DATA in demo mode, which is the mode this platform runs in by # default. Derived from the liabilities-and-equity side so the synthetic balance sheet # actually balances (A = L + E) instead of being an independently invented number: other_liabilities = revenue * 0.05 total_assets = shareholders_equity + total_debt + current_liabilities + other_liabilities # A conventional maturity split; the demo profiles carry no maturity schedule of their own. short_term_debt = total_debt * 0.15 long_term_debt = total_debt - short_term_debt lease_liabilities = revenue * 0.02 years.append(dict( period_end=date(year, 12, 31), period_type="FY", filing_date=date(year + 1, 2, 15), currency=profile.currency, revenue=round(revenue, 1), cogs=round(cogs, 1) if cogs else None, gross_profit=round(gross_profit, 1) if gross_profit else None, operating_income=round(operating_income, 1), ebit=round(ebit, 1), ebitda=round(ebitda, 1), net_income=round(net_income, 1), eps_diluted=round(eps_diluted, 4) if eps_diluted else None, tax_expense=round(tax_expense, 1), pretax_income=round(pretax_income, 1), interest_expense=round(interest_expense, 1), depreciation_and_amortization=round(d_and_a, 1), cash_and_equivalents=round(cash, 1), short_term_investments=0.0, total_debt=round(total_debt, 1), shareholders_equity=round(shareholders_equity, 1), total_assets=round(total_assets, 1), short_term_debt=round(short_term_debt, 1), long_term_debt=round(long_term_debt, 1), lease_liabilities=round(lease_liabilities, 1), minority_interest=0.0, preferred_equity=0.0, goodwill=round(revenue * 0.1, 1), intangible_assets=round(revenue * 0.05, 1), current_assets=round(current_assets, 1), current_liabilities=round(current_liabilities, 1), receivables=round(revenue * 0.08, 1), inventory=round(revenue * 0.06, 1) if gross_margin else 0.0, shares_outstanding=round(shares, 2), diluted_shares=round(shares, 2), operating_cash_flow=round(operating_cash_flow, 1), capital_expenditure=round(capex, 1), dividends_paid=round(dividends_paid, 1), buybacks=round(buybacks, 1), stock_issuance=0.0, stock_based_compensation=round(revenue * 0.01, 1), )) return years class DemoDataAdapter(ProviderAdapter): name = "DEMO" tier = "DEMO" def _profile(self, ticker: str) -> DemoSeedProfile: p = _INDEX_BY_TICKER.get(ticker.upper()) if not p: raise ProviderNotFoundError(f"No DEMO seed profile for {ticker}. Available: {sorted(_INDEX_BY_TICKER)}") return p def get_company_profile(self, ticker: str, exchange_mic: Optional[str] = None) -> ProviderCompanyProfile: p = self._profile(ticker) return ProviderCompanyProfile( ticker=p.ticker, exchange_mic=p.exchange_mic, legal_name=p.legal_name, display_name=p.legal_name, country_iso2=p.country_iso2, sector=p.sector_code, industry=p.industry_code, currency=p.currency, beta=p.beta, description="DEMO DATA — synthetic financials, not reported figures. See docs/DATA_SOURCES.md §9.", ) def _statements(self, ticker: str) -> list[dict]: return _generate_annual_series(self._profile(ticker)) def get_income_statements(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: return self._to_periods(ticker, limit) def get_balance_sheets(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: return self._to_periods(ticker, limit) def get_cash_flows(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: return self._to_periods(ticker, limit) def _to_periods(self, ticker: str, limit: int) -> list[ProviderFinancialPeriod]: rows = self._statements(ticker)[-limit:] p = self._profile(ticker) return [ ProviderFinancialPeriod( period_end=r["period_end"], period_type=r["period_type"], filing_date=r["filing_date"], currency=p.currency, line_items=r, ) for r in reversed(rows) # most recent first, matching FinancialSnapshot's convention ] def get_prices(self, ticker: str, start: date, end: date) -> list[ProviderPriceBar]: p = self._profile(ticker) latest = self._statements(ticker)[-1] eps = latest["eps_diluted"] or 1.0 base_price = max(eps * p.pe_assumption, 1.0) rng = _seeded_random(ticker + "-price") bars = [] d = start price = base_price * 0.85 # start a bit below today's synthetic price, drift up to it while d <= end: if d.weekday() < 5: price = max(price * (1 + rng.uniform(-0.01, 0.011)), 0.01) bars.append(ProviderPriceBar(date=d, open=price, high=price * 1.01, low=price * 0.99, close=price, adjusted_close=price, volume=1_000_000, currency=p.currency)) d += timedelta(days=1) if bars: bars[-1] = ProviderPriceBar(bars[-1].date, bars[-1].open, bars[-1].high, bars[-1].low, base_price, base_price, bars[-1].volume, p.currency) return bars def get_estimates(self, ticker: str) -> list[ProviderEstimateRow]: rows = self._statements(ticker) last_eps = rows[-1]["eps_diluted"] p = self._profile(ticker) if last_eps is None: return [] forward = last_eps * (1 + p.revenue_cagr) return [ProviderEstimateRow( period_end=date(date.today().year + 1, 12, 31), metric="eps", consensus_value=round(forward, 4), num_analysts=12, as_of_date=date.today(), )] def get_dividends(self, ticker: str) -> list[ProviderDividendRow]: p = self._profile(ticker) if not p.dividend_payout_pct: return [] rows = self._statements(ticker) out = [] for r in rows[-5:]: if r["dividends_paid"] and r["shares_outstanding"]: dps = r["dividends_paid"] / r["shares_outstanding"] out.append(ProviderDividendRow( ex_date=date(r["period_end"].year, 11, 15), pay_date=date(r["period_end"].year, 12, 1), amount_per_share=round(dps, 4), currency=p.currency, )) return out def list_universe(self, exchange_mic: Optional[str] = None, country_iso2: Optional[str] = None) -> list[str]: return [ p.ticker for p in DEMO_SEED_PROFILES if (exchange_mic is None or p.exchange_mic == exchange_mic) and (country_iso2 is None or p.country_iso2 == country_iso2) ] ============================================================ 5. FMP ADAPTER — CONCRETE IMPLEMENTATION / PERIOD SHAPE ============================================================ """ Financial Modeling Prep adapter — primary provider (docs/DATA_SOURCES.md §3). Endpoint paths and field names follow FMP's documented v3 REST API. This adapter has NOT been exercised against the live API in this build environment (network egress here is restricted to a package-registry allowlist that does not include financialmodelingprep.com — see docs/DATA_SOURCES.md §3 and TROUBLESHOOTING.md). Before relying on this in production: run `scripts/probe_provider.py --provider fmp --ticker AAPL` with a real `FMP_API_KEY` and diff the actual response shape against the field mapping below — FMP's field names do shift between API versions and this mapping was written from documented/known conventions, not a live response. """ from __future__ import annotations from datetime import date, datetime from typing import Optional import httpx from app.adapters.base import ProviderAdapter, ProviderAuthError, ProviderNotFoundError, ProviderRateLimitError from app.adapters.validation import ValidationReport, validate_line_items, validate_price_bar from app.adapters.schemas import ( ProviderCompanyProfile, ProviderDividendRow, ProviderEstimateRow, ProviderFinancialPeriod, ProviderPriceBar, ) # FMP field name -> our canonical LineItems field name (see app/engines/types.py::LineItems) _INCOME_MAP = { "revenue": "revenue", "costOfRevenue": "cogs", "grossProfit": "gross_profit", "operatingIncome": "operating_income", "ebitda": "ebitda", "netIncome": "net_income", "eps": "eps_basic", "epsdiluted": "eps_diluted", "incomeTaxExpense": "tax_expense", "incomeBeforeTax": "pretax_income", "interestExpense": "interest_expense", "depreciationAndAmortization": "depreciation_and_amortization", } _BALANCE_MAP = { "cashAndCashEquivalents": "cash_and_equivalents", "shortTermInvestments": "short_term_investments", "totalDebt": "total_debt", "shortTermDebt": "short_term_debt", "longTermDebt": "long_term_debt", "capitalLeaseObligations": "lease_liabilities", "totalAssets": "total_assets", "totalCurrentAssets": "current_assets", "totalCurrentLiabilities": "current_liabilities", "totalStockholdersEquity": "shareholders_equity", "minorityInterest": "minority_interest", "preferredStock": "preferred_equity", "goodwill": "goodwill", "intangibleAssets": "intangible_assets", "netReceivables": "receivables", "inventory": "inventory", } _CASHFLOW_MAP = { "operatingCashFlow": "operating_cash_flow", "capitalExpenditure": "capital_expenditure", "freeCashFlow": "free_cash_flow", "dividendsPaid": "dividends_paid", "commonStockRepurchased": "buybacks", "commonStockIssued": "stock_issuance", "netIncome": "net_income", } def _map_fields(raw: dict, mapping: dict) -> dict: out = {} for fmp_key, canonical_key in mapping.items(): val = raw.get(fmp_key) if val is not None: # capex/buybacks come back negative (cash outflow) from FMP by convention — store as # positive magnitudes, matching LineItems' documented convention. if canonical_key in ("capital_expenditure", "dividends_paid", "buybacks") and val < 0: val = -val out[canonical_key] = val return out def _parse_date(s: Optional[str]) -> Optional[date]: if not s: return None return datetime.strptime(s[:10], "%Y-%m-%d").date() class FMPAdapter(ProviderAdapter): name = "FMP" tier = "PRIMARY" def __init__( self, api_key: str, base_url: str = "https://financialmodelingprep.com/api/v3", timeout: float = 20.0, transport: Optional[httpx.BaseTransport] = None, ): self._api_key = api_key self._base_url = base_url.rstrip("/") # `transport` is exposed purely for testing (httpx.MockTransport) — production callers # never pass it, so real requests always go over the network via the default transport. self._client = httpx.Client(base_url=self._base_url, timeout=timeout, transport=transport) def _get(self, path: str, **params) -> list | dict: params["apikey"] = self._api_key resp = self._client.get(path, params=params) if resp.status_code == 401: raise ProviderAuthError(f"FMP auth failed for {path}") if resp.status_code == 429: raise ProviderRateLimitError(f"FMP rate limit hit for {path}") if resp.status_code == 404: raise ProviderNotFoundError(f"FMP 404 for {path}") resp.raise_for_status() return resp.json() def get_company_profile(self, ticker: str, exchange_mic: Optional[str] = None) -> ProviderCompanyProfile: data = self._get(f"/profile/{ticker}") if not data: raise ProviderNotFoundError(f"No FMP profile for {ticker}") row = data[0] return ProviderCompanyProfile( ticker=ticker, exchange_mic=row.get("exchangeShortName"), legal_name=row.get("companyName", ticker), display_name=row.get("companyName", ticker), country_iso2=row.get("country"), sector=row.get("sector"), industry=row.get("industry"), currency=row.get("currency", "USD"), isin=row.get("isin"), website=row.get("website"), description=row.get("description"), beta=row.get("beta"), ) def _get_statements(self, endpoint: str, mapping: dict, ticker: str, period: str, limit: int) -> list[ProviderFinancialPeriod]: """AUDIT FIX (Part D): line items are now type-coerced and sanity-checked before leaving the adapter. A period is never rejected wholesale -- one impossible field does not invalidate the other thirty -- but a non-numeric or arithmetically-impossible value (a negative revenue, a negative share count) is dropped with a recorded reason instead of travelling into the metrics engine as a string or a nonsense number. `self.last_statement_validation` holds the report for the most recent call.""" rows = self._get(f"/{endpoint}/{ticker}", period=period, limit=limit) report = ValidationReport() out = [] for row in rows: period_end = _parse_date(row.get("date")) out.append(ProviderFinancialPeriod( period_end=period_end, period_type="FY" if period == "annual" else row.get("period", "Q"), filing_date=_parse_date(row.get("fillingDate") or row.get("acceptedDate")), currency=row.get("reportedCurrency", "USD"), line_items=validate_line_items( _map_fields(row, mapping), report, row_key=str(period_end), ), )) self.last_statement_validation = report return out def get_income_statements(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: return self._get_statements("income-statement", _INCOME_MAP, ticker, period, limit) def get_balance_sheets(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: return self._get_statements("balance-sheet-statement", _BALANCE_MAP, ticker, period, limit) def get_cash_flows(self, ticker: str, period: str = "annual", limit: int = 11) -> list[ProviderFinancialPeriod]: return self._get_statements("cash-flow-statement", _CASHFLOW_MAP, ticker, period, limit) def get_prices(self, ticker: str, start: date, end: date) -> list[ProviderPriceBar]: """AUDIT FIX (Part D — docs/AUDIT_VALIDATION_D.md): see EODHDAdapter.get_prices() for the full reasoning. Same policy applied here so the two adapters cannot diverge: reject the row, not the batch, and record why. `self.last_price_validation` holds the report.""" data = self._get(f"/historical-price-full/{ticker}", **{"from": start.isoformat(), "to": end.isoformat()}) rows = data.get("historical", []) if isinstance(data, dict) else [] report = ValidationReport() bars = [] for r in rows: clean = validate_price_bar( { "date": r.get("date"), "open": r.get("open"), "high": r.get("high"), "low": r.get("low"), "close": r.get("close"), "adjusted_close": r.get("adjClose"), "volume": r.get("volume"), }, report, ) if clean is None: continue bars.append(ProviderPriceBar( date=_parse_date(clean["date"]), open=clean["open"], high=clean["high"], low=clean["low"], close=clean["close"], adjusted_close=clean["adjusted_close"], volume=clean["volume"], currency="USD", )) self.last_price_validation = report return bars def get_estimates(self, ticker: str) -> list[ProviderEstimateRow]: rows = self._get(f"/analyst-estimates/{ticker}") out = [] for r in rows: period_end = _parse_date(r.get("date")) for metric_key, fmp_key in (("eps", "estimatedEpsAvg"), ("revenue", "estimatedRevenueAvg")): if r.get(fmp_key) is not None: out.append(ProviderEstimateRow( period_end=period_end, metric=metric_key, consensus_value=r[fmp_key], num_analysts=r.get("numberAnalystEstimatedRevenue"), as_of_date=date.today(), )) return out def get_dividends(self, ticker: str) -> list[ProviderDividendRow]: data = self._get(f"/historical-price-full/stock_dividend/{ticker}") rows = data.get("historical", []) if isinstance(data, dict) else [] return [ ProviderDividendRow( ex_date=_parse_date(r["date"]), pay_date=_parse_date(r.get("paymentDate")), amount_per_share=r.get("dividend", 0.0), currency="USD", ) for r in rows ] def list_universe(self, exchange_mic: Optional[str] = None, country_iso2: Optional[str] = None) -> list[str]: rows = self._get("/stock-screener", exchange=exchange_mic, country=country_iso2, limit=1000) return [r["symbol"] for r in rows if r.get("symbol")] ============================================================ 6. RESOLVER — FINAL PATCHED CONTRACT ============================================================ """Source hierarchy conflict resolution (spec §8/§9, docs/DATA_SOURCES.md §4).""" from __future__ import annotations from dataclasses import dataclass from typing import Optional # Lower index = higher priority. SOURCE_HIERARCHY = ["OFFICIAL_FILING", "PRIMARY", "SECONDARY", "CALCULATED"] @dataclass(frozen=True) class ConflictResolution: selected_value: Optional[float] selected_source_tier: Optional[str] selection_reason: str is_conflicting: bool def resolve_conflict(candidates: list[tuple[str, Optional[float]]], tolerance_pct: float = 0.005) -> ConflictResolution: """ `candidates` = [(source_tier, value), ...] for the same (security, period, field). Selects the highest-priority tier's value; flags CONFLICTING if a lower-priority source disagrees by more than `tolerance_pct` (default 0.5%) — small rounding/currency-conversion differences are not treated as a real conflict. """ present = [(tier, val) for tier, val in candidates if val is not None] if not present: return ConflictResolution(None, None, "No source provided a value", is_conflicting=False) present.sort(key=lambda c: SOURCE_HIERARCHY.index(c[0]) if c[0] in SOURCE_HIERARCHY else 99) selected_tier, selected_value = present[0] conflicting = False for tier, val in present[1:]: if selected_value == 0: if val != 0: conflicting = True elif abs(val - selected_value) / abs(selected_value) > tolerance_pct: conflicting = True reason = f"Highest-priority available source: {selected_tier}" if conflicting: reason += f"; disagreement > {tolerance_pct:.1%} with a lower-priority source — see data_quality row" return ConflictResolution(selected_value, selected_tier, reason, is_conflicting=conflicting) @dataclass(frozen=True) class FieldConflict: field: str value_a: Optional[float] value_b: Optional[float] selected_value: Optional[float] selection_reason: str def resolve_line_items( primary: dict, secondary: Optional[dict], fields: list[str], primary_tier: str, secondary_tier: Optional[str] = None, ) -> tuple[dict, list[FieldConflict]]: """ StockLab overhaul, final engineering pass, Part A6 -- the statement-level counterpart to `resolve_conflict()` above, which only ever resolves ONE field at a time. This applies it across a whole LineItems-shaped dict (`app/adapters/schemas.py`) so a real ingestion call site can resolve an entire income statement / balance sheet / cash flow / shares record between two providers in one call, without each call site re-implementing the per-field loop. `primary`/`secondary` are dicts of {field_name: value} for the SAME (security, period) from two different provider adapters -- the same shape `ingest.py`'s existing single-source loop already reads from `period.line_items`. `fields` is the list of field names to consider (the caller passes the exact tuple it already iterates for that statement type, e.g. ingest.py's existing `("revenue", "cogs", ...)` tuples -- kept as an explicit argument here rather than hardcoded, so this function has no per-statement-type knowledge of its own). Returns `(resolved, conflicts)`: - `resolved` contains exactly the fields at least one source provided a value for (mirrors the existing `if field in li: setattr(...)` "only touch fields that were actually provided" contract every single-source call site already relies on -- a field neither source reported is simply absent from `resolved`, never defaulted to 0/None-overwrite). - `conflicts` lists only the fields where `resolve_conflict()` flagged real (> tolerance_pct) disagreement -- the caller persists these as `DataQuality` rows; the vast majority of fields on any real statement are expected to end up here with zero conflicts. When `secondary` is None (no second source was fetched for this call -- the default, `Settings.PROVIDER_SECONDARY = "NONE"`), this degrades to exactly "use primary's value for every field primary provided", byte-identical in output to what single-source ingestion already does. There is deliberately no behavior change for that, the common, case -- dual- source resolution only activates when a caller actually has two dicts to pass in. """ resolved: dict = {} conflicts: list[FieldConflict] = [] for field in fields: val_a = primary.get(field) val_b = secondary.get(field) if secondary is not None else None if val_a is None and val_b is None: continue candidates = [(primary_tier, val_a)] if secondary is not None and secondary_tier is not None: candidates.append((secondary_tier, val_b)) result = resolve_conflict(candidates) if result.selected_value is not None: resolved[field] = result.selected_value # Secondary-only provenance must be observable. A NULL in the # primary provider followed by a value from the secondary provider # is not a numeric conflict, but it is still a cross-provider merge. if ( secondary is not None and secondary_tier is not None and val_a is None and val_b is not None ): conflicts.append( FieldConflict( field, val_a, val_b, result.selected_value, "Primary field missing; value supplied by secondary source", ) ) elif result.is_conflicting: conflicts.append( FieldConflict( field, val_a, val_b, result.selected_value, result.selection_reason, ) ) return resolved, conflicts def merge_statement_line_items( income: dict, balance: Optional[dict] = None, cash_flow: Optional[dict] = None, ) -> dict: """Merge one provider's three statement responses for the SAME reporting period into the one flat line-items dict `ingest_security()` writes to the four statement tables. AUDIT FIX (StockLab final engineering pass, Part A9 — a critical bug found while wiring TTM). Before this existed, `ingest_security()` called ONLY `adapter.get_income_statements()` and then wrote `BalanceSheet`, `CashFlow` and `Shares` rows out of that single response's line items. `adapter.get_balance_sheets()` and `adapter.get_cash_flows()` were never called anywhere in the application (verified by grep across `app/` and `scripts/`: the only remaining matches were comments). Against `DemoDataAdapter` this was invisible, because its three statement methods all return the SAME fully-merged dict containing every field. Against `FMPAdapter` — the real primary provider — `get_income_statements()` maps only `_INCOME_MAP`, so every balance-sheet, cash-flow and share-count column would have been written NULL for real data. See docs/AUDIT_INGEST_STATEMENTS_A9.md. Merge rule: first non-None value wins, in income -> balance -> cash-flow order. The income statement is therefore authoritative for fields more than one endpoint reports (FMP's cash-flow response also carries `netIncome`, for example), and a `None` from a later statement can never erase a real value from an earlier one. """ merged: dict = {} for source in (income, balance, cash_flow): if not source: continue for key, value in source.items(): if value is None: continue if merged.get(key) is None: merged[key] = value return merged ============================================================ 7. VALIDATION — FINAL PATCHED CONTRACT ============================================================ """ Shared provider-response validation (StockLab final engineering pass, Part D). ## Why this module exists Writing `tests/test_eodhd_adapter.py` in Part A7 pinned four behaviours of the adapter layer that are individually defensible and collectively incoherent: 1. `get_prices()` does **no type validation**: a JSON body with `"close": "not-a-number"` is passed straight through into `ProviderPriceBar.close` as a `str`. It travels all the way to a metric calculation before anything notices. 2. `get_prices()` fails the **whole batch** on one malformed bar: the list comprehension raises `KeyError` on a row missing `"date"`, so two well-formed bars are discarded with the third. 3. `list_universe()` does the **opposite** — silently drops rows without a `Code` key and returns the rest, with no record that anything was dropped. 4. Nothing anywhere checks whether a parsed value is *possible*: a negative price, a `high` below its `low`, a `close` outside the day's range, a negative share count. So the same layer both over-reacts and under-reacts, and in neither case does the caller learn what happened. Part A7 deliberately left this alone — fixing one adapter method while its counterpart stayed inconsistent would have been worse than the gap — and recorded it as Part D's job. This module is that job. ## The policy this module implements **Reject the row, not the batch, and always say so.** A malformed bar is dropped with a recorded `ValidationIssue` naming the field, the offending value and the reason. The caller gets the good rows AND a `ValidationReport` it can log, count, or refuse to proceed on. That is strictly better than both existing behaviours: no silent data loss, and no batch thrown away over one bad row. **Coerce narrowly, never creatively.** A numeric string becomes a float, because providers really do return `"12.5"`. Everything else — `None`, `""`, `"n/a"`, a bool, a list, `NaN`, `Infinity` — becomes `None` with an issue recorded. `True` is explicitly not `1.0`: Python would happily do that arithmetic and produce a silently wrong number. **Sanity-check what is checkable, and nothing else.** A negative price is impossible; a `high` below a `low` is impossible; a `close` outside `[low, high]` is impossible. Those are checked. What a "reasonable" P/E or revenue growth is, is a judgment about companies, not about data integrity — this module does not have opinions about those, deliberately. Pure, dependency-free (stdlib only), so it is genuinely testable in this environment. """ from __future__ import annotations import math from dataclasses import dataclass, field from typing import Any, Optional SEVERITY_ERROR = "ERROR" # the row cannot be used SEVERITY_WARNING = "WARNING" # the row is usable but something is off @dataclass(frozen=True) class ValidationIssue: field_name: str value: Any reason: str severity: str = SEVERITY_ERROR row_key: Optional[str] = None # e.g. the bar's date, so an issue can be traced to a row def __str__(self) -> str: # pragma: no cover - logging aid where = f" [{self.row_key}]" if self.row_key else "" return f"{self.severity}{where} {self.field_name}={self.value!r}: {self.reason}" @dataclass class ValidationReport: """Accumulates issues across a batch. Never raises — the caller decides what to do.""" issues: list[ValidationIssue] = field(default_factory=list) rows_seen: int = 0 rows_accepted: int = 0 def add(self, issue: ValidationIssue) -> None: self.issues.append(issue) @property def rows_rejected(self) -> int: return self.rows_seen - self.rows_accepted @property def errors(self) -> list[ValidationIssue]: return [i for i in self.issues if i.severity == SEVERITY_ERROR] @property def warnings(self) -> list[ValidationIssue]: return [i for i in self.issues if i.severity == SEVERITY_WARNING] @property def ok(self) -> bool: return not self.errors @property def rejection_rate(self) -> float: """0.0-1.0. A caller should treat a high rate as a provider/mapping problem, not as data: losing 40% of a price history silently is how a backtest ends up quietly wrong.""" return (self.rows_rejected / self.rows_seen) if self.rows_seen else 0.0 def summary(self) -> str: return ( f"{self.rows_accepted}/{self.rows_seen} rows accepted, " f"{len(self.errors)} errors, {len(self.warnings)} warnings" ) def coerce_float( value: Any, field_name: str, report: Optional[ValidationReport] = None, row_key: Optional[str] = None, required: bool = False, ) -> Optional[float]: """Return `value` as a float, or `None` with an issue recorded. Accepts int, float and numeric strings (providers really do return `"12.5"`). Rejects everything else, including `bool` — `True` would otherwise become `1.0` and produce a silently wrong number that no downstream check could catch. """ if value is None or value == "": if required and report is not None: report.add(ValidationIssue(field_name, value, "required field is missing", row_key=row_key)) return None if isinstance(value, bool): if report is not None: report.add(ValidationIssue(field_name, value, "boolean is not a number", row_key=row_key)) return None if isinstance(value, (int, float)): number = float(value) elif isinstance(value, str): try: number = float(value.strip()) except ValueError: if report is not None: report.add(ValidationIssue(field_name, value, "not parseable as a number", row_key=row_key)) return None else: if report is not None: report.add(ValidationIssue(field_name, value, f"unsupported type {type(value).__name__}", row_key=row_key)) return None if math.isnan(number) or math.isinf(number): if report is not None: report.add(ValidationIssue(field_name, value, "NaN/Infinity is not a usable value", row_key=row_key)) return None return number def validate_price_bar( raw: dict, report: ValidationReport, date_key: str = "date", row_key: Optional[str] = None, ) -> Optional[dict]: """Validate one OHLCV row. Returns a clean dict of floats, or `None` if the row is unusable. `raw` uses this codebase's canonical keys (`open`/`high`/`low`/`close`/`adjusted_close`/ `volume`); an adapter maps its provider's names before calling. The row is REJECTED when: the date is missing, `close` is missing or unparseable, or `close` is not strictly positive. Everything else produces a WARNING and a `None` for that field — an unusable `volume` is not a reason to throw away a valid price. """ report.rows_seen += 1 key = row_key or str(raw.get(date_key)) if not raw.get(date_key): report.add(ValidationIssue(date_key, raw.get(date_key), "row has no date", row_key=key)) return None close = coerce_float(raw.get("close"), "close", report, key, required=True) if close is None: return None if close <= 0: report.add(ValidationIssue("close", close, "price must be strictly positive", row_key=key)) return None out: dict = {date_key: raw[date_key], "close": close} for name in ("open", "high", "low", "adjusted_close"): value = coerce_float(raw.get(name), name, report, key) if value is not None and value <= 0: report.add(ValidationIssue(name, value, "price must be strictly positive", SEVERITY_WARNING, key)) value = None out[name] = value volume = coerce_float(raw.get("volume"), "volume", report, key) if volume is not None and volume < 0: report.add(ValidationIssue("volume", volume, "volume cannot be negative", SEVERITY_WARNING, key)) volume = None out["volume"] = volume high, low = out["high"], out["low"] if high is not None and low is not None and high < low: # Impossible, and it means the two are swapped or mismapped. Both are dropped rather than # guessing which one is wrong; `close` is still usable, which is what matters downstream. report.add(ValidationIssue("high/low", (high, low), "high is below low", SEVERITY_WARNING, key)) out["high"] = out["low"] = high = low = None if high is not None and close > high: report.add(ValidationIssue("close", close, f"close is above high ({high})", SEVERITY_WARNING, key)) if low is not None and close < low: report.add(ValidationIssue("close", close, f"close is below low ({low})", SEVERITY_WARNING, key)) report.rows_accepted += 1 return out #: Line items that cannot be negative in any real filing. Kept deliberately short: every entry is #: an arithmetic impossibility, not a judgment about what a healthy company looks like. Equity, #: net income, operating income, FCF and retained earnings are all legitimately negative and are #: NOT here. NON_NEGATIVE_LINE_ITEMS = ( "revenue", "total_assets", "current_assets", "current_liabilities", "cash_and_equivalents", "short_term_investments", "total_debt", "short_term_debt", "long_term_debt", "lease_liabilities", "goodwill", "intangible_assets", "inventory", "receivables", "shares_outstanding", "diluted_shares", ) def validate_line_items( line_items: dict, report: ValidationReport, row_key: Optional[str] = None, ) -> dict: """Coerce and sanity-check one period's line items. Returns the cleaned dict. Unlike a price bar, a financial period is never rejected wholesale: one impossible field does not invalidate the other thirty, and a period with a bad `inventory` still supports every metric that does not use inventory. Bad fields are dropped individually with an issue recorded. """ report.rows_seen += 1 out: dict = {} for name, value in line_items.items(): number = coerce_float(value, name, report, row_key) if number is None: continue if name in NON_NEGATIVE_LINE_ITEMS and number < 0: report.add(ValidationIssue(name, number, "cannot be negative in a real filing", SEVERITY_WARNING, row_key)) continue out[name] = number # Cross-field checks: each is an identity that must hold, not a heuristic. revenue = out.get("revenue") cogs = out.get("cogs") gross_profit = out.get("gross_profit") # Accounting identity: # gross_profit = revenue - cogs # # The provider resolver already treats <=0.5% financial differences # as non-conflicting. Use the same tolerance here so ordinary provider # rounding does not create a false integrity failure. if revenue is not None and cogs is not None and gross_profit is not None: expected_gross_profit = revenue - cogs if expected_gross_profit == 0: inconsistent = abs(gross_profit) > 0.01 else: inconsistent = ( abs(gross_profit - expected_gross_profit) / abs(expected_gross_profit) > 0.005 ) if inconsistent: report.add( ValidationIssue( "gross_profit", gross_profit, ( f"does not reconcile with revenue ({revenue}) " f"- cogs ({cogs}) = {expected_gross_profit}" ), SEVERITY_WARNING, row_key, ) ) # Do not allow an internally inconsistent dependent value # into the persisted merged statement. del out["gross_profit"] gross_profit = None if revenue is not None and gross_profit is not None and gross_profit > revenue: report.add(ValidationIssue("gross_profit", gross_profit, f"exceeds revenue ({revenue})", SEVERITY_WARNING, row_key)) total_assets, current_assets = out.get("total_assets"), out.get("current_assets") if total_assets is not None and current_assets is not None and current_assets > total_assets: report.add(ValidationIssue("current_assets", current_assets, f"exceeds total assets ({total_assets})", SEVERITY_WARNING, row_key)) diluted, outstanding = out.get("diluted_shares"), out.get("shares_outstanding") if diluted is not None and outstanding is not None and diluted < outstanding: # Diluted counts every share that COULD exist, so it can never be below basic/outstanding. report.add(ValidationIssue("diluted_shares", diluted, f"below shares outstanding ({outstanding})", SEVERITY_WARNING, row_key)) report.rows_accepted += 1 return out ============================================================ 8. DATABASE SESSION ============================================================ """SQLAlchemy engine/session setup. FastAPI dependency `get_db` yields one session per request. AUDIT (StockLab overhaul, Part 25 — PostgreSQL production optimization): connection pooling was previously hardcoded to `pool_pre_ping=True` only, with SQLAlchemy's library defaults for everything else (`pool_size=5`, `max_overflow=10`, no `pool_recycle`, `pool_timeout=30`) — workable for local dev but not documented or tunable for a small production server. Now sourced from `Settings` (`DB_POOL_SIZE`/`DB_MAX_OVERFLOW`/`DB_POOL_TIMEOUT_SECONDS`/`DB_POOL_RECYCLE_SECONDS`/ `DB_POOL_PRE_PING`, all env-overridable) so the target Debian 13 / Podman deployment can size the pool for its actual RAM/connection budget without a code change — see docs/DEPLOYMENT.md §9. """ from __future__ import annotations from collections.abc import Generator from sqlalchemy import create_engine from sqlalchemy.orm import Session, sessionmaker from app.core.config import get_settings settings = get_settings() engine = create_engine( settings.DATABASE_URL, pool_pre_ping=settings.DB_POOL_PRE_PING, pool_size=settings.DB_POOL_SIZE, max_overflow=settings.DB_MAX_OVERFLOW, pool_timeout=settings.DB_POOL_TIMEOUT_SECONDS, pool_recycle=settings.DB_POOL_RECYCLE_SECONDS, future=True, ) SessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False, future=True) def get_db() -> Generator[Session, None, None]: db = SessionLocal() try: yield db finally: db.close() ============================================================ 9. CONFIGURATION — DATABASE / PROVIDERS / INGEST FLAGS ============================================================ """ Central application configuration (spec §57 — every threshold/weight configurable; §67 env management). Everything sensitive comes from the environment, never hard-coded, never committed (.env.example documents every variable with a safe placeholder). """ from __future__ import annotations from functools import lru_cache from typing import Literal, Optional from pydantic import model_validator from pydantic_settings import BaseSettings, SettingsConfigDict # AUDIT (StockLab overhaul, security audit -- JWT secret strength): the placeholder default below # is intentionally obvious/guessable so a forgotten override is loud in dev, not silently "secure # enough." _validate_production_safety() below refuses to start the app if this placeholder (or # anything shorter than a reasonable minimum) reaches production -- see that validator for why. _INSECURE_DEFAULT_JWT_SECRET = "CHANGE_ME_INSECURE_DEV_ONLY" _MIN_PRODUCTION_JWT_SECRET_LENGTH = 32 class Settings(BaseSettings): model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8", extra="ignore") ENVIRONMENT: Literal["development", "staging", "production"] = "development" DEBUG: bool = False # --- Database --- DATABASE_URL: str = "postgresql+psycopg://stocklab:stocklab@localhost:5432/stocklab" # --- Redis --- REDIS_URL: str = "redis://localhost:6379/0" # --- Auth --- JWT_SECRET_KEY: str = _INSECURE_DEFAULT_JWT_SECRET JWT_ALGORITHM: str = "HS256" ACCESS_TOKEN_EXPIRE_MINUTES: int = 30 REFRESH_TOKEN_EXPIRE_DAYS: int = 14 # --- Data providers (see docs/DATA_SOURCES.md) --- PROVIDER_PRIMARY: Literal["FMP", "EODHD", "DEMO"] = "DEMO" PROVIDER_SECONDARY: Literal["FMP", "EODHD", "NONE"] = "NONE" FMP_API_KEY: Optional[str] = None FMP_BASE_URL: str = "https://financialmodelingprep.com/api/v3" EODHD_API_KEY: Optional[str] = None EODHD_BASE_URL: str = "https://eodhd.com/api" PROVIDER_REQUEST_TIMEOUT_SECONDS: float = 20.0 PROVIDER_MAX_RETRIES: int = 3 # --- Reporting basis (final engineering pass, Part A9 — docs/AUDIT_TTM_A9.md) --- # docs/AUDIT_METRICS.md cross-cutting finding #2: every metric this platform labels "TTM" is # really "most recent completed fiscal year", because only annual periods are ingested and # build_snapshot_from_db() filters on period_type == "FY". These two settings are the opt-in # path to a real trailing-twelve-month basis. BOTH default to the existing behaviour, so a # deployment that changes nothing behaves exactly as before. INGEST_QUARTERLY_PERIODS: bool = False # When true, ingestion additionally fetches period="quarter" statements and stores them as # Q1..Q4 FinancialPeriod rows alongside the annual ones. Costs one extra provider call per # statement type per security, which is why it is opt-in. FINANCIAL_BASIS: Literal["ANNUAL", "TTM"] = "ANNUAL" # "TTM" makes build_snapshot_from_db() aggregate the four most recent quarters via # app/engines/ttm.py for the CURRENT period. Requires INGEST_QUARTERLY_PERIODS=true to have # any data to aggregate; with no usable four-quarter window the snapshot falls back to the # annual basis and the reason is logged rather than silently swallowed. # --- Valuation assumptions (ASSUMPTION-tagged, never invented per company — VALUATION.md §1) --- DEFAULT_RISK_FREE_RATE: float = 0.045 DEFAULT_EQUITY_RISK_PREMIUM: float = 0.045 DEFAULT_BETA_IF_MISSING: float = 1.0 MIN_COST_OF_DEBT: float = 0.01 # AUDIT (StockLab overhaul Part 9/11): these four were bare literals inline in # app/workers/recompute.py before this pass (tax_rate=0.21, growth fallback=0.04, # reinvestment fallback=0.05, terminal_growth=0.025) — moved to configurable settings so # they're visible/overridable rather than buried in worker code, per spec §57. DEFAULT_CORPORATE_TAX_RATE: float = 0.21 # used only when the metrics engine's own effective/fallback tax rate is unavailable DEFAULT_DCF_REVENUE_GROWTH_FALLBACK: float = 0.04 # used only when the security has no 3Y revenue CAGR DEFAULT_DCF_REINVESTMENT_RATE_FALLBACK: float = 0.05 # used only when capex/revenue is unavailable DEFAULT_DCF_TERMINAL_GROWTH: float = 0.025 # capped at DEFAULT_RISK_FREE_RATE by cap_terminal_growth_at_risk_free() # --- DCF trajectory / scenario derivation (final engineering pass, Part B1) --- DCF_GROWTH_FADE: bool = True # When true, the near-term revenue growth rate fades linearly toward DEFAULT_DCF_TERMINAL_GROWTH # over DCF_GROWTH_FADE_YEARS instead of being held flat for all 10 explicit years. Holding a # 3Y-CAGR-derived growth rate flat for a decade is the single least defensible assumption in # the previous DCF; "growth fades to the terminal rate" is the standard convention. DCF_GROWTH_FADE_YEARS: int = 10 DCF_SCENARIO_MODE: Literal["hybrid", "additive_pp", "multiplicative"] = "hybrid" # "hybrid" (default, Part B1): Bear/Bull shift by max(floor_pp, fraction x base_rate), which # measurably gives mature businesses a narrower fair-value band and hypergrowth a wider one. # "additive_pp": fixed percentage-point shifts. "multiplicative": the pre-B1 fixed-factor # behaviour, retained so a deployment can reproduce older valuations. The measured spreads for # all three are tabulated in docs/AUDIT_DCF_B1.md. # --- Database connection pooling (spec: PostgreSQL production optimization) --- DB_POOL_SIZE: int = 5 DB_MAX_OVERFLOW: int = 10 DB_POOL_TIMEOUT_SECONDS: int = 30 DB_POOL_RECYCLE_SECONDS: int = 1800 DB_POOL_PRE_PING: bool = True # --- Celery worker tuning (spec: Redis/Celery production review — Part 25) --- # Conservative defaults sized for a small single-server deployment (docs/DEPLOYMENT.md §9), # not a large cluster — override via env for a bigger box. CELERY_WORKER_CONCURRENCY: int = 2 CELERY_WORKER_PREFETCH_MULTIPLIER: int = 1 # don't let one worker hoard tasks off the queue CELERY_TASK_ACKS_LATE: bool = True # requeue a task if the worker dies mid-execution CELERY_TASK_TIME_LIMIT_SECONDS: int = 900 # hard kill after 15 min (a stuck task shouldn't run forever) CELERY_TASK_SOFT_TIME_LIMIT_SECONDS: int = 780 # soft warning 2 min before the hard limit # --- Scoring defaults (spec §17, overridable per screener) --- SCORE_WEIGHT_QUALITY: float = 0.25 SCORE_WEIGHT_FINANCIAL_HEALTH: float = 0.20 SCORE_WEIGHT_GROWTH: float = 0.20 SCORE_WEIGHT_COMPETITIVE_ADVANTAGE: float = 0.15 SCORE_WEIGHT_VALUATION: float = 0.20 MIN_PEER_GROUP_SIZE: int = 8 # --- Peer groups and industry reference multiples (final engineering pass, Part B2) --- INDUSTRY_MULTIPLE_MIN_GROUP_SIZE: int = 5 # An industry needs at least this many companies reporting a given multiple before its median # is published as an INDUSTRY_MEDIAN reference. A "median P/E" from two companies is not an # industry reference; industries below the threshold are simply absent and the caller keeps # its self-historical fallback. See docs/AUDIT_PEER_GROUPS_B2.md. MULTIPLES_REFERENCE_PREFERENCE: Literal["industry_then_self", "self_then_industry", "self_only"] = "industry_then_self" # Which reference anchors a multiples fair value when both are available. "self_only" # reproduces the pre-B2 behaviour exactly. # --- Margin of safety defaults (VALUATION.md §5) --- DEFAULT_STRONG_BUY_MOS: float = 0.35 DEFAULT_BUY_MOS: float = 0.20 DEFAULT_OVERVALUED_PREMIUM: float = 0.15 # --- Rate limiting (SECURITY.md) --- RATE_LIMIT_PER_MINUTE: int = 120 # AUDIT FIX (StockLab overhaul, final engineering pass, Part A3): three more, deliberately # separate from the general default above and from each other -- each protects a route with a # genuinely different risk/benefit shape (docs/AUDIT_SECURITY_A3.md has the full reasoning for # every route, including the ones left undecorated). All three configurable per spec §57, same # as RATE_LIMIT_PER_MINUTE. SCREENER_RATE_LIMIT_PER_MINUTE: int = 30 # Tighter than the general default: /v1/screeners/run is unauthenticated, and unlike a fixed # read its cost scales with the number of filters a caller supplies (each ScreenFilter compiles # to its own correlated EXISTS subquery, app/engines/screening/executor.py) -- a small number of # concurrent callers with maximal filter lists is a real DoS-shaped risk in a way a simple GET # isn't. SEARCH_RATE_LIMIT_PER_MINUTE: int = 30 # Same tier as screener, same reason: /v1/search is unauthenticated and its `ILIKE '%q%'` # leading-wildcard pattern (app/api/v1/search.py) cannot use a plain btree index -- a real, # already-documented scan-cost concern (docs/AUDIT_PERFORMANCE.md finding #4) that a rate limit # mitigates until pg_trgm is adopted. WATCHLIST_WRITE_RATE_LIMIT_PER_MINUTE: int = 60 # Looser than screener/search: watchlist add/remove is authenticated (get_current_user) and # scoped to the caller's own rows only -- the risk is DB write/commit churn from a scripted # caller, not an unauthenticated amplification vector, so a generous limit is enough. # --- CORS --- CORS_ALLOWED_ORIGINS: str = "http://localhost:3000" # --- Screener cache --- SCREENER_CACHE_TTL_SECONDS: int = 900 @property def cors_origins_list(self) -> list[str]: return [o.strip() for o in self.CORS_ALLOWED_ORIGINS.split(",") if o.strip()] @model_validator(mode="after") def _validate_production_safety(self) -> "Settings": """AUDIT FIX (StockLab overhaul, security audit): before this pass, ENVIRONMENT=production with a forgotten/default JWT_SECRET_KEY started successfully and silently signed tokens with a public, well-known string -- anyone could forge a valid access token for any user ID. This is a fail-closed startup check, not a runtime request check: it raises once, at process start, rather than degrading security silently. Only applies when ENVIRONMENT == "production" -- development/staging keep the friendly placeholder so local setup doesn't require generating a real secret first.""" if self.ENVIRONMENT == "production": if self.JWT_SECRET_KEY == _INSECURE_DEFAULT_JWT_SECRET: raise ValueError( "JWT_SECRET_KEY is still the insecure default placeholder. Generate a real " "secret before running in production, e.g.: " "python3 -c \"import secrets; print(secrets.token_urlsafe(48))\"" ) if len(self.JWT_SECRET_KEY) < _MIN_PRODUCTION_JWT_SECRET_LENGTH: raise ValueError( f"JWT_SECRET_KEY is only {len(self.JWT_SECRET_KEY)} characters -- production " f"requires at least {_MIN_PRODUCTION_JWT_SECRET_LENGTH}. Generate a real " "secret, e.g.: python3 -c \"import secrets; print(secrets.token_urlsafe(48))\"" ) if self.DEBUG: raise ValueError( "DEBUG=true is not allowed when ENVIRONMENT=production (would risk verbose " "error responses). Set DEBUG=false." ) return self @lru_cache def get_settings() -> Settings: return Settings() ============================================================ 10. FINANCIAL MODELS — COMPLETE ============================================================ """ Point-in-time financial statement tables (spec §12, §44 / docs/DATA_MODEL.md). Never updated in place — a restatement is a new row with a later `filing_date`, not a mutation of the original. `financial_periods` is the parent row; each statement table is 1:1 with it via `financial_period_id`. This split (rather than one wide table) keeps each statement's columns independently nullable/auditable and mirrors how providers deliver them. """ from __future__ import annotations from datetime import date from typing import Optional from sqlalchemy import Date, ForeignKey, String, UniqueConstraint from sqlalchemy.orm import Mapped, mapped_column, relationship from app.models.base import Base, TimestampMixin, uuid_pk class FinancialPeriod(Base, TimestampMixin): __tablename__ = "financial_periods" __table_args__ = ( UniqueConstraint("security_id", "period_end", "period_type", "filing_date", name="uq_financial_period"), ) id: Mapped[str] = uuid_pk() security_id: Mapped[str] = mapped_column(ForeignKey("securities.id"), index=True) period_end: Mapped[date] = mapped_column(Date, index=True) period_type: Mapped[str] = mapped_column(String(8)) # FY | Q1 | Q2 | Q3 | Q4 | TTM filing_date: Mapped[Optional[date]] = mapped_column(Date, nullable=True, index=True) published_at: Mapped[Optional[date]] = mapped_column(Date, nullable=True) retrieved_at: Mapped[Optional[date]] = mapped_column(Date, nullable=True) currency: Mapped[str] = mapped_column(String(3)) source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) methodology_version: Mapped[str] = mapped_column(String(16), default="v1") income_statement: Mapped[Optional["IncomeStatement"]] = relationship(back_populates="period", uselist=False) balance_sheet: Mapped[Optional["BalanceSheet"]] = relationship(back_populates="period", uselist=False) cash_flow: Mapped[Optional["CashFlow"]] = relationship(back_populates="period", uselist=False) shares: Mapped[Optional["Shares"]] = relationship(back_populates="period", uselist=False) def _fk_period(): return mapped_column(ForeignKey("financial_periods.id"), unique=True, index=True) class IncomeStatement(Base, TimestampMixin): __tablename__ = "income_statements" id: Mapped[str] = uuid_pk() financial_period_id: Mapped[str] = _fk_period() revenue: Mapped[Optional[float]] = mapped_column(nullable=True) cogs: Mapped[Optional[float]] = mapped_column(nullable=True) gross_profit: Mapped[Optional[float]] = mapped_column(nullable=True) operating_income: Mapped[Optional[float]] = mapped_column(nullable=True) ebit: Mapped[Optional[float]] = mapped_column(nullable=True) ebitda: Mapped[Optional[float]] = mapped_column(nullable=True) net_income: Mapped[Optional[float]] = mapped_column(nullable=True) eps_basic: Mapped[Optional[float]] = mapped_column(nullable=True) eps_diluted: Mapped[Optional[float]] = mapped_column(nullable=True) tax_expense: Mapped[Optional[float]] = mapped_column(nullable=True) pretax_income: Mapped[Optional[float]] = mapped_column(nullable=True) interest_expense: Mapped[Optional[float]] = mapped_column(nullable=True) depreciation_and_amortization: Mapped[Optional[float]] = mapped_column(nullable=True) stock_based_compensation: Mapped[Optional[float]] = mapped_column(nullable=True) adjusted_net_income: Mapped[Optional[float]] = mapped_column(nullable=True) period: Mapped["FinancialPeriod"] = relationship(back_populates="income_statement") class BalanceSheet(Base, TimestampMixin): __tablename__ = "balance_sheets" id: Mapped[str] = uuid_pk() financial_period_id: Mapped[str] = _fk_period() cash_and_equivalents: Mapped[Optional[float]] = mapped_column(nullable=True) short_term_investments: Mapped[Optional[float]] = mapped_column(nullable=True) total_debt: Mapped[Optional[float]] = mapped_column(nullable=True) short_term_debt: Mapped[Optional[float]] = mapped_column(nullable=True) long_term_debt: Mapped[Optional[float]] = mapped_column(nullable=True) lease_liabilities: Mapped[Optional[float]] = mapped_column(nullable=True) total_assets: Mapped[Optional[float]] = mapped_column(nullable=True) current_assets: Mapped[Optional[float]] = mapped_column(nullable=True) current_liabilities: Mapped[Optional[float]] = mapped_column(nullable=True) shareholders_equity: Mapped[Optional[float]] = mapped_column(nullable=True) minority_interest: Mapped[Optional[float]] = mapped_column(nullable=True) preferred_equity: Mapped[Optional[float]] = mapped_column(nullable=True) goodwill: Mapped[Optional[float]] = mapped_column(nullable=True) intangible_assets: Mapped[Optional[float]] = mapped_column(nullable=True) receivables: Mapped[Optional[float]] = mapped_column(nullable=True) inventory: Mapped[Optional[float]] = mapped_column(nullable=True) working_capital: Mapped[Optional[float]] = mapped_column(nullable=True) period: Mapped["FinancialPeriod"] = relationship(back_populates="balance_sheet") class CashFlow(Base, TimestampMixin): __tablename__ = "cash_flows" id: Mapped[str] = uuid_pk() financial_period_id: Mapped[str] = _fk_period() operating_cash_flow: Mapped[Optional[float]] = mapped_column(nullable=True) capital_expenditure: Mapped[Optional[float]] = mapped_column(nullable=True) free_cash_flow: Mapped[Optional[float]] = mapped_column(nullable=True) # cached; recomputed, not authoritative dividends_paid: Mapped[Optional[float]] = mapped_column(nullable=True) buybacks: Mapped[Optional[float]] = mapped_column(nullable=True) stock_issuance: Mapped[Optional[float]] = mapped_column(nullable=True) period: Mapped["FinancialPeriod"] = relationship(back_populates="cash_flow") class Shares(Base, TimestampMixin): __tablename__ = "shares" id: Mapped[str] = uuid_pk() financial_period_id: Mapped[str] = _fk_period() shares_outstanding: Mapped[Optional[float]] = mapped_column(nullable=True) diluted_shares: Mapped[Optional[float]] = mapped_column(nullable=True) period: Mapped["FinancialPeriod"] = relationship(back_populates="shares") class Dividend(Base, TimestampMixin): __tablename__ = "dividends" id: Mapped[str] = uuid_pk() security_id: Mapped[str] = mapped_column(ForeignKey("securities.id"), index=True) ex_date: Mapped[date] = mapped_column(Date, index=True) pay_date: Mapped[Optional[date]] = mapped_column(Date, nullable=True) amount_per_share: Mapped[float] = mapped_column() currency: Mapped[str] = mapped_column(String(3)) is_special: Mapped[bool] = mapped_column(default=False) source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) class Buyback(Base, TimestampMixin): __tablename__ = "buybacks" id: Mapped[str] = uuid_pk() security_id: Mapped[str] = mapped_column(ForeignKey("securities.id"), index=True) period_end: Mapped[date] = mapped_column(Date, index=True) gross_amount: Mapped[Optional[float]] = mapped_column(nullable=True) shares_repurchased: Mapped[Optional[float]] = mapped_column(nullable=True) source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) class Estimate(Base, TimestampMixin): """Forward analyst consensus estimates — always ESTIMATED, never VERIFIED.""" __tablename__ = "estimates" __table_args__ = (UniqueConstraint("security_id", "period_end", "metric", "as_of_date", name="uq_estimate"),) id: Mapped[str] = uuid_pk() security_id: Mapped[str] = mapped_column(ForeignKey("securities.id"), index=True) period_end: Mapped[date] = mapped_column(Date, index=True) metric: Mapped[str] = mapped_column(String(32)) # "eps" | "revenue" | ... consensus_value: Mapped[Optional[float]] = mapped_column(nullable=True) num_analysts: Mapped[Optional[int]] = mapped_column(nullable=True) as_of_date: Mapped[date] = mapped_column(Date, index=True) source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) class Guidance(Base, TimestampMixin): __tablename__ = "guidance" id: Mapped[str] = uuid_pk() security_id: Mapped[str] = mapped_column(ForeignKey("securities.id"), index=True) period_end: Mapped[date] = mapped_column(Date, index=True) metric: Mapped[str] = mapped_column(String(32)) low: Mapped[Optional[float]] = mapped_column(nullable=True) high: Mapped[Optional[float]] = mapped_column(nullable=True) issued_date: Mapped[date] = mapped_column(Date, index=True) revision_direction: Mapped[Optional[str]] = mapped_column(String(8), nullable=True) # UP | DOWN | MAINTAINED source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) ============================================================ 11. COMPANY / SECURITY MODELS — COMPLETE ============================================================ """Companies (the legal/economic entity) and Securities (1:N listed lines).""" from __future__ import annotations from datetime import date from typing import TYPE_CHECKING, Optional from sqlalchemy import Boolean, Date, ForeignKey, String, UniqueConstraint from sqlalchemy.orm import Mapped, mapped_column, relationship from app.models.base import Base, TimestampMixin, uuid_pk if TYPE_CHECKING: # pragma: no cover - import cycle guard # SQLAlchemy resolves these relationship targets from its own registry at mapper-configuration # time, so they are never imported at runtime; this block exists so static analysis can see # that the string annotations below refer to real classes. from app.models.reference import Country, Exchange, Industry, Sector class Company(Base, TimestampMixin): __tablename__ = "companies" id: Mapped[str] = uuid_pk() legal_name: Mapped[str] = mapped_column(String(256), index=True) display_name: Mapped[str] = mapped_column(String(128), index=True) country_id: Mapped[str] = mapped_column(ForeignKey("countries.id"), index=True) sector_id: Mapped[str] = mapped_column(ForeignKey("sectors.id"), index=True) industry_id: Mapped[str] = mapped_column(ForeignKey("industries.id"), index=True) website: Mapped[Optional[str]] = mapped_column(String(256), nullable=True) description: Mapped[Optional[str]] = mapped_column(nullable=True) is_active: Mapped[bool] = mapped_column(Boolean, default=True) country: Mapped["Country"] = relationship() sector: Mapped["Sector"] = relationship() industry: Mapped["Industry"] = relationship() securities: Mapped[list["Security"]] = relationship(back_populates="company") class Security(Base, TimestampMixin): """A listed line for a company — a company can have multiple (home listing + ADR, dual listing, secondary listing).""" __tablename__ = "securities" __table_args__ = (UniqueConstraint("ticker", "exchange_id", name="uq_security_ticker_exchange"),) id: Mapped[str] = uuid_pk() company_id: Mapped[str] = mapped_column(ForeignKey("companies.id"), index=True) exchange_id: Mapped[str] = mapped_column(ForeignKey("exchanges.id"), index=True) ticker: Mapped[str] = mapped_column(String(32), index=True) isin: Mapped[Optional[str]] = mapped_column(String(12), nullable=True, index=True) cusip: Mapped[Optional[str]] = mapped_column(String(9), nullable=True) security_type: Mapped[str] = mapped_column(String(16), default="COMMON") # COMMON | ADR | PREFERRED | ETF currency: Mapped[str] = mapped_column(String(3)) is_primary_listing: Mapped[bool] = mapped_column(Boolean, default=True) listing_date: Mapped[Optional[date]] = mapped_column(Date, nullable=True) delisting_date: Mapped[Optional[date]] = mapped_column(Date, nullable=True) status: Mapped[str] = mapped_column(String(16), default="ACTIVE") # ACTIVE | DELISTED | SUSPENDED discovery_status: Mapped[Optional[str]] = mapped_column(String(24), nullable=True) # DISCOVERED | UNDER_REVIEW | QUALIFIED | WATCHLIST | INVESTABLE | REJECTED (spec §6) company: Mapped["Company"] = relationship(back_populates="securities") exchange: Mapped["Exchange"] = relationship() ============================================================ 12. GOVERNANCE MODELS — Source / DataQuality ============================================================ """Data governance: provider/filing registry, data-quality/conflict tracking, users, audit log.""" from __future__ import annotations from datetime import datetime from typing import Optional from sqlalchemy import JSON, Boolean, DateTime, ForeignKey, String, func from sqlalchemy.orm import Mapped, mapped_column from app.models.base import Base, TimestampMixin, uuid_pk class Source(Base, TimestampMixin): """One row per (provider/filing) ingestion event — the top of the source hierarchy chain (docs/DATA_SOURCES.md §4) and the archive of the raw payload for replay (§44).""" __tablename__ = "sources" id: Mapped[str] = uuid_pk() provider: Mapped[str] = mapped_column(String(32), index=True) # "FMP" | "EODHD" | "SEC_EDGAR" | "DEMO" provider_tier: Mapped[str] = mapped_column(String(16)) # OFFICIAL_FILING | PRIMARY | SECONDARY | CALCULATED is_demo: Mapped[bool] = mapped_column(Boolean, default=False, index=True) # DATA_SOURCES.md §9 — hard separation retrieved_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) endpoint: Mapped[Optional[str]] = mapped_column(String(256), nullable=True) raw_payload: Mapped[Optional[dict]] = mapped_column(JSON, nullable=True) # deliberate JSON exception, DATA_MODEL.md content_hash: Mapped[Optional[str]] = mapped_column(String(64), nullable=True, index=True) class DataQuality(Base, TimestampMixin): """Per-data-point quality/conflict record (spec §9).""" __tablename__ = "data_quality" id: Mapped[str] = uuid_pk() entity_table: Mapped[str] = mapped_column(String(64)) # e.g. "income_statements" entity_id: Mapped[str] = mapped_column(String(64), index=True) field: Mapped[str] = mapped_column(String(64)) status: Mapped[str] = mapped_column(String(16)) # VERIFIED|CALCULATED|ESTIMATED|ASSUMPTION|MISSING|STALE|CONFLICTING source_a_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) value_a: Mapped[Optional[float]] = mapped_column(nullable=True) source_b_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) value_b: Mapped[Optional[float]] = mapped_column(nullable=True) selected_value: Mapped[Optional[float]] = mapped_column(nullable=True) selected_source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) selection_reason: Mapped[Optional[str]] = mapped_column(nullable=True) class RevokedToken(Base): """Revoked/rotated refresh-token store (StockLab overhaul, final engineering pass, Part A4). Deliberately scoped to refresh tokens only -- access tokens are short-lived (Settings.ACCESS_TOKEN_EXPIRE_MINUTES, 30 min default) and validated on every authenticated request (app/api/v1/deps.py::get_current_user); adding a DB round-trip to that path for every request would be the "incompatible auth rewrite" this pass's instructions said to avoid, for a marginal gain given how short-lived access tokens already are. Refresh tokens are long-lived (Settings.REFRESH_TOKEN_EXPIRE_DAYS, 14 days default) and used rarely (once per session renewal) -- the theft window that actually matters, and where a DB check is cheap relative to how often it runs. No FK constraint enforced on user_id on purpose: a token can outlive its user row in edge cases (revoke-then-delete-account ordering), and this table's only real purpose is jti lookup, not joining back to User -- keeping it nullable/unconstrained avoids a fragile cross-table dependency for a security-relevant write path that should never itself fail to write. """ __tablename__ = "revoked_tokens" id: Mapped[str] = uuid_pk() jti: Mapped[str] = mapped_column(String(64), unique=True, index=True) user_id: Mapped[Optional[str]] = mapped_column(String(64), nullable=True, index=True) token_type: Mapped[str] = mapped_column(String(16), default="refresh") # only "refresh" is ever written today revoked_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) # Copy of the original token's own `exp` claim -- lets prune_expired_revoked_tokens() (see # app/core/token_revocation.py) delete rows once decode_token() would reject that jti's token # on expiry alone anyway, so this table doesn't grow forever. expires_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), index=True) reason: Mapped[Optional[str]] = mapped_column(String(32), nullable=True) # "logout" | "rotated" | None class User(Base, TimestampMixin): __tablename__ = "users" id: Mapped[str] = uuid_pk() email: Mapped[str] = mapped_column(String(256), unique=True, index=True) hashed_password: Mapped[str] = mapped_column(String(256)) display_name: Mapped[Optional[str]] = mapped_column(String(128), nullable=True) is_active: Mapped[bool] = mapped_column(Boolean, default=True) is_admin: Mapped[bool] = mapped_column(Boolean, default=False) class AuditLogEntry(Base): __tablename__ = "audit_log" id: Mapped[str] = uuid_pk() occurred_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now(), index=True) actor_user_id: Mapped[Optional[str]] = mapped_column(ForeignKey("users.id"), nullable=True) event_type: Mapped[str] = mapped_column(String(64), index=True) # LOGIN | LOGIN_FAILED | TOKEN_REFRESH | LOGOUT | ... detail: Mapped[Optional[dict]] = mapped_column(JSON, nullable=True) # never secrets — see SECURITY.md ip_address: Mapped[Optional[str]] = mapped_column(String(64), nullable=True) ============================================================ 13. EXISTING PROVIDER-RESILIENCE TEST CONTRACT ============================================================ """ Provider resilience audit tests (StockLab overhaul, Part 23): source-hierarchy conflict resolution (real, tested — but see docs/AUDIT_PROVIDER_RESILIENCE.md for the finding that nothing in this codebase currently calls it with two real sources) and the retry/backoff policy added to `app/workers/ingest.py::ingest_security_task` this pass (not directly testable here — Celery itself isn't installed in this sandbox, see the docs — so this file covers what IS testable without Celery/DB: the conflict-resolution engine itself). """ from __future__ import annotations from app.adapters.resolver import merge_statement_line_items, resolve_conflict, resolve_line_items def test_resolve_conflict_picks_highest_priority_source(): result = resolve_conflict([("SECONDARY", 100.0), ("OFFICIAL_FILING", 100.3), ("PRIMARY", 100.1)]) assert result.selected_value == 100.3 assert result.selected_source_tier == "OFFICIAL_FILING" assert result.is_conflicting is False # all three within 0.5% of the selected value def test_resolve_conflict_flags_real_disagreement(): # OFFICIAL_FILING says 100, PRIMARY says 120 -> 20% apart, well over the 0.5% tolerance. result = resolve_conflict([("OFFICIAL_FILING", 100.0), ("PRIMARY", 120.0)]) assert result.selected_value == 100.0 assert result.is_conflicting is True def test_resolve_conflict_tolerates_small_rounding_differences(): # 100.0 vs 100.3 = 0.3% apart, under the default 0.5% tolerance -> not flagged as conflicting. result = resolve_conflict([("OFFICIAL_FILING", 100.0), ("PRIMARY", 100.3)]) assert result.is_conflicting is False def test_resolve_conflict_handles_zero_value_without_dividing_by_zero(): result = resolve_conflict([("OFFICIAL_FILING", 0.0), ("PRIMARY", 0.0)]) assert result.is_conflicting is False result2 = resolve_conflict([("OFFICIAL_FILING", 0.0), ("PRIMARY", 5.0)]) assert result2.is_conflicting is True # 0 vs non-zero is a real disagreement, not a /0 crash def test_resolve_conflict_no_sources_available(): result = resolve_conflict([("PRIMARY", None)]) assert result.selected_value is None assert result.is_conflicting is False def test_resolve_conflict_unknown_tier_sorts_last(): result = resolve_conflict([("SOME_UNKNOWN_TIER", 999.0), ("SECONDARY", 50.0)]) assert result.selected_value == 50.0 # SECONDARY, a known tier, wins over an unranked one # --- app/adapters/resolver.py::resolve_line_items (StockLab overhaul, Part A6) ------------------- # The statement-level (multi-field) counterpart to resolve_conflict() above -- see that function's # docstring for the full design. These tests are what makes A6's "real, tested" claim honest: the # actual ingest.py wiring that CALLS this (fetching two live adapters, writing DataQuality rows) # cannot be executed in this sandbox (no network, no sqlalchemy) and is DOCUMENTED / NOT TESTED -- # see docs/AUDIT_PROVIDER_CONFLICT_A6.md -- but this pure resolution logic has real coverage. _FIELDS = ("revenue", "net_income", "ebitda") def test_resolve_line_items_no_secondary_passes_through_primary_unchanged(): # the default case: PROVIDER_SECONDARY=NONE -- must be byte-identical to single-source ingestion primary = {"revenue": 100.0, "net_income": 10.0} resolved, conflicts = resolve_line_items(primary, None, _FIELDS, primary_tier="PRIMARY") assert resolved == {"revenue": 100.0, "net_income": 10.0} assert conflicts == [] def test_resolve_line_items_field_absent_from_both_sources_is_omitted(): primary = {"revenue": 100.0} # net_income, ebitda absent from both secondary = {"revenue": 100.0} resolved, conflicts = resolve_line_items(primary, secondary, _FIELDS, "PRIMARY", "SECONDARY") assert "net_income" not in resolved assert "ebitda" not in resolved assert resolved == {"revenue": 100.0} def test_resolve_line_items_agreeing_sources_no_conflicts(): primary = {"revenue": 100.0, "net_income": 10.0, "ebitda": 20.0} secondary = {"revenue": 100.2, "net_income": 10.0, "ebitda": 20.05} # all within 0.5% tolerance resolved, conflicts = resolve_line_items(primary, secondary, _FIELDS, "PRIMARY", "SECONDARY") assert resolved == {"revenue": 100.0, "net_income": 10.0, "ebitda": 20.0} # primary wins tier order assert conflicts == [] def test_resolve_line_items_flags_real_disagreement_on_one_field_only(): primary = {"revenue": 100.0, "net_income": 10.0} secondary = {"revenue": 100.0, "net_income": 50.0} # net_income disagrees by 400% resolved, conflicts = resolve_line_items(primary, secondary, _FIELDS, "PRIMARY", "SECONDARY") assert resolved["net_income"] == 10.0 # PRIMARY still wins tier order despite the conflict assert len(conflicts) == 1 assert conflicts[0].field == "net_income" assert conflicts[0].value_a == 10.0 and conflicts[0].value_b == 50.0 def test_resolve_line_items_official_filing_secondary_overrides_primary_tier(): # if the "secondary" adapter is actually the higher-priority tier for this field (e.g. an # OFFICIAL_FILING source plugged in as `secondary`), its value must still win -- tier order, # not argument position, decides. primary = {"revenue": 100.0} # tier PRIMARY secondary = {"revenue": 105.0} # tier OFFICIAL_FILING -- higher priority resolved, conflicts = resolve_line_items(primary, secondary, ("revenue",), "PRIMARY", "OFFICIAL_FILING") assert resolved["revenue"] == 105.0 def test_resolve_line_items_one_source_missing_a_field_other_has_it(): primary = {"revenue": 100.0} # net_income missing from primary entirely secondary = {"revenue": 100.0, "net_income": 10.0} resolved, conflicts = resolve_line_items(primary, secondary, _FIELDS, "PRIMARY", "SECONDARY") assert resolved["net_income"] == 10.0 assert len(conflicts) == 1 assert conflicts[0].field == "net_income" assert conflicts[0].value_a is None assert conflicts[0].value_b == 10.0 assert conflicts[0].selected_value == 10.0 assert conflicts[0].selection_reason == "Primary field missing; value supplied by secondary source" # --- Part A9: merge_statement_line_items (the critical ingestion bug fix) --- def test_merge_combines_all_three_statements(): """The bug this guards: before Part A9, only the income dict reached the BalanceSheet/CashFlow tables, so every balance-sheet and cash-flow column was NULL for FMP-sourced data.""" merged = merge_statement_line_items( {"revenue": 1000.0, "net_income": 100.0}, {"total_assets": 5000.0, "shareholders_equity": 2000.0}, {"operating_cash_flow": 150.0, "capital_expenditure": 40.0}, ) assert merged == { "revenue": 1000.0, "net_income": 100.0, "total_assets": 5000.0, "shareholders_equity": 2000.0, "operating_cash_flow": 150.0, "capital_expenditure": 40.0, } def test_merge_income_statement_wins_on_a_shared_field(): """FMP's cash-flow response also carries netIncome. The income statement is authoritative.""" merged = merge_statement_line_items( {"net_income": 100.0}, None, {"net_income": 999.0, "operating_cash_flow": 150.0}, ) assert merged["net_income"] == 100.0 assert merged["operating_cash_flow"] == 150.0 def test_merge_later_none_never_erases_an_earlier_real_value(): merged = merge_statement_line_items({"net_income": 100.0}, {"net_income": None}, None) assert merged["net_income"] == 100.0 def test_merge_fills_a_field_the_income_statement_left_none(): merged = merge_statement_line_items({"net_income": None}, {"net_income": 42.0}, None) assert merged["net_income"] == 42.0 def test_merge_drops_none_valued_keys_entirely(): """A key present with a None value must not appear in the merged dict -- resolve_line_items treats an absent field and a None field differently only in that a None-valued key would otherwise be offered as a candidate.""" merged = merge_statement_line_items({"revenue": 1000.0, "cogs": None}, None, None) assert "cogs" not in merged def test_merge_with_only_an_income_statement_is_identity(): """The DEMO adapter returns the same fully-merged dict from all three statement methods, so the Part A9 fix must be a no-op for demo data. This is the degenerate form of that.""" income = {"revenue": 1000.0, "total_assets": 5000.0, "operating_cash_flow": 150.0} assert merge_statement_line_items(income, None, None) == income def test_merge_of_three_identical_dicts_is_idempotent(): """The exact DEMO shape: get_income_statements/get_balance_sheets/get_cash_flows all return the same dict. Merging it with itself twice must return it unchanged -- proving the ingestion fix cannot alter demo-mode output.""" d = {"revenue": 1000.0, "total_assets": 5000.0, "operating_cash_flow": 150.0, "diluted_shares": 100.0, "net_income": 90.0} assert merge_statement_line_items(d, dict(d), dict(d)) == d def test_merge_of_all_empty_inputs_is_empty(): assert merge_statement_line_items({}, None, None) == {} ALL_TESTS = [obj for name, obj in list(globals().items()) if name.startswith("test_") and callable(obj)] if __name__ == "__main__": passed, failed = 0, [] for fn in ALL_TESTS: try: fn() passed += 1 print(f"PASS {fn.__name__}") except AssertionError as e: failed.append(fn.__name__) print(f"FAIL {fn.__name__}: {e}") print(f"\n{passed}/{len(ALL_TESTS)} passed") if failed: raise SystemExit(1) ============================================================ 14. DEMO TEST CONTRACT ============================================================ """Sanity tests for DemoDataAdapter — internal consistency of the synthetic series, and that every row it produces is traceable back to demo status (spec: never blend demo with production).""" from __future__ import annotations from app.adapters.demo import DEMO_SEED_PROFILES, DemoDataAdapter, YEARS_OF_HISTORY from app.adapters.base import ProviderNotFoundError adapter = DemoDataAdapter() def test_unknown_ticker_raises_not_found(): try: adapter.get_company_profile("NOT_A_REAL_DEMO_TICKER") raised = False except ProviderNotFoundError: raised = True assert raised def test_income_statements_have_full_history_and_are_deterministic(): a = adapter.get_income_statements("AAPL") b = adapter.get_income_statements("AAPL") assert len(a) == YEARS_OF_HISTORY assert [p.line_items["revenue"] for p in a] == [p.line_items["revenue"] for p in b] # deterministic seed # most-recent-first assert a[0].period_end > a[-1].period_end def test_gross_profit_equals_revenue_minus_cogs_when_present(): periods = adapter.get_income_statements("MSFT") for p in periods: li = p.line_items if li.get("cogs") is not None and li.get("gross_profit") is not None: assert abs(li["revenue"] - li["cogs"] - li["gross_profit"]) < 0.2 # independent per-field rounding to 1dp def test_bank_has_no_cogs_gross_margin_not_meaningful_upstream(): periods = adapter.get_income_statements("JPM") assert all(p.line_items.get("cogs") is None for p in periods) def test_every_seed_profile_is_reachable_and_profile_matches(): for seed in DEMO_SEED_PROFILES: profile = adapter.get_company_profile(seed.ticker) assert profile.ticker == seed.ticker assert "DEMO" in (profile.description or "") def test_list_universe_filters_by_country(): us_tickers = adapter.list_universe(country_iso2="US") assert "AAPL" in us_tickers assert "SAP" not in us_tickers def test_demo_statement_methods_all_return_the_same_merged_line_items(): """Part A9 evidence. DemoDataAdapter.get_income_statements/get_balance_sheets/get_cash_flows all return ONE fully-merged dict per period. That is exactly why the ingestion bug fixed in Part A9 -- writing BalanceSheet/CashFlow rows out of the income response -- was invisible in demo mode and would have been catastrophic against FMP. Pinned so the demo adapter cannot quietly diverge and start hiding the same class of bug again.""" from app.adapters.demo import DemoDataAdapter a = DemoDataAdapter() ticker = sorted(a.list_universe())[0] inc = {(p.period_end, p.period_type): p.line_items for p in a.get_income_statements(ticker)} bal = {(p.period_end, p.period_type): p.line_items for p in a.get_balance_sheets(ticker)} cf = {(p.period_end, p.period_type): p.line_items for p in a.get_cash_flows(ticker)} assert inc and inc.keys() == bal.keys() == cf.keys() for key in inc: assert inc[key] == bal[key] == cf[key] def test_demo_income_response_already_carries_balance_and_cash_flow_fields(): from app.adapters.demo import DemoDataAdapter a = DemoDataAdapter() ticker = sorted(a.list_universe())[0] li = a.get_income_statements(ticker)[0].line_items for field in ("total_assets", "total_debt", "shareholders_equity", "operating_cash_flow", "capital_expenditure", "diluted_shares"): assert li.get(field) is not None, f"{field} missing from demo income response" def test_demo_total_assets_is_populated_for_every_seed_profile(): """Part A9 finding: the demo generator never emitted total_assets, so ROA, asset turnover, the goodwill-concentration flag and the acquisition-driven-growth flag were permanently skipped in demo mode -- the mode this platform runs in by default.""" from app.adapters.demo import DemoDataAdapter a = DemoDataAdapter() for ticker in a.list_universe(): for period in a.get_income_statements(ticker): assert period.line_items.get("total_assets") is not None, ticker def test_demo_balance_sheet_balances(): """total_assets is derived from the liabilities-and-equity side precisely so the synthetic balance sheet is internally consistent rather than an independently invented number.""" from app.adapters.demo import DemoDataAdapter a = DemoDataAdapter() for ticker in a.list_universe(): li = a.get_income_statements(ticker)[0].line_items implied = ( li["shareholders_equity"] + li["total_debt"] + li["current_liabilities"] + li["revenue"] * 0.05 ) assert abs(li["total_assets"] - implied) < max(1.0, abs(implied) * 1e-6), ticker def test_demo_debt_maturity_split_sums_to_total_debt(): from app.adapters.demo import DemoDataAdapter a = DemoDataAdapter() for ticker in a.list_universe(): li = a.get_income_statements(ticker)[0].line_items assert abs((li["short_term_debt"] + li["long_term_debt"]) - li["total_debt"]) < 1.0, ticker ALL_TESTS = [obj for name, obj in list(globals().items()) if name.startswith("test_") and callable(obj)] if __name__ == "__main__": passed, failed = 0, [] for fn in ALL_TESTS: try: fn() passed += 1 print(f"PASS {fn.__name__}") except AssertionError as e: failed.append(fn.__name__) print(f"FAIL {fn.__name__}: {e}") print(f"\n{passed}/{len(ALL_TESTS)} passed") if failed: raise SystemExit(1) ============================================================ 15. INGESTION-SPECIFIC TEST REFERENCES ============================================================ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_seed_reference_data.py:6:Running `PYTHONPATH=. python3 scripts/seed_demo.py` — or calling `ingest_security(db, adapter, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:13:from app.adapters.fmp import FMPAdapter /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:41: adapter = FMPAdapter(api_key="fake", transport=_mock_transport(handler)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:52: adapter = FMPAdapter(api_key="fake", transport=_mock_transport(handler)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:65: adapter = FMPAdapter(api_key="bad-key", transport=_mock_transport(handler)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:79: adapter = FMPAdapter(api_key="fake", transport=_mock_transport(handler)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:130: return FMPAdapter(api_key="k", transport=httpx.MockTransport(handler)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_fmp_adapter.py:135: left as an assertion in a document. `ingest_security()` used to write the BalanceSheet, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/tests/test_provider_resilience.py:5:`app/workers/ingest.py::ingest_security_task` this pass (not directly testable here — Celery ============================================================ 16. SOURCE/SCHEMA FIELDS USED BY _record_conflicts ============================================================ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:13: selected_value: Optional[float] /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:15: selection_reason: str /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:31: selected_tier, selected_value = present[0] /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:35: if selected_value == 0: /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:38: elif abs(val - selected_value) / abs(selected_value) > tolerance_pct: /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:44: return ConflictResolution(selected_value, selected_tier, reason, is_conflicting=conflicting) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:50: value_a: Optional[float] /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:51: value_b: Optional[float] /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:52: selected_value: Optional[float] /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:53: selection_reason: str /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:100: if result.selected_value is not None: /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:101: resolved[field] = result.selected_value /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:117: result.selected_value, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:127: result.selected_value, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py:128: result.selection_reason, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/engines/scoring/data_quality.py:40: DataQualityStatus.CONFLICTING: 15.0, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/engines/scoring/data_quality.py:79: conflicting = sum(1 for s in field_statuses if s == DataQualityStatus.CONFLICTING) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/engines/scoring/confidence.py:38: DataQualityStatus.CONFLICTING: 20.0, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/engines/valuation/multiples.py:58: # AUDIT FIX (StockLab final engineering pass, Part B3 -- docs/AUDIT_FAIR_VALUE_B3.md). /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/engines/valuation/margin_of_safety.py:50: # AUDIT FIX (StockLab final engineering pass, Part B3 -- docs/AUDIT_FAIR_VALUE_B3.md). /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/engines/valuation/blend.py:59: AUDIT FIX (StockLab final engineering pass, Part B3 -- docs/AUDIT_FAIR_VALUE_B3.md). This was /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/ingest.py:273: entity_table=entity_table, entity_id=entity_id, field=c.field, status="CONFLICTING", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/ingest.py:274: source_a_id=source_a.id, value_a=c.value_a, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/ingest.py:275: source_b_id=source_b.id if source_b is not None else None, value_b=c.value_b, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/ingest.py:276: selected_value=c.selected_value, selected_source_id=source_a.id if c.selected_value == c.value_a else (source_b.id if source_b is not None else None), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/ingest.py:277: selection_reason=c.selection_reason, /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/recompute.py:543: # AUDIT FIX (StockLab final engineering pass, Part B3 -- docs/AUDIT_FAIR_VALUE_B3.md). /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:38: status: Mapped[str] = mapped_column(String(16)) # VERIFIED|CALCULATED|ESTIMATED|ASSUMPTION|MISSING|STALE|CONFLICTING /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:39: source_a_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:40: value_a: Mapped[Optional[float]] = mapped_column(nullable=True) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:41: source_b_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:42: value_b: Mapped[Optional[float]] = mapped_column(nullable=True) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:43: selected_value: Mapped[Optional[float]] = mapped_column(nullable=True) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:44: selected_source_id: Mapped[Optional[str]] = mapped_column(ForeignKey("sources.id"), nullable=True) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py:45: selection_reason: Mapped[Optional[str]] = mapped_column(nullable=True) ============================================================ 17. MIGRATION — DATAQUALITY / FINANCIAL TABLE DEFINITIONS ============================================================ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0004_alert_events_emerging.py:38: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0004_alert_events_emerging.py:41: sa.Column("security_id", sa.String(36), sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0004_alert_events_emerging.py:61: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0004_alert_events_emerging.py:64: sa.Column("security_id", sa.String(36), sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0004_alert_events_emerging.py:80: # The Emerging Opportunities ranking query: highest emerging score among EMERGING companies. /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0003_revoked_tokens.py:32: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:31: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:45: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:55: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:62: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:71: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:82: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:83: "sources", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:94: op.create_index("ix_sources_provider", "sources", ["provider"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:95: op.create_index("ix_sources_is_demo", "sources", ["is_demo"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:96: op.create_index("ix_sources_content_hash", "sources", ["content_hash"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:98: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:99: "companies", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:111: op.create_index("ix_companies_legal_name", "companies", ["legal_name"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:112: op.create_index("ix_companies_display_name", "companies", ["display_name"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:113: op.create_index("ix_companies_country_id", "companies", ["country_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:114: op.create_index("ix_companies_sector_id", "companies", ["sector_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:115: op.create_index("ix_companies_industry_id", "companies", ["industry_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:117: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:118: "securities", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:120: sa.Column("company_id", UUID, sa.ForeignKey("companies.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:135: op.create_index("ix_securities_company_id", "securities", ["company_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:136: op.create_index("ix_securities_exchange_id", "securities", ["exchange_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:137: op.create_index("ix_securities_ticker", "securities", ["ticker"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:138: op.create_index("ix_securities_isin", "securities", ["isin"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:140: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:141: "financial_periods", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:143: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:150: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:155: op.create_index("ix_financial_periods_security_id", "financial_periods", ["security_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:156: op.create_index("ix_financial_periods_period_end", "financial_periods", ["period_end"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:157: op.create_index("ix_financial_periods_filing_date", "financial_periods", ["filing_date"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:160: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:163: sa.Column("financial_period_id", UUID, sa.ForeignKey("financial_periods.id"), nullable=False, unique=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:169: _statement_table("income_statements", [ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:180: _statement_table("balance_sheets", [ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:191: _statement_table("cash_flows", [ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:196: _statement_table("shares", [ /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:197: sa.Column("shares_outstanding", sa.Float, nullable=True), sa.Column("diluted_shares", sa.Float, nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:200: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:203: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:209: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:215: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:218: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:221: sa.Column("shares_repurchased", sa.Float, nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:222: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:228: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:231: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:237: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:245: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:248: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:255: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:262: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:265: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:271: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:278: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:281: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:286: sa.Column("source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:293: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:296: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:315: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:318: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:334: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:337: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:358: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:361: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:380: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:381: "data_quality", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:387: sa.Column("source_a_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:389: sa.Column("source_b_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:392: sa.Column("selected_source_id", UUID, sa.ForeignKey("sources.id"), nullable=True), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:396: op.create_index("ix_data_quality_entity_id", "data_quality", ["entity_id"]) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:398: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:411: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:422: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:431: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:435: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:443: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:447: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:457: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:467: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:471: sa.Column("security_id", UUID, sa.ForeignKey("securities.id"), nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:472: sa.Column("shares", sa.Float, nullable=False), /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:481: op.create_table( /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:497: "screen_results", "screeners", "data_quality", "scores", "valuation", "metric_history", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:499: "shares", "cash_flows", "balance_sheets", "income_statements", "financial_periods", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0001_initial_schema.py:500: "securities", "companies", "sources", "users", "industries", "sectors", "exchanges", "countries", /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0002_confidence_data_quality_scores.py:1:"""add confidence_score/data_quality_score columns to scores /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0002_confidence_data_quality_scores.py:8:`compute_data_quality_score()` (docs/FINAL_REPORT.md §7's "Score/Confidence/Data Quality wiring /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0002_confidence_data_quality_scores.py:32: op.add_column("scores", sa.Column("data_quality_score", sa.Float, nullable=True)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0002_confidence_data_quality_scores.py:33: op.add_column("scores", sa.Column("data_quality_components", sa.JSON, nullable=True)) /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0002_confidence_data_quality_scores.py:37: op.drop_column("scores", "data_quality_components") /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/alembic/versions/0002_confidence_data_quality_scores.py:38: op.drop_column("scores", "data_quality_score") ============================================================ 18. EXACT CALL GRAPH AROUND resolve_line_items ============================================================ 1-""" 2-Ingestion pipeline (spec §45/§52, docs/ARCHITECTURE.md §4). One security's worth of raw 3-provider data -> normalized point-in-time rows in Postgres, source-hierarchy conflict 4-resolution applied, then a `financials.updated` signal that triggers recompute (app/workers/ 5-recompute.py). Demo and live data are never blended for the same security (docs/DATA_SOURCES.md 6-§9) — enforced in `_assert_no_demo_live_mix` below, not just documented. 7- 8-AUDIT FIX (StockLab overhaul, final engineering pass, Part A6, docs/AUDIT_PROVIDER_CONFLICT_A6.md): 9-`resolve_conflict` used to be imported here but never called — this file ingested exactly one 10-provider's data per call, so there was nothing for it to resolve BETWEEN (docs/ 11-AUDIT_PROVIDER_RESILIENCE.md's finding). `ingest_security()` now optionally accepts a 12-`secondary_adapter`; when given, both providers' line items are resolved field-by-field via 13:`resolve_line_items()` (app/adapters/resolver.py) and real disagreements are persisted as 14-`DataQuality` rows. When `secondary_adapter` is None — the default, and what every existing caller 15-still passes — behavior is byte-identical to before this pass: nothing about the single-source path 16-changed. 17-""" 18-from __future__ import annotations 19-from datetime import date, timedelta 20- 21- 22-from sqlalchemy.orm import Session 23- 24-from app.adapters.base import ( 25- ProviderAdapter, ProviderAuthError, ProviderNotFoundError, ProviderRateLimitError, 26-) 27-from app.adapters.resolver import merge_statement_line_items, resolve_line_items 28-from app.adapters.validation import ValidationReport, validate_line_items 29-from app.core.config import get_settings 30-from app.core.db import SessionLocal 31-from app.core.logging import get_logger 32-from app.engines.identity import resolve_company_identity 33-from app.engines.reference_data import canonical_mic, exchange_for_mic 34-from app.models import ( 35- BalanceSheet, CashFlow, Company, Country, DataQuality, Estimate, Exchange, FinancialPeriod, 36- Industry, IncomeStatement, Price, Sector, Security, Shares, Source, 37-) 38-from app.workers.celery_app import celery_app 39- 40-logger = get_logger(__name__) 41- 42- 43-class ReferenceDataError(RuntimeError): 44- """A NOT NULL reference row could not be resolved, so ingestion must not proceed. 45- 46- AUDIT FIX (DEMO ingestion defect). `get_or_create_company()` used to build `Company` and 47- `Exchange` rows whose `country_id` was a conditional expression falling back to a null. Both 48- columns are **NOT NULL** 49- (`alembic/versions/0001_initial_schema.py`, `companies.country_id` and `exchanges.country_id`), 50- so that `else None` branch could never succeed: it handed Postgres a NULL for a NOT NULL 51- column and the whole transaction died with 52- 53- psycopg.errors.NotNullViolation: null value in column "country_id" of relation "exchanges" 54- 55- several statements later, with a traceback pointing at `db.flush()` rather than at the missing 56- reference row. The code was written as if the column were nullable; it is not. 57- 58- Failing here instead is the correct behaviour and not merely a nicer message: 59- 60- * an exchange row without a country is **not representable** in this schema, so there is no 61- value to fall back to — inventing one (e.g. the company's country) is exactly the 62- shared-reference-table corruption §18 forbids and that a previous pass removed; 63- * the error names the ticker, the field and the specific remedy, so the fix is a one-line 64- addition to a reference table rather than a database forensics session; 65- * it is raised **before** any row is added to the session, so the caller's transaction is 66- still clean and one bad ticker cannot poison the rest of a seed run. 67- """ 68- 69- 70-# Field tuples pulled out to module level (StockLab overhaul, Part A6) -- previously inline in the 71:# single setattr loop below; now shared between the single-source path and resolve_line_items() 72-# calls so the two can never drift out of sync with each other. 73-_INCOME_FIELDS = ( 74- "revenue", "cogs", "gross_profit", "operating_income", "ebit", "ebitda", "net_income", 75- "eps_diluted", "tax_expense", "pretax_income", "interest_expense", 76- "depreciation_and_amortization", "stock_based_compensation", 77-) 78-_BALANCE_FIELDS = ( 79- "cash_and_equivalents", "short_term_investments", "total_debt", "shareholders_equity", 80- "minority_interest", "preferred_equity", "goodwill", "intangible_assets", 81- "current_assets", "current_liabilities", "receivables", "inventory", 82-) 83-_CASH_FLOW_FIELDS = ("operating_cash_flow", "capital_expenditure", "dividends_paid", "buybacks", "stock_issuance") 84-_SHARES_FIELDS = ("shares_outstanding", "diluted_shares") 85- 86- 87-def _get_or_create_source(db: Session, provider: str, tier: str, is_demo: bool, endpoint: str) -> Source: 88- src = Source(provider=provider, provider_tier=tier, is_demo=is_demo, endpoint=endpoint) 89- db.add(src) 90- db.flush() 91- return src 92- 93- 94-def _assert_no_demo_live_mix(db: Session, security_id: str, is_demo: bool) -> None: 95- """Refuse to write demo rows for a security that already has live rows, or vice versa 96- (docs/DATA_SOURCES.md §9 — hard separation, enforced here, not just at read time).""" 97- existing = ( 98- db.query(FinancialPeriod.id) 99- .join(Source, FinancialPeriod.source_id == Source.id) 100- .filter(FinancialPeriod.security_id == security_id, Source.is_demo == (not is_demo)) 101- .first() 102- ) 103- if existing: 104- raise RuntimeError( 105- f"Refusing to ingest {'demo' if is_demo else 'live'} data for security {security_id}: " 106- f"it already has {'live' if is_demo else 'demo'} financial_periods rows." -- 229- timezone=ref.timezone if ref else "UTC", 230- ) 231- db.add(exchange) 232- db.flush() 233- if identity.is_cross_border_listing: 234- # Not an error — an ADR or a foreign listing. Logged because it is the case where a single 235- # `country` column would have been misleading, and because it is worth being able to count 236- # how much of the universe is cross-border. 237- logger.info("ingest.company.cross_border_listing", ticker=ticker, 238- company_country=identity.company_country.iso2, 239- exchange_country=identity.exchange_country.iso2) 240- 241- security = db.query(Security).filter_by(ticker=ticker, exchange_id=exchange.id).one_or_none() 242- if security is not None: 243- return security 244- 245- company = Company( 246- legal_name=profile.legal_name, display_name=profile.display_name, 247- country_id=country.id, sector_id=sector.id, industry_id=industry.id, 248- website=profile.website, description=profile.description, 249- ) 250- db.add(company) 251- db.flush() 252- 253- security = Security( 254- company_id=company.id, exchange_id=exchange.id, ticker=ticker, 255- isin=profile.isin, currency=profile.currency, discovery_status="QUALIFIED", 256- ) 257- db.add(security) 258- db.flush() 259- return security 260- 261- 262-def _record_conflicts(db: Session, entity_table: str, entity_id: str, conflicts: list, 263- source_a: Source, source_b: Source) -> None: 264: """Persist resolve_line_items()'s flagged disagreements as DataQuality rows (StockLab 265- overhaul, Part A6). entity_id is the owning FinancialPeriod's id for every statement type -- 266- IncomeStatement/BalanceSheet/CashFlow/Shares are each a 1:1 extension of FinancialPeriod with 267- no independently meaningful id of their own, so `entity_table` (e.g. "income_statements") is 268- what disambiguates which statement a given row's conflicting field belongs to, not entity_id 269- itself. A deliberate simplification, not an oversight -- see docs/AUDIT_PROVIDER_CONFLICT_A6.md 270- for the alternative (a dedicated per-statement-row id) and why it wasn't needed for this.""" 271- for c in conflicts: 272- db.add(DataQuality( 273- entity_table=entity_table, entity_id=entity_id, field=c.field, status="CONFLICTING", 274- source_a_id=source_a.id, value_a=c.value_a, 275- source_b_id=source_b.id if source_b is not None else None, value_b=c.value_b, 276- selected_value=c.selected_value, selected_source_id=source_a.id if c.selected_value == c.value_a else (source_b.id if source_b is not None else None), 277- selection_reason=c.selection_reason, 278- )) 279- 280- 281-def _fetch_statements_by_period(adapter, method, ticker: str, label: str) -> dict: 282- """Fetch one statement type and index it by (period_end, period_type). 283- 284- Returns an empty dict -- never raises -- when the adapter has not implemented this statement 285- (EODHDAdapter's balance-sheet/cash-flow methods are `NotImplementedError` today). An adapter 286- with a partial implementation must degrade to "that statement's columns stay empty for this 287- run", exactly as before this fix, rather than crashing every ingestion task. 288- """ 289- try: 290- rows = method(ticker) 291- except NotImplementedError: 292- logger.warning( 293- "ingest.security.statement_not_implemented", 294- ticker=ticker, provider=adapter.name, statement=label, 295- ) 296- return {} 297- return {(r.period_end, r.period_type): r.line_items for r in rows} 298- 299- -- 331- volume=bar.volume, 332- currency=bar.currency, 333- source_id=source.id, 334- ) 335- db.add(row) 336- written += 1 337- else: 338- row.open = bar.open 339- row.high = bar.high 340- row.low = bar.low 341- row.close = bar.close 342- row.adjusted_close = bar.adjusted_close 343- row.volume = bar.volume 344- row.currency = bar.currency 345- row.source_id = source.id 346- 347- logger.info( 348- "ingest.prices.complete", 349- ticker=ticker, 350- provider=adapter.name, 351- bars=len(bars), 352- inserted=written, 353- ) 354- return written 355- 356- 357-def ingest_security(db: Session, adapter: ProviderAdapter, ticker: str, 358- secondary_adapter: ProviderAdapter | None = None, 359- include_quarterly: bool = False) -> str: 360- """Ingest one security's financials from `adapter` (the primary/only source for this call). 361- 362- AUDIT FIX (StockLab overhaul, final engineering pass, Part A6): `secondary_adapter` is new and 363- optional. When None -- the default, and what every caller in this codebase passed before this 364- pass -- this function's behavior is completely unchanged: every field comes straight from 365- `adapter`, exactly as before. When given a second adapter, each period's line items are 366: resolved field-by-field against the primary's via `resolve_line_items()` (source-hierarchy 367- tie-break, real disagreements flagged) instead of blindly overwritten by whichever ran last. 368- See docs/AUDIT_PROVIDER_CONFLICT_A6.md for the full design, the period-alignment limitation, 369- and honest TESTED/NOT TESTED status -- none of the dual-source path has been executed against 370- real provider data or a real database in this build environment. 371- """ 372- is_demo = adapter.tier == "DEMO" 373- security = get_or_create_company(db, adapter, ticker) 374- _assert_no_demo_live_mix(db, security.id, is_demo) 375- 376- source = _get_or_create_source(db, adapter.name, adapter.tier, is_demo, f"ingest:{ticker}") 377- 378- # AUDIT FIX (StockLab final engineering pass, Part A9 -- CRITICAL, see 379- # docs/AUDIT_INGEST_STATEMENTS_A9.md). Before this fix, this function fetched ONLY the income 380- # statement and then wrote the BalanceSheet, CashFlow and Shares rows out of that same 381- # response's line items. adapter.get_balance_sheets() and adapter.get_cash_flows() were never 382- # called anywhere in the application. That is invisible in DEMO mode (DemoDataAdapter returns 383- # one fully-merged dict from all three methods) and catastrophic against FMP, whose income 384- # response carries only _INCOME_MAP's fields -- every balance-sheet, cash-flow and share-count 385- # column would have been written NULL for real provider data, taking most metrics, all five 386- # pillar scores and every valuation down with it. 387- income_periods = adapter.get_income_statements(ticker) 388- balance_by_period = _fetch_statements_by_period(adapter, adapter.get_balance_sheets, ticker, "balance_sheets") 389- cash_by_period = _fetch_statements_by_period(adapter, adapter.get_cash_flows, ticker, "cash_flows") 390- 391- # Part A9: optional quarterly ingestion. Off by default (Settings.INGEST_QUARTERLY_PERIODS). 392- # Quarterly rows are stored as ordinary FinancialPeriod rows with period_type Q1..Q4 -- the 393- # (security_id, period_end, period_type, filing_date) unique constraint keeps them distinct 394- # from the annual rows, and every existing query that filters period_type == "FY" is 395- # unaffected. app/engines/ttm.py turns four of them into a TTM basis at recompute time. 396- if include_quarterly: 397- try: 398- quarterly_income = adapter.get_income_statements(ticker, period="quarter") 399- except NotImplementedError: 400- logger.warning("ingest.security.quarterly_not_implemented", ticker=ticker, 401- provider=adapter.name) -- 449- # instant this function tried to call the secondary adapter's statement method. Treated 450- # as "no secondary statement data available this run", not a crash: falls back to 451- # exactly the primary-only behavior secondary_adapter=None already has for every period. 452- logger.warning("ingest.security.secondary_statements_not_implemented", ticker=ticker, 453- secondary_provider=secondary_adapter.name) 454- 455- for p in income_periods: 456- fp = db.query(FinancialPeriod).filter_by( 457- security_id=security.id, period_end=p.period_end, period_type=p.period_type, filing_date=p.filing_date, 458- ).one_or_none() 459- if fp is None: 460- fp = FinancialPeriod( 461- security_id=security.id, period_end=p.period_end, period_type=p.period_type, 462- filing_date=p.filing_date, currency=p.currency, source_id=source.id, 463- ) 464- db.add(fp) 465- db.flush() 466- 467- key = (p.period_end, p.period_type) 468- li = merge_statement_line_items( 469- p.line_items, 470- balance_by_period.get(key), 471- cash_by_period.get(key), 472- ) 473- secondary_p = secondary_by_period.get(key) 474- li_b = ( 475- merge_statement_line_items( 476- secondary_p.line_items, secondary_balance.get(key), secondary_cash.get(key) 477- ) 478- if secondary_p is not None 479- else None 480- ) 481- secondary_tier = secondary_adapter.tier if (secondary_adapter is not None and li_b is not None) else None 482- 483- income = fp.income_statement or IncomeStatement(financial_period_id=fp.id) 484: resolved, conflicts = resolve_line_items( 485- li, li_b, _INCOME_FIELDS, adapter.tier, secondary_tier 486- ) 487- merged_report = ValidationReport() 488- resolved = validate_line_items( 489- resolved, 490- merged_report, 491- row_key=f"{ticker}:{p.period_end}:{p.period_type}:income", 492- ) 493- for field, value in resolved.items(): 494- setattr(income, field, value) 495- db.add(income) 496- if conflicts: 497- _record_conflicts(db, "income_statements", fp.id, conflicts, source, secondary_source) 498- 499- balance = fp.balance_sheet or BalanceSheet(financial_period_id=fp.id) 500: resolved, conflicts = resolve_line_items( 501- li, li_b, _BALANCE_FIELDS, adapter.tier, secondary_tier 502- ) 503- merged_report = ValidationReport() 504- resolved = validate_line_items( 505- resolved, 506- merged_report, 507- row_key=f"{ticker}:{p.period_end}:{p.period_type}:balance", 508- ) 509- for field, value in resolved.items(): 510- setattr(balance, field, value) 511- db.add(balance) 512- if conflicts: 513- _record_conflicts(db, "balance_sheets", fp.id, conflicts, source, secondary_source) 514- 515- cash_flow = fp.cash_flow or CashFlow(financial_period_id=fp.id) 516: resolved, conflicts = resolve_line_items( 517- li, li_b, _CASH_FLOW_FIELDS, adapter.tier, secondary_tier 518- ) 519- merged_report = ValidationReport() 520- resolved = validate_line_items( 521- resolved, 522- merged_report, 523- row_key=f"{ticker}:{p.period_end}:{p.period_type}:cash_flow", 524- ) 525- for field, value in resolved.items(): 526- setattr(cash_flow, field, value) 527- db.add(cash_flow) 528- if conflicts: 529- _record_conflicts(db, "cash_flows", fp.id, conflicts, source, secondary_source) 530- 531- shares = fp.shares or Shares(financial_period_id=fp.id) 532: resolved, conflicts = resolve_line_items( 533- li, li_b, _SHARES_FIELDS, adapter.tier, secondary_tier 534- ) 535- merged_report = ValidationReport() 536- resolved = validate_line_items( 537- resolved, 538- merged_report, 539- row_key=f"{ticker}:{p.period_end}:{p.period_type}:shares", 540- ) 541- for field, value in resolved.items(): 542- setattr(shares, field, value) 543- db.add(shares) 544- if conflicts: 545- _record_conflicts(db, "shares", fp.id, conflicts, source, secondary_source) 546- 547- # AUDIT FIX (final master pass, §21) — the fourth "engine built, never wired" defect. 548- # `ProviderAdapter.get_estimates()` is implemented by FMP and by the demo adapter, 549- # `ProviderEstimateRow` exists, and the `estimates` table exists with a unique constraint on 550- # (security_id, period_end, metric, as_of_date). Nothing ever called it, nothing ever wrote a 551- # row, and `build_snapshot_from_db()` never passed `forward_eps_estimate` — so it was None on 552- # every real run and `forward_pe` and `eps_growth_forward` were PERMANENTLY NULL. `forward_pe` 553- # is a member of the Valuation pillar, so every Valuation score was computed from six of its 554- # seven metrics with nothing saying so. 555- # Writes are idempotent against that unique constraint: an estimate already stored for the 556- # same (period_end, metric, as_of_date) is updated in place, never duplicated on re-ingestion. 557- estimates_written = _ingest_estimates(db, adapter, security.id, ticker, source) 558- prices_written = _ingest_prices(db, adapter, security.id, ticker, source) 559- 560- db.commit() 561- logger.info("ingest.security.complete", ticker=ticker, provider=adapter.name, is_demo=is_demo, 562- periods=len(income_periods), estimates=estimates_written, 563- prices=prices_written, 564- secondary_provider=(secondary_adapter.name if secondary_adapter else None)) 565- return security.id 566- 567- ============================================================ 19. EXACT _record_conflicts IMPLEMENTATION ============================================================ 247- country_id=country.id, sector_id=sector.id, industry_id=industry.id, 248- website=profile.website, description=profile.description, 249- ) 250- db.add(company) 251- db.flush() 252- 253- security = Security( 254- company_id=company.id, exchange_id=exchange.id, ticker=ticker, 255- isin=profile.isin, currency=profile.currency, discovery_status="QUALIFIED", 256- ) 257- db.add(security) 258- db.flush() 259- return security 260- 261- 262:def _record_conflicts(db: Session, entity_table: str, entity_id: str, conflicts: list, 263- source_a: Source, source_b: Source) -> None: 264- """Persist resolve_line_items()'s flagged disagreements as DataQuality rows (StockLab 265- overhaul, Part A6). entity_id is the owning FinancialPeriod's id for every statement type -- 266- IncomeStatement/BalanceSheet/CashFlow/Shares are each a 1:1 extension of FinancialPeriod with 267- no independently meaningful id of their own, so `entity_table` (e.g. "income_statements") is 268- what disambiguates which statement a given row's conflicting field belongs to, not entity_id 269- itself. A deliberate simplification, not an oversight -- see docs/AUDIT_PROVIDER_CONFLICT_A6.md 270- for the alternative (a dedicated per-statement-row id) and why it wasn't needed for this.""" 271- for c in conflicts: 272- db.add(DataQuality( 273- entity_table=entity_table, entity_id=entity_id, field=c.field, status="CONFLICTING", 274- source_a_id=source_a.id, value_a=c.value_a, 275- source_b_id=source_b.id if source_b is not None else None, value_b=c.value_b, 276- selected_value=c.selected_value, selected_source_id=source_a.id if c.selected_value == c.value_a else (source_b.id if source_b is not None else None), 277- selection_reason=c.selection_reason, 278- )) 279- 280- 281-def _fetch_statements_by_period(adapter, method, ticker: str, label: str) -> dict: 282- """Fetch one statement type and index it by (period_end, period_type). 283- 284- Returns an empty dict -- never raises -- when the adapter has not implemented this statement 285- (EODHDAdapter's balance-sheet/cash-flow methods are `NotImplementedError` today). An adapter 286- with a partial implementation must degrade to "that statement's columns stay empty for this 287- run", exactly as before this fix, rather than crashing every ingestion task. 288- """ 289- try: 290- rows = method(ticker) 291- except NotImplementedError: 292- logger.warning( 293- "ingest.security.statement_not_implemented", 294- ticker=ticker, provider=adapter.name, statement=label, 295- ) 296- return {} 297- return {(r.period_end, r.period_type): r.line_items for r in rows} 298- 299- 300- 301-def _ingest_prices( 302- db: Session, 303- adapter: ProviderAdapter, 304- security_id: str, 305- ticker: str, 306- source: Source, 307-) -> int: 308- """Ingest daily OHLCV bars for one security, idempotently by security/date.""" 309- end = date.today() 310- start = end - timedelta(days=365) 311- 312- bars = adapter.get_prices(ticker, start, end) 313- written = 0 314- 315- for bar in bars: 316- row = ( 317- db.query(Price) 318- .filter_by(security_id=security_id, date=bar.date) 319- .one_or_none() 320- ) 321- 322- if row is None: 323- row = Price( 324- security_id=security_id, 325- date=bar.date, 326- open=bar.open, 327- high=bar.high, 328- low=bar.low, 329- close=bar.close, 330- adjusted_close=bar.adjusted_close, 331- volume=bar.volume, 332- currency=bar.currency, 333- source_id=source.id, 334- ) 335- db.add(row) 336- written += 1 337- else: 338- row.open = bar.open 339- row.high = bar.high 340- row.low = bar.low 341- row.close = bar.close 342- row.adjusted_close = bar.adjusted_close ============================================================ 20. CURRENT SOURCE HASHES ============================================================ 782820e1a686029f751c148d327dbef76784fb6ef1d52107aefa9c3a75c08af6 /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/workers/ingest.py 183842352c4c3b0ec9e2fabffa043ff6fbb125c2a0b602c9da53477d17b7bfb8 /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/resolver.py 4ff976fb3b68594b7991c8d172143278670ae3685d49ad61f5b945db61b4f19a /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/validation.py f4fd2dc266df3047c21414ee64cc11587c80384a80f9dbab84d6d353df062501 /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/base.py 233f1caff590a88a0ebfcc897047c92426bff31a40ce3b69d890410a27ac7560 /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/adapters/schemas.py bc24958a9e68a43e4fd6af4ca7695149523aef931cd4ecee1264f47a3d00a977 /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/financials.py 14d5836dd9c63e44e3b49ed1249bdc643cdb9e154678efb5772f7ce4677d5e8e /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/companies.py 1bc468b9c94aa8a313411b2d44cf47cce58d884956ac54c32ffdad751567d78b /data/files/bmw/stocklab-release-candidate-sl003-final-20260910-111453/app/models/governance.py ============================================================ 21. PRODUCTION SAFETY CHECK ============================================================ stocklab-redis docker.io/library/redis:7-alpine Up 34 hours (healthy) stocklab-db docker.io/library/postgres:16-alpine Up 34 hours (healthy) stocklab-release-qa-redis docker.io/library/redis:7 Up 27 hours stocklab-frontend localhost/stocklab-frontend:release-qa-20260909-182234 Up 25 hours (healthy) stocklab-backend localhost/stocklab-backend:screener-percent-final-20260909-194639 Up 22 hours (healthy) stocklab-worker localhost/stocklab-backend:screener-percent-final-20260909-194639 Up 22 hours (healthy) stocklab-beat localhost/stocklab-backend:screener-percent-final-20260909-194639 Up 22 hours Restart counters: stocklab-backend restart=0 stocklab-worker restart=0 stocklab-beat restart=0 stocklab-db restart=0 stocklab-redis restart=0 stocklab-frontend restart=0 Ports: LISTEN 0 2048 127.0.0.1:18000 0.0.0.0:* users:(("uvicorn",pid=1574355,fd=11)) LISTEN 0 511 127.0.0.1:6379 0.0.0.0:* users:(("redis-server",pid=1133711,fd=6)) LISTEN 0 200 127.0.0.1:5432 0.0.0.0:* users:(("postgres",pid=1141225,fd=6)) LISTEN 0 511 127.0.0.1:13000 0.0.0.0:* users:(("next-server (v1",pid=1497838,fd=18)) ============================================================ 22. RESULT ============================================================ CONTRACT_INSPECTION=COMPLETE READ_ONLY=TRUE NO_DATABASE_CONTAINER_STARTED=TRUE NO_DATABASE_WRITE=TRUE NO_PRODUCTION_RESTART=TRUE NO_PRODUCTION_MUTATION=TRUE LOG=/data/files/bmw/stocklab-sl004-contract-inspection-20260910-191856.log ============================================================