"""
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
import httpx
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"))

    if not income_periods:
        raise RuntimeError(
            f"Provider returned no financial periods for {ticker}; "
            "refusing to commit stale data or trigger recompute."
        )

    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, httpx.TransportError),
    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)
