Files
signal-platform/app/services/sec_fundamentals_importer.py
T
dennisthiessenandClaude Opus 5 77570557db feat(sec): cap how long the fundamentals import can stay deferred
MISSING_XBRL_RETRY_DAYS bounds how long ONE filing blocks promotion. It does
not bound the import as a whole, and the two come apart because a blocking
filing is only queued by promote(), which a deferred run never reaches. During
a rolling supply of unresolvable filings — earnings season, when SEC's
Company-Facts aggregation lags furthest — each new arrival restarts the 3-day
clock before the previous one clears, and nothing is written at all: not the
good rows, not the gap rows that would stop those filings blocking again.

Add an aggregate ceiling. Once promotions have been stale for
PROMOTION_CEILING_DAYS (7), every unresolved filing is aged past the retry
window in place, so promote() queues them all through the path that already
exists, source_max_date advances, and _missing() keeps queued rows aged-out on
later runs. The import self-heals instead of compounding.

Deliberately not the alternative of queueing gap rows on a deferred run: that
would drop the grace period to a single run for every filing, including the
common case of a Company-Facts lag that resolves in a day, and it needs a write
on a run that failed validation.

The per-filing window is untouched, a never-promoted source never trips (that
is initial setup, not a wedge), and affected symbols stay barred from setups
either way since setup_blocked_ciks ignores the window. A forced promotion
raises promotion_ceiling_forced so the safety valve is never silent.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-11 16:58:46 +02:00

1087 lines
50 KiB
Python

"""SEC fundamentals importer (workstream A, phase A3).
A ``SourceImporter`` (see ``app/services/data_import.py``) that populates the
immutable ``fundamental_snapshots`` from SEC Company Facts and back-fills
``tickers.cik/sic/sic_description``. EDGAR-daily-index driven: it fetches
companyfacts only for tracked issuers that filed since the last run (full-history
backfill on the first run / for newly-added issuers). Shadow only — nothing reads
snapshots until A4.
Guardrails (design + reviews):
- ``detect_revision`` caches the resolved universe + the exact tracked index rows
and composes the revision from them; ``stage`` consumes those same cached inputs
(it does not refetch the index/universe) so promoted data always matches the
computed revision.
- Resolution is read-only in ``stage`` (proposals only); ticker writes happen in
``promote`` via ``apply_ticker_updates``.
- ``validate`` runs the **index↔Company-Facts consistency gate** before any write:
a tracked XBRL index accession missing from Company Facts fails the run (the two
are separate SEC products that can lag) so we retry rather than record a
null/partial snapshot. Non-XBRL amendments are skipped with a recorded reason.
A failure here blocks every later run (``source_max_date`` only advances on a
promoted run), so it names the offending filings in the alert and separates the
causes — ``not_in_companyfacts`` (facts lag) vs ``not_in_submissions`` (the
index row is absent from the issuer's own filing list, which no retry fixes).
- **Co-registrant recovery**, because "missing from Company Facts" is often not
missing at all: SEC files some combined parent/subsidiary filings' XBRL under
the co-registrant's CIK, so the ticker-carrying parent's own file never gets
that accession. The daily index lists every co-registrant of an accession, so
the facts are found there and re-stamped to the real filer — guarded by a
share-count continuity check so a subsidiary's standalone numbers can never be
stored as the parent's. Confirmed 2026-07-27 (NEE via FPL, DOW via Dow Chemical)
and it is not transient: an NEE filing misattributed in 2014 is still misfiled.
- **Bounded blocking.** Anything still unresolvable after ``MISSING_XBRL_RETRY_DAYS``
stops failing the whole import and enters a durable retry queue. The scheduled
importer retries queued accessions automatically, while the affected issuer is
excluded from actionable setups until its filing is recovered.
- ``promote`` inserts snapshots ``ON CONFLICT (accession) DO NOTHING`` (immutable),
reports differing existing accessions, and applies ticker updates in the same
transaction.
- ``reparse=True`` is the one exception to immutability, and it is deliberate:
it restages every accession with the current parser and **rewrites** the rows
that now reconstruct differently. Immutability protects SEC's record (one row
per accession, amendments retained) — but the stored row is *our* reconstruction,
so after a parser fix, keeping it is preserving a stale cache, not history.
Manually invoked through ``scripts/reparse_fundamentals.py``; never scheduled.
"""
from __future__ import annotations
import json
import logging
from collections import Counter, defaultdict
from dataclasses import dataclass, field, replace
from datetime import date, datetime, time, timedelta, timezone
from typing import Any, Callable
from sqlalchemy import delete, exists, select, update
from app.database import insert_for_session
from app.models.data_import_run import DataImportRun
from app.models.fundamental_snapshot import FundamentalSnapshot
from app.models.sec_filing_gap import SecFilingGap
from app.models.system_event import SystemEvent
from app.services import fundamentals_quality_service
from app.services import sec_facts_parser as parser
from app.services import sec_universe
from app.services.data_import import STATUS_PROMOTED, ValidationResult
from app.services.sec_client import SecClient, SecError, cik10
from app.services.sec_facts_parser import FilingMeta, SnapshotRow
from app.services.sec_universe import ResolvedUniverse
logger = logging.getLogger(__name__)
SOURCE = "sec_facts"
_XBRL_FORMS = {"10-K", "10-Q", "10-K/A", "10-Q/A"}
# On the one-time backfill, require this fraction of tracked issuers to yield at
# least one snapshot (guards a broken fetch/parse from promoting a hollow table).
MIN_BACKFILL_COVERAGE = 0.5
# How long an index accession may stay unresolvable before the run stops failing
# on it. Genuine index↔facts lag clears within a day (a weekend stretches it to
# three); past that it is misfiled, not late, and blocking forever costs more
# than the missing filing does — see the unresolved-filing guardrail below.
MISSING_XBRL_RETRY_DAYS = 3
# Aggregate ceiling on deferral. MISSING_XBRL_RETRY_DAYS bounds how long ONE
# filing blocks; it does not bound how long the import as a whole can stay
# deferred. Those differ because a blocking filing is only queued by promote(),
# which a deferred run never reaches — so during a rolling supply of
# unresolvable filings (earnings season, when SEC's Company-Facts aggregation is
# furthest behind) each new arrival restarts the clock before the previous one
# clears, and nothing is written at all: not the good rows, not the gap rows
# that would stop those filings blocking again.
#
# Once promotions have been stale this long, every unresolved filing is treated
# as past the window. promote() then queues them all (see the _past_retry_window
# call there), source_max_date advances, and _missing() forces queued rows
# aged-out on later runs so they never block again — the import self-heals
# through the paths that already exist.
#
# Well above MISSING_XBRL_RETRY_DAYS so ordinary overlapping blocks never trip
# it. Affected symbols stay barred from setups either way: setup_blocked_ciks is
# built from every missing filing regardless of window.
PROMOTION_CEILING_DAYS = 7
FILING_GAP_ESCALATE_DAYS = 14
# Share-count band a co-registrant-recovered row must land in, relative to the
# issuer's own last snapshot. Wide enough for buybacks/issuance, nowhere near
# wide enough to let a subsidiary shell's token float through (see _shares_continuous).
RECOVERY_SHARES_MIN = 0.5
RECOVERY_SHARES_MAX = 2.0
_SNAPSHOT_COLS = (
"cik", "accession", "form", "filed_date", "accepted_at", "period_start",
"period_end", "fiscal_year", "fiscal_period", "revenue", "net_income",
"operating_income", "diluted_eps", "cfo", "capex", "depreciation_amortization",
"cash_and_st_investments", "total_debt", "shares_outstanding",
"shares_outstanding_date", "weighted_avg_diluted_shares",
)
# Compare ALL source fields (every column except the accession key) to flag a
# differing existing accession — immutable, so we report, never mutate.
_COMPARE_COLS = tuple(c for c in _SNAPSHOT_COLS if c != "accession")
@dataclass
class StagedFundamentals:
resolved: ResolvedUniverse
sic_updates: list[tuple[int, str | None, str | None]] = field(default_factory=list)
rows: list[SnapshotRow] = field(default_factory=list)
skipped_filings: list[dict[str, str]] = field(default_factory=list)
field_issues: list[dict[str, str]] = field(default_factory=list)
skipped_non_xbrl: list[dict[str, str]] = field(default_factory=list)
# Index rows we could not resolve to Company Facts, with a per-row reason
# (not_in_companyfacts | not_in_submissions | ...) — see _missing().
missing_xbrl: list[dict[str, Any]] = field(default_factory=list)
# Accessions parsed out of a co-registrant's Company Facts file.
recovered: list[dict[str, Any]] = field(default_factory=list)
invalid_payloads: list[dict[str, str]] = field(default_factory=list)
existing_accessions: set[str] = field(default_factory=set)
# Tracked issuers whose registrant has NO XBRL 10-K/10-Q at all: they can
# never yield a snapshot, so this is a resolution problem (a ticker pointed
# at a successor shell), not missing data. See sec_universe.CIK_OVERRIDES_KEY.
no_xbrl_filings: list[dict[str, Any]] = field(default_factory=list)
discrepancies: list[dict[str, Any]] = field(default_factory=list)
backfill: bool = False
issuers_fetched: int = 0
issuers_with_rows: int = 0
def _now() -> datetime:
return datetime.now(timezone.utc)
class SecFundamentalsImporter:
source = SOURCE
def __init__(
self,
*,
client_factory: Callable[[], SecClient] | None = None,
today: date | None = None,
reparse: bool = False,
) -> None:
self._client_factory = client_factory or (lambda: SecClient())
self.today = today or _now().date()
# Reparse: re-derive every stored accession with the CURRENT parser and
# rewrite the ones that now reconstruct differently. Snapshots are
# immutable with respect to SEC (one row per accession, amendments kept),
# but the stored row is *our reconstruction* — when a parser bug is fixed,
# leaving it stale is not immutability, it is a stale cache. Manually
# invoked via scripts/reparse_fundamentals.py; never scheduled.
self.reparse = reparse
# cached by detect_revision, consumed by stage:
self._resolved: ResolvedUniverse | None = None
self._index_rows: list[dict[str, Any]] = []
# accession -> the OTHER CIKs the daily index lists it under (co-registrants
# of a combined filing). Only populated for accessions a tracked issuer filed.
self._coregistrants: dict[str, list[int]] = {}
self._retry_rows: list[dict[str, Any]] = []
self._latest_index_date: date | None = None
self._backfill = False
# Set by validate() when the aggregate ceiling forced the block open;
# read by promote() to alert that it did.
self._ceiling_tripped: dict[str, Any] | None = None
# -- SourceImporter protocol -------------------------------------------
async def detect_revision(self, db) -> str | None:
async with self._client_factory() as client:
self._resolved = await sec_universe.resolve_ciks(db, client)
last_processed = await self._last_processed_index_date(db)
self._latest_index_date = await client.latest_index_date(self.today)
if self._latest_index_date is None:
raise SecError("no EDGAR daily index available")
# Reparse needs every accession restaged, not just those filed since
# the last run — the facts a fixed parser now accepts were never
# stored, so a reparse cannot be served from the database.
self._coregistrants = {}
if last_processed is None or self.reparse:
self._backfill = True
self._index_rows = []
else:
self._backfill = False
self._index_rows = await self._collect_index_rows(
client, last_processed, self._latest_index_date
)
content = sec_universe.index_content_hash(self._index_rows)
revision = sec_universe.compose_revision(
self._latest_index_date, content, self._resolved.symbol_to_cik
)
self._retry_rows = []
if not self._backfill:
self._retry_rows = await self._retry_backlog(
db,
set(self._resolved.cik_to_ticker_ids),
)
# Company Facts can change while the daily index revision stays fixed.
# Returning None deliberately bypasses the framework's no-op gate so a
# scheduled run retries every active gap.
return None if self._retry_rows else revision
async def stage(self, db) -> StagedFundamentals:
assert self._resolved is not None, "detect_revision must run first"
resolved = self._resolved
staged = StagedFundamentals(resolved=resolved, backfill=self._backfill)
cik_to_tids = resolved.cik_to_ticker_ids
# Whole index rows (not bare accessions): form + index date are what make
# an unresolvable filing diagnosable without re-walking the index by hand.
filed_by_cik: dict[int, list[dict[str, Any]]] = defaultdict(list)
for r in self._index_rows:
if r["cik"] in cik_to_tids:
filed_by_cik[r["cik"]].append(r)
# Promoted-around filings live in a small durable retry queue, including
# the one-time migration backfill. Merge them into normal incremental
# work so the scheduled importer heals them without operator action.
if not self._backfill:
seen = {
(int(cik), row["accession"])
for cik, rows in filed_by_cik.items()
for row in rows
}
for row in self._retry_rows:
cik = int(row["cik"])
key = (cik, row["accession"])
if key in seen:
continue
filed_by_cik[cik].append(row)
seen.add(key)
coregistrants = [int(value) for value in row.get("coregistrants") or []]
if coregistrants:
self._coregistrants[row["accession"]] = coregistrants
existing = await self._ciks_with_snapshots(db, set(cik_to_tids))
if self._backfill:
backfill_ciks = set(cik_to_tids)
else:
# Newly added issuers (resolved but no snapshots yet) get a full-history
# backfill; issuers that already have history are handled incrementally.
backfill_ciks = {c for c in cik_to_tids if c not in existing}
incremental_ciks = set(filed_by_cik) - backfill_ciks
# Continuity reference for co-registrant recovery, read once up front.
last_shares = await self._last_shares_outstanding(db, set(cik_to_tids))
async with self._client_factory() as client:
for cik in sorted(backfill_ciks | incremental_ciks):
is_backfill = cik in backfill_ciks
await self._stage_issuer(
client, cik, is_backfill, filed_by_cik, staged, last_shares
)
# Read-only discrepancy detection: an accession we reconstructed that is
# already stored, differing in ANY source field (immutable → report in
# validation, event on promote, never mutate). Also gives promote the
# existing set so its insert count is dialect-independent.
if staged.rows:
existing = await self._existing_by_accession(db, [r.accession for r in staged.rows])
staged.existing_accessions = set(existing)
for row in staged.rows:
old = existing.get(row.accession)
if old is not None:
fields = _diff_fields(row, old)
if fields:
staged.discrepancies.append({"accession": row.accession, "fields": fields})
return staged
async def _stage_issuer(
self, client, cik, is_backfill, filed_by_cik, staged, last_shares
) -> None:
cf = await client.companyfacts(cik)
bad = _companyfacts_structure_error(cf)
if bad is not None:
# Malformed payload (missing facts/units structure) — record separately
# and fail validation, rather than letting it degrade to skipped rows.
staged.invalid_payloads.append({"cik": cik10(cik), "reason": bad})
staged.issuers_fetched += 1
return
sub = await client.submissions(cik, include_history=is_backfill)
xbrl_meta, nonxbrl = _filing_meta(sub)
if not xbrl_meta:
staged.no_xbrl_filings.append(
{"cik": cik10(cik), "name": sub.get("name"), "tickers": sub.get("tickers")}
)
fiscal_year_end = sub.get("fiscal_year_end")
recovered_rows: list[SnapshotRow] = []
index_rows = {
row["accession"]: row for row in filed_by_cik.get(cik, [])
}
if is_backfill:
accns = set(xbrl_meta)
else:
present = parser.companyfacts_accessions(cf)
accns = set()
for index_row in filed_by_cik.get(cik, []):
accn = index_row["accession"]
if accn in nonxbrl:
staged.skipped_non_xbrl.append({"cik": cik10(cik), "accession": accn})
elif accn not in xbrl_meta:
# The daily index lists it but the issuer's own filing list does
# not (submissions lagging the index, no usable period metadata,
# or a co-registrant filing). NOT a Company-Facts lag — separate
# cause, separate fix, so it gets its own reason.
staged.missing_xbrl.append(
_missing(
cik,
index_row,
"not_in_submissions",
self.today,
self._coregistrants.get(accn),
)
)
elif accn in present:
accns.add(accn)
else:
# Filed, XBRL, but absent from this issuer's Company Facts. Try
# the co-registrant file before treating it as missing data.
row, source_cik = await self._recover_from_coregistrant(
client, cik, accn, xbrl_meta, fiscal_year_end,
last_shares.get(cik10(cik)), staged,
)
if row is not None:
recovered_rows.append(row)
staged.recovered.append({
"cik": cik10(cik),
"accession": accn,
"source_cik": source_cik,
"form": index_row.get("form"),
})
else:
staged.missing_xbrl.append(_missing(
cik, index_row,
# Found, but it did not look like this issuer's own
# numbers — say so; it is not the same as absent.
"coregistrant_facts_rejected" if source_cik
else "not_in_companyfacts",
self.today,
self._coregistrants.get(accn),
))
# fiscalYearEnd (MMDD) is what lets the parser derive period identity from
# reportDate instead of SEC's unreliable fy/fp fields.
result = parser.parse_snapshots(cf, xbrl_meta, accns, fiscal_year_end=fiscal_year_end)
for skipped in result.skipped_filings:
index_row = index_rows.get(skipped["accession"])
if index_row is not None:
# Facts are present but our parser cannot construct a snapshot.
# A new index row keeps the normal grace period before promotion;
# a row already read from the queue retains its _retry_queue marker
# so later imports promote and retry without wedging the index.
staged.missing_xbrl.append(_missing(
cik,
index_row,
"parser_unusable",
self.today,
self._coregistrants.get(skipped["accession"]),
))
staged.rows.extend(result.rows)
staged.rows.extend(recovered_rows)
staged.skipped_filings.extend(result.skipped_filings)
staged.field_issues.extend(result.field_issues)
staged.issuers_fetched += 1
if result.rows or recovered_rows:
staged.issuers_with_rows += 1
# SIC proposal for this issuer's tickers (read-only; applied in promote).
sic = str(sub["sic"]) if sub.get("sic") else None
desc = sub.get("sic_description")
for tid in staged.resolved.cik_to_ticker_ids.get(cik, []):
staged.sic_updates.append((tid, sic, desc))
async def _recover_from_coregistrant(
self, client, cik: int, accn: str, xbrl_meta, fiscal_year_end, reference, staged,
) -> tuple[SnapshotRow | None, str | None]:
"""Look for ``accn``'s facts in a co-registrant's Company Facts file.
SEC sometimes files a combined parent/subsidiary filing's XBRL under the
co-registrant's CIK rather than the filer's — the ticker-carrying parent's
own file simply never gets that accession. Verified 2026-07-27 for NEE
(facts under Florida Power & Light) and DOW (under Dow Chemical); an NEE
filing misattributed the same way in **2014** is still misattributed, so
this does not self-correct and no amount of retrying recovers it.
Returns ``(row, source_cik)`` on success, ``(None, source_cik)`` when the
facts were found but rejected by the continuity guard, ``(None, None)``
when no co-registrant has them.
Incremental path only: the co-registrant map comes from the daily index,
which a backfill/reparse does not walk. A reparse therefore recovers a
filing only once SEC re-files it under the filer's own CIK.
"""
for co in self._coregistrants.get(accn, []):
try:
cf_co = await client.companyfacts(co)
except SecError:
continue # a co-registrant shell often has no facts file at all
if _companyfacts_structure_error(cf_co) is not None:
continue
result = parser.parse_snapshots(
cf_co, xbrl_meta, {accn}, fiscal_year_end=fiscal_year_end
)
if not result.rows:
continue
row = result.rows[0]
if not _shares_continuous(row.shares_outstanding, reference):
return None, cik10(co)
# A recovered row is the one most worth flagging, so its parser caveats
# travel with it rather than being dropped on the way out.
staged.field_issues.extend(result.field_issues)
# parse_snapshots stamps the CIK of the payload it read — re-stamp to
# the issuer that actually filed, or the row lands under the shell.
return replace(row, cik=cik10(cik)), cik10(co)
return None, None
async def validate(self, db, staged: StagedFundamentals) -> ValidationResult:
messages: list[str] = []
# Consistency gate — before any write. Only filings still inside the retry
# window block: a failure here stops every later run too (source_max_date
# advances on promotion alone), so blocking forever on a filing SEC has
# misfiled would cost far more than the one filing it withholds. Older
# ones are carried by promote() as a warning instead. The message names
# the filings: "which ones" has to be in the alert itself, not merely
# reconstructible by re-walking the index.
blocking = _within_retry_window(staged.missing_xbrl)
aged_out = _past_retry_window(staged.missing_xbrl)
# ...unless promotions have been stale past the aggregate ceiling, in
# which case the deferral has cost more than the filings it withholds.
# Ageing them here (not just locally) is deliberate: promote() re-derives
# the queue from the same list, so this is what gets them queued.
self._ceiling_tripped = None
if blocking and db is not None and await self._promotions_stale(db):
for item in staged.missing_xbrl:
item["age_days"] = max(
item.get("age_days", 0), MISSING_XBRL_RETRY_DAYS + 1
)
self._ceiling_tripped = {
"forced": len(blocking),
"unresolved": len(staged.missing_xbrl),
}
blocking = _within_retry_window(staged.missing_xbrl)
aged_out = _past_retry_window(staged.missing_xbrl)
if blocking:
messages.append(
f"{len(blocking)} tracked XBRL filing(s) unresolved within the "
f"{MISSING_XBRL_RETRY_DAYS}-day retry window "
f"({_reason_counts(blocking)}) — retry: {_missing_detail(blocking)}"
)
# Malformed companyfacts payloads must fail, not degrade to skipped rows.
if staged.invalid_payloads:
messages.append(
f"{len(staged.invalid_payloads)} issuer(s) returned a malformed "
"companyfacts payload (missing facts structure)"
)
accns = [r.accession for r in staged.rows]
if len(accns) != len(set(accns)):
messages.append("duplicate accession in staged snapshots")
if staged.backfill:
n_issuers = len(staged.resolved.cik_to_ticker_ids)
coverage = staged.issuers_with_rows / n_issuers if n_issuers else 0.0
if coverage < MIN_BACKFILL_COVERAGE:
messages.append(
f"backfill coverage {coverage:.0%} < {MIN_BACKFILL_COVERAGE:.0%}"
)
summary = {
"backfill": staged.backfill,
"issuers_fetched": staged.issuers_fetched,
"issuers_with_rows": staged.issuers_with_rows,
"snapshot_rows": len(staged.rows),
"skipped_filings": len(staged.skipped_filings),
"field_issues": len(staged.field_issues),
"skipped_non_xbrl": len(staged.skipped_non_xbrl),
"no_xbrl_filings": staged.no_xbrl_filings[:50],
"no_xbrl_filings_count": len(staged.no_xbrl_filings),
"no_xbrl_ciks": sorted({
str(item["cik"])
for item in staged.no_xbrl_filings
if item.get("cik")
}),
"missing_xbrl": staged.missing_xbrl[:50],
"missing_xbrl_count": len(staged.missing_xbrl),
"missing_xbrl_blocking": len(blocking),
# Present only when the aggregate ceiling forced this run through, so
# a promoted run that carries known-unresolved filings says so.
"promotion_ceiling_tripped": self._ceiling_tripped,
"recovered_from_coregistrant": staged.recovered[:50],
"recovered_count": len(staged.recovered),
# Complete compact gate input; detailed audit lists above stay capped.
"setup_blocked_ciks": sorted({
str(item["cik"])
for item in [*staged.missing_xbrl, *staged.no_xbrl_filings]
if item.get("cik")
}),
"invalid_payloads": staged.invalid_payloads,
"cik_updates": len(staged.resolved.cik_updates),
# differing existing accessions (immutable — kept, reported here)
"discrepancies": staged.discrepancies[:50],
"discrepancy_count": len(staged.discrepancies),
}
return ValidationResult(
ok=not messages,
summary=summary,
source_max_date=self._latest_index_date,
messages=messages,
# Company-Facts absence is usually publication lag, but can also be a
# permanent co-registrant misfile that the daily index did not expose.
# Defer quietly at first; the framework warns if promotions stay stale.
retryable=(
len(messages) == 1
and bool(blocking)
and all(
m.get("reason") in {"not_in_companyfacts", "parser_unusable"}
for m in blocking
)
),
deferred_alert_after_days=MISSING_XBRL_RETRY_DAYS,
deferred_alert_messages=(
[
f"{len(aged_out)} tracked SEC filing(s) remain unresolved past "
f"the {MISSING_XBRL_RETRY_DAYS}-day retry window. They will "
f"enter automatic retry and block affected symbols from setups: "
f"{_missing_detail(aged_out)}"
]
if aged_out
else []
),
)
async def promote(self, db, staged: StagedFundamentals, run_id: int) -> dict[str, int]:
inserted = 0
updated = 0
# Only accessions whose reconstruction actually changed are rewritten;
# an unchanged stored row is left completely alone.
changed = {d["accession"] for d in staged.discrepancies} if self.reparse else set()
for row in staged.rows:
if row.accession in staged.existing_accessions:
if row.accession in changed:
# Write the FULL column set (_row_values covers _SNAPSHOT_COLS)
# so a rewritten row is never half old-parse, half new-parse.
# created_at stays at the original insert; import_run_id
# attributes the rewrite.
values = _row_values(row, run_id)
values.pop("created_at", None)
await db.execute(
update(FundamentalSnapshot)
.where(FundamentalSnapshot.accession == row.accession)
.values(**values)
)
updated += 1
continue # otherwise immutable — keep the original row
stmt = insert_for_session(db, FundamentalSnapshot).values(**_row_values(row, run_id))
stmt = stmt.on_conflict_do_nothing(index_elements=["accession"]) # race belt-and-suspenders
await db.execute(stmt)
inserted += 1
# Synchronize the retry queue in the snapshot-promotion transaction.
existing_gaps = (await db.execute(select(SecFilingGap))).scalars().all()
existing_gap_accessions = {gap.accession for gap in existing_gaps}
resolved_accessions = {row.accession for row in staged.rows}
# A filing now classified non-XBRL can never yield a snapshot and is no
# longer a fundamentals completeness gap.
resolved_accessions.update(
item["accession"] for item in staged.skipped_non_xbrl
)
queue_resolved = 0
if resolved_accessions:
result = await db.execute(
delete(SecFilingGap).where(
SecFilingGap.accession.in_(resolved_accessions)
)
)
queue_resolved = int(result.rowcount or 0)
now = _now()
tolerated = _past_retry_window(staged.missing_xbrl)
for gap in tolerated:
stmt = insert_for_session(db, SecFilingGap).values(
cik=gap["cik"],
accession=gap["accession"],
form=gap.get("form"),
index_date=gap.get("index_date"),
reason=gap["reason"],
coregistrant_ciks_json=json.dumps(gap.get("coregistrants") or []),
first_seen_at=now,
last_attempted_at=now,
)
await db.execute(
stmt.on_conflict_do_update(
index_elements=["accession"],
set_={
"cik": stmt.excluded.cik,
"form": stmt.excluded.form,
"index_date": stmt.excluded.index_date,
"reason": stmt.excluded.reason,
"coregistrant_ciks_json": stmt.excluded.coregistrant_ciks_json,
"last_attempted_at": stmt.excluded.last_attempted_at,
},
)
)
# Remove gaps made irrelevant by a later valid 10-K/10-Q. Quality reads
# already ignore them; physical cleanup keeps the queue small.
active_ids = {gap.id for gap in await fundamentals_quality_service.active_gaps(db)}
obsolete_ids = {
gap.id for gap in existing_gaps
if gap.id not in active_ids and gap.accession not in resolved_accessions
}
if obsolete_ids:
result = await db.execute(
delete(SecFilingGap).where(SecFilingGap.id.in_(obsolete_ids))
)
queue_resolved += int(result.rowcount or 0)
newly_queued = [
gap for gap in tolerated
if gap["accession"] not in existing_gap_accessions
]
# Warn (in-transaction, so it commits atomically with the promotion) when
# any existing accession reconstructed differently — kept immutable.
if staged.discrepancies:
accns = ", ".join(d["accession"] for d in staged.discrepancies[:10])
disposition = (
f"REWRITTEN by reparse run {run_id}" if self.reparse else "kept immutable"
)
db.add(SystemEvent(
severity="warning",
source="sec_facts",
code="snapshot_reparse" if self.reparse else "snapshot_discrepancy",
message=(
f"{len(staged.discrepancies)} stored accession(s) reconstructed "
f"differently; {disposition}: {accns}"
)[:4000],
dedup_key=f"sec_facts:discrepancy:{run_id}",
created_at=_now(),
))
# A ceiling-forced promotion is the safety valve firing — it must be
# visible, or the import silently starts carrying known-unresolved
# filings. The affected symbols stay barred from setups regardless.
if self._ceiling_tripped:
db.add(SystemEvent(
severity="warning",
source="sec_facts",
code="promotion_ceiling_forced",
message=(
f"Promoted with {self._ceiling_tripped['unresolved']} unresolved "
f"filing(s) — {self._ceiling_tripped['forced']} still inside the "
f"{MISSING_XBRL_RETRY_DAYS}-day retry window — because nothing had "
f"promoted in {PROMOTION_CEILING_DAYS} days. They are queued for "
"retry and their symbols remain blocked from setups."
)[:4000],
dedup_key=f"sec_facts:promotion_ceiling_forced:{run_id}",
created_at=now,
))
# Persistent current gaps get one actionable escalation rather than a
# daily warning. The nullable marker makes this durable and noise-free.
escalation_cutoff = now - timedelta(days=FILING_GAP_ESCALATE_DAYS)
aged_gaps = (
await db.execute(
select(SecFilingGap).where(
SecFilingGap.first_seen_at <= escalation_cutoff,
SecFilingGap.escalated_at.is_(None),
)
)
).scalars().all()
if aged_gaps:
named = ", ".join(
f"{gap.cik}/{gap.accession} ({gap.reason})"
for gap in aged_gaps[:10]
)
db.add(SystemEvent(
severity="warning",
source="sec_facts",
code="filing_gap_aged",
message=(
f"{len(aged_gaps)} SEC filing gap(s) remain unresolved after "
f"{FILING_GAP_ESCALATE_DAYS} days; affected setups remain paused. "
f"Review the filing/CIK mapping or parser: {named}"
)[:4000],
dedup_key=f"sec_facts:filing_gap_aged:{run_id}",
created_at=now,
))
await db.execute(
update(SecFilingGap)
.where(SecFilingGap.id.in_([gap.id for gap in aged_gaps]))
.values(escalated_at=now)
)
# Recovered rows are real data from an unexpected place — record where they
# came from, so a wrong recovery is auditable rather than invisible.
if staged.recovered:
named = ", ".join(
f"{r['accession']} <- CIK {r['source_cik']}" for r in staged.recovered[:10]
)
logger.info(
"sec_facts: recovered %d filing(s) from co-registrants: %s",
len(staged.recovered),
named,
)
# One warning when a gap first enters automatic retry. Repeating it every
# day adds noise; the queue remains the durable actionable state.
if newly_queued:
symbols_by_cik: dict[str, list[str]] = defaultdict(list)
for symbol, cik in staged.resolved.symbol_to_cik.items():
symbols_by_cik[cik10(cik)].append(symbol)
named = ", ".join(
f"{'/'.join(symbols_by_cik.get(gap['cik'], [])) or gap['cik']}"
f"/{gap['accession']}"
for gap in newly_queued[:10]
)
db.add(SystemEvent(
severity="warning",
source="sec_facts",
code="unresolved_filing",
message=(
f"{len(newly_queued)} filing(s) entered automatic SEC retry. "
f"Affected symbols are blocked from new actionable setups until "
f"their filing is recovered: {named}"
)[:4000],
dedup_key=f"sec_facts:unresolved_filing:{run_id}",
created_at=_now(),
))
# A new registrant may have no XBRL filing yet. Keep it out of actionable
# setups, but log it instead of raising a recurring operator warning.
if staged.no_xbrl_filings:
named = ", ".join(
f"{e['cik']} ({e.get('name') or '?'})" for e in staged.no_xbrl_filings[:10]
)
logger.info(
"sec_facts: %d registrant(s) have no XBRL history yet: %s",
len(staged.no_xbrl_filings),
named,
)
ticker_counts = await sec_universe.apply_ticker_updates(
db, staged.resolved, staged.sic_updates
)
return {
"inserted": inserted,
"updated": updated,
"existing_unchanged": len(staged.existing_accessions) - updated,
"discrepancies": len(staged.discrepancies),
"retry_queue_added": len(newly_queued),
"retry_queue_resolved": queue_resolved,
**ticker_counts,
}
# -- helpers -----------------------------------------------------------
async def _retry_backlog(
self,
db,
tracked_ciks: set[int],
) -> list[dict[str, Any]]:
"""Active typed gaps; migration 028 owns historical bootstrap."""
if not tracked_ciks:
return []
tracked = {cik10(cik) for cik in tracked_ciks}
candidates: dict[str, dict[str, Any]] = {}
queued = await fundamentals_quality_service.active_gaps(db, tracked)
for gap in queued:
try:
coregistrants = json.loads(gap.coregistrant_ciks_json or "[]")
except (TypeError, ValueError):
coregistrants = []
candidates[gap.accession] = {
"cik": gap.cik,
"accession": gap.accession,
"form": gap.form,
"index_date": gap.index_date,
"reason": gap.reason,
"coregistrants": coregistrants,
"_retry_queue": True,
}
if not candidates:
return []
resolved = set(
(
await db.execute(
select(FundamentalSnapshot.accession).where(
FundamentalSnapshot.accession.in_(list(candidates))
)
)
).scalars().all()
)
return [
item
for accession, item in candidates.items()
if accession not in resolved
]
async def _promotions_stale(self, db) -> bool:
"""Has nothing promoted within ``PROMOTION_CEILING_DAYS``?
Only true for a source that HAS promoted before. A never-promoted import
is initial setup, not a wedge: forcing its first promotion through would
mask a misconfiguration rather than recover from a transient SEC gap.
Measured from ``self.today`` rather than the wall clock, so the ceiling
honors the same injected date that ages the filings it releases.
"""
cutoff = datetime.combine(
self.today - timedelta(days=PROMOTION_CEILING_DAYS),
time.min,
tzinfo=timezone.utc,
)
ever, recent = (
await db.execute(
select(
exists().where(
DataImportRun.source == SOURCE,
DataImportRun.status == STATUS_PROMOTED,
),
exists().where(
DataImportRun.source == SOURCE,
DataImportRun.status == STATUS_PROMOTED,
DataImportRun.started_at >= cutoff,
),
)
)
).one()
return bool(ever) and not bool(recent)
async def _last_processed_index_date(self, db) -> date | None:
return (
await db.execute(
select(DataImportRun.source_max_date)
.where(DataImportRun.source == SOURCE, DataImportRun.status == STATUS_PROMOTED)
.order_by(DataImportRun.id.desc())
.limit(1)
)
).scalar_one_or_none()
async def _collect_index_rows(
self, client: SecClient, last_processed: date, latest: date
) -> list[dict[str, Any]]:
# Walk EVERY unprocessed date. No cap — dropping the older part of a long
# outage while still advancing source_max_date would permanently lose
# those filings. A large gap is one-time cost, not silent data loss.
tracked = set(self._resolved.cik_to_ticker_ids) if self._resolved else set()
gap = (latest - last_processed).days
if gap > 60:
logger.warning("sec_facts: %d-day index gap since %s; walking all", gap, last_processed)
rows: list[dict[str, Any]] = []
day = last_processed + timedelta(days=1)
while day <= latest:
# Group the whole day first: a combined filing is listed once per
# co-registrant CIK, and those sibling CIKs are the only pointer to
# where SEC may have put the XBRL (see _recover_from_coregistrant).
by_accession: dict[str, list[dict[str, Any]]] = defaultdict(list)
for r in await client.daily_index(day):
if r["form"] in _XBRL_FORMS:
by_accession[r["accession"]].append(r)
for accession, group in by_accession.items():
filers = {r["cik"] for r in group}
tracked_filers = filers & tracked
if not tracked_filers:
continue
siblings = sorted(filers - tracked_filers)
if siblings:
self._coregistrants[accession] = siblings
for r in group:
if r["cik"] in tracked_filers:
r["index_date"] = day # not hashed (revision uses cik/accession)
rows.append(r)
day += timedelta(days=1)
return rows
async def _ciks_with_snapshots(self, db, ciks: set[int]) -> set[int]:
if not ciks:
return set()
cik_strs = [cik10(c) for c in ciks]
found = (
await db.execute(
select(FundamentalSnapshot.cik)
.where(FundamentalSnapshot.cik.in_(cik_strs))
.distinct()
)
).scalars().all()
return {int(c) for c in found}
async def _last_shares_outstanding(self, db, ciks: set[int]) -> dict[str, float]:
"""Latest known shares outstanding per tracked issuer — the continuity
reference co-registrant recovery is checked against."""
if not ciks:
return {}
rows = (
await db.execute(
select(FundamentalSnapshot.cik, FundamentalSnapshot.shares_outstanding)
.where(
FundamentalSnapshot.cik.in_([cik10(c) for c in ciks]),
FundamentalSnapshot.shares_outstanding.is_not(None),
)
# Last write per cik wins, so the sort must be total: an amendment
# and its original share a period_end, and an undefined tie there
# would make recovery non-deterministic across runs and dialects.
.order_by(
FundamentalSnapshot.period_end,
FundamentalSnapshot.filed_date,
FundamentalSnapshot.accession,
)
)
).all()
return {cik: float(shares) for cik, shares in rows}
async def _existing_by_accession(self, db, accessions: list[str]) -> dict[str, FundamentalSnapshot]:
if not accessions:
return {}
rows = (
await db.execute(
select(FundamentalSnapshot).where(FundamentalSnapshot.accession.in_(accessions))
)
).scalars().all()
return {r.accession: r for r in rows}
def _companyfacts_structure_error(cf: Any) -> str | None:
"""None if the payload is structurally sound, else a reason string. Checks the
top-level ``facts`` mapping AND that every concept carries a ``units`` mapping —
a missing/non-dict units would silently drop that concept's facts otherwise."""
if not isinstance(cf, dict) or not isinstance(cf.get("facts"), dict):
return "missing facts structure"
for concepts in cf["facts"].values():
if not isinstance(concepts, dict):
return "malformed taxonomy structure"
for body in concepts.values():
if not isinstance(body, dict) or not isinstance(body.get("units"), dict):
return "missing units structure"
return None
def _missing(
cik: int,
row: dict[str, Any],
reason: str,
today: date,
coregistrants: list[int] | None = None,
) -> dict[str, Any]:
"""One unresolvable index row, carrying everything needed to look the filing
up by hand (EDGAR accession + the index date it was seen on) and to decide
whether it is still young enough to be worth blocking on."""
index_date = row.get("index_date")
age_days = (
(today - index_date).days if isinstance(index_date, date) else 0
)
if row.get("_retry_queue"):
age_days = max(age_days, MISSING_XBRL_RETRY_DAYS + 1)
return {
"cik": cik10(cik),
"accession": row["accession"],
"form": row.get("form"),
"index_date": index_date,
# A newly observed row without a date blocks safely. A durable queue row
# has already passed the bounded window and is forced aged-out above.
"age_days": age_days,
"reason": reason,
"coregistrants": list(coregistrants or []),
}
def _within_retry_window(missing: list[dict[str, Any]]) -> list[dict[str, Any]]:
return [m for m in missing if m.get("age_days", 0) <= MISSING_XBRL_RETRY_DAYS]
def _past_retry_window(missing: list[dict[str, Any]]) -> list[dict[str, Any]]:
return [m for m in missing if m.get("age_days", 0) > MISSING_XBRL_RETRY_DAYS]
def _shares_continuous(shares: float | None, reference: float | None) -> bool:
"""Does a co-registrant-recovered share count look like this issuer's own?
The failure worth preventing is storing a subsidiary's standalone facts as the
parent's. A co-registrant shell holds a token float — Florida Power & Light
against NextEra's 2.09bn shares — so any sane band separates them while still
tolerating buybacks and issuance. With no history to compare against (a newly
tracked issuer) or no share count at all, recovery is refused, not guessed.
"""
if not shares or not reference:
return False
return RECOVERY_SHARES_MIN <= shares / reference <= RECOVERY_SHARES_MAX
def _reason_counts(missing: list[dict[str, Any]]) -> str:
counts = Counter(m["reason"] for m in missing)
return ", ".join(f"{reason}={n}" for reason, n in sorted(counts.items()))
def _missing_detail(missing: list[dict[str, Any]], limit: int = 10) -> str:
detail = ", ".join(
f"{m['cik']}/{m['accession']} {m.get('form') or '?'} "
f"[{m.get('index_date') or '?'}] {m['reason']}"
for m in missing[:limit]
)
if len(missing) > limit:
detail += f", +{len(missing) - limit} more"
return detail
def _filing_meta(sub: dict[str, Any]) -> tuple[dict[str, FilingMeta], set[str]]:
"""(xbrl_meta, nonxbrl_accessions) from a submissions payload. xbrl_meta only
includes 10-K/10-Q(/A) filings that are XBRL and have full period metadata."""
xbrl: dict[str, FilingMeta] = {}
nonxbrl: set[str] = set()
for f in sub.get("filings", []):
if f["form"] not in _XBRL_FORMS:
continue
if not f.get("is_xbrl"):
nonxbrl.add(f["accession"])
continue
if not (f.get("report_date") and f.get("filing_date") and f.get("acceptance_datetime")):
continue
xbrl[f["accession"]] = FilingMeta(
report_date=date.fromisoformat(f["report_date"]),
filing_date=date.fromisoformat(f["filing_date"]),
accepted_at=_parse_dt(f["acceptance_datetime"]),
form=f["form"],
)
return xbrl, nonxbrl
def _parse_dt(value: str) -> datetime:
return datetime.fromisoformat(value.replace("Z", "+00:00"))
def _row_values(row: SnapshotRow, run_id: int) -> dict[str, Any]:
values = {col: getattr(row, col) for col in _SNAPSHOT_COLS}
values["import_run_id"] = run_id
values["created_at"] = _now()
return values
def _diff_fields(row: SnapshotRow, old: FundamentalSnapshot) -> list[str]:
"""Source fields where a re-parsed row differs from the stored row."""
return [
col for col in _COMPARE_COLS
if not _same_value(getattr(row, col), getattr(old, col))
]
def _same_value(parsed: Any, stored: Any) -> bool:
"""Compare a freshly parsed value against its stored round-trip.
Datetimes need care: every timestamp here is UTC by construction, but
``DateTime(timezone=True)`` only preserves tzinfo on Postgres — SQLite hands
back a naive value. Comparing representations would report an unchanged row
as differing, which would both spam the discrepancy warning and make a
reparse rewrite every row it touched. Compare instants instead.
"""
if isinstance(parsed, datetime) and isinstance(stored, datetime):
return _as_utc(parsed) == _as_utc(stored)
return parsed == stored
def _as_utc(value: datetime) -> datetime:
return value if value.tzinfo is not None else value.replace(tzinfo=timezone.utc)