Compare commits

...
3 Commits
Author SHA1 Message Date
dennisthiessen 7bcdf77ef9 fix(sec): name aged filings in deferred warnings
Deploy / lint (push) Successful in 10s
Deploy / test (push) Successful in 2m3s
Deploy / deploy (push) Successful in 42s
2026-07-31 14:58:22 +02:00
dennisthiessen c8c660e63d fix(sec): warn when deferred imports stay stale 2026-07-31 13:27:31 +02:00
dennisthiessen f58f8b0818 fix(sec): defer expected Company Facts lag without alerting 2026-07-31 12:40:29 +02:00
8 changed files with 283 additions and 20 deletions
+4 -2
View File
@@ -10,7 +10,8 @@ class DataImportRun(Base):
"""One row per bulk-import attempt (SEC facts / Dolt earnings / Dolt stocks). """One row per bulk-import attempt (SEC facts / Dolt earnings / Dolt stocks).
Lean audit record for the batch import framework: every attempt is logged, Lean audit record for the batch import framework: every attempt is logged,
whether it promoted, was a ``no_op`` (unchanged revision), or ``failed``. whether it promoted, was a ``no_op`` (unchanged revision), was ``deferred``
for an expected retry, or ``failed``.
``row_counts`` and ``validation`` hold JSON strings (repo convention — see ``row_counts`` and ``validation`` hold JSON strings (repo convention — see
``fundamental_data.unavailable_fields_json``), not JSONB; the validation ``fundamental_data.unavailable_fields_json``), not JSONB; the validation
blob carries reconciliation/discrepancy summaries so no separate conflicts blob carries reconciliation/discrepancy summaries so no separate conflicts
@@ -28,7 +29,7 @@ class DataImportRun(Base):
source: Mapped[str] = mapped_column(String(32), nullable=False) source: Mapped[str] = mapped_column(String(32), nullable=False)
# Dolt commit hash, or SEC archive SHA-256. Null until known. # Dolt commit hash, or SEC archive SHA-256. Null until known.
revision: Mapped[str | None] = mapped_column(String(64), nullable=True) revision: Mapped[str | None] = mapped_column(String(64), nullable=True)
# running | validated | promoted | no_op | failed # running | validated | promoted | no_op | deferred | failed
status: Mapped[str] = mapped_column(String(16), nullable=False) status: Mapped[str] = mapped_column(String(16), nullable=False)
source_max_date: Mapped[date | None] = mapped_column(Date, nullable=True) source_max_date: Mapped[date | None] = mapped_column(Date, nullable=True)
row_counts_json: Mapped[str | None] = mapped_column(Text, nullable=True) row_counts_json: Mapped[str | None] = mapped_column(Text, nullable=True)
@@ -39,4 +40,5 @@ class DataImportRun(Base):
completed_at: Mapped[datetime | None] = mapped_column( completed_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True DateTime(timezone=True), nullable=True
) )
# Failure detail, or the non-error reason when status is deferred.
error_details: Mapped[str | None] = mapped_column(Text, nullable=True) error_details: Mapped[str | None] = mapped_column(Text, nullable=True)
+12 -2
View File
@@ -43,7 +43,12 @@ from app.services import (
fundamentals_parity_service, fundamentals_parity_service,
fundamental_data_refresh_service, fundamental_data_refresh_service,
) )
from app.services.data_import import STATUS_FAILED, SourceImporter, run_import from app.services.data_import import (
STATUS_DEFERRED,
STATUS_FAILED,
SourceImporter,
run_import,
)
from app.services.dolt_earnings_importer import DoltEarningsImporter from app.services.dolt_earnings_importer import DoltEarningsImporter
from app.services.sec_fundamentals_importer import SecFundamentalsImporter from app.services.sec_fundamentals_importer import SecFundamentalsImporter
from app.services.alert_service import dispatch_alerts from app.services.alert_service import dispatch_alerts
@@ -945,7 +950,7 @@ async def _run_shadow_import(job_name: str, importer: SourceImporter) -> bool:
"""Run an importer and return whether its scheduled job was enabled. """Run an importer and return whether its scheduled job was enabled.
The SEC wrapper uses the return value to run its activated local cache step The SEC wrapper uses the return value to run its activated local cache step
after failed, no-op, promoted, or source-locked attempts while still honoring after deferred, failed, no-op, promoted, or source-locked attempts while honoring
the job-level disable switch. the job-level disable switch.
""" """
_log_event(logging.INFO, "job_start", job=job_name) _log_event(logging.INFO, "job_start", job=job_name)
@@ -967,6 +972,11 @@ async def _run_shadow_import(job_name: str, importer: SourceImporter) -> bool:
revision = f" · {run.revision[:12]}" if run.revision else "" revision = f" · {run.revision[:12]}" if run.revision else ""
message = f"{run.status}{revision}" message = f"{run.status}{revision}"
if run.status == STATUS_DEFERRED:
message = run.error_details or message
_log_event(logging.INFO, "job_deferred", job=job_name, message=message)
_runtime_finish(job_name, "deferred", processed=0, total=1, message=message)
return True
if run.status == STATUS_FAILED: if run.status == STATUS_FAILED:
message = run.error_details or message message = run.error_details or message
_log_event(logging.ERROR, "job_error", job=job_name, message=message) _log_event(logging.ERROR, "job_error", job=job_name, message=message)
+78 -6
View File
@@ -30,10 +30,10 @@ import hashlib
import json import json
import logging import logging
from dataclasses import dataclass, field from dataclasses import dataclass, field
from datetime import date, datetime, timezone from datetime import date, datetime, timedelta, timezone
from typing import Any, Protocol, runtime_checkable from typing import Any, Protocol, runtime_checkable
from sqlalchemy import select, text from sqlalchemy import exists, select, text
from sqlalchemy.engine import Engine # noqa: F401 (typing only) from sqlalchemy.engine import Engine # noqa: F401 (typing only)
from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession
@@ -48,6 +48,7 @@ STATUS_RUNNING = "running"
STATUS_VALIDATED = "validated" STATUS_VALIDATED = "validated"
STATUS_PROMOTED = "promoted" STATUS_PROMOTED = "promoted"
STATUS_NO_OP = "no_op" STATUS_NO_OP = "no_op"
STATUS_DEFERRED = "deferred"
STATUS_FAILED = "failed" STATUS_FAILED = "failed"
_MAX_ERROR_LEN = 4000 _MAX_ERROR_LEN = 4000
@@ -68,6 +69,12 @@ class ValidationResult:
summary: dict[str, Any] = field(default_factory=dict) summary: dict[str, Any] = field(default_factory=dict)
source_max_date: date | None = None source_max_date: date | None = None
messages: list[str] = field(default_factory=list) messages: list[str] = field(default_factory=list)
# Expected source-side lag: retry without an immediate error alert. Sources
# can bound the quiet period with deferred_alert_after_days. Only meaningful
# when ok=False.
retryable: bool = False
deferred_alert_after_days: int | None = None
deferred_alert_messages: list[str] = field(default_factory=list)
@runtime_checkable @runtime_checkable
@@ -123,19 +130,46 @@ async def _last_promoted_revision(db: AsyncSession, source: str) -> str | None:
return row.scalar_one_or_none() return row.scalar_one_or_none()
async def _promotion_state_since(
db: AsyncSession, source: str, cutoff: datetime
) -> str:
promoted = (
DataImportRun.source == source,
DataImportRun.status == STATUS_PROMOTED,
)
ever, recent = (
await db.execute(
select(
exists().where(*promoted),
exists().where(*promoted, DataImportRun.started_at >= cutoff),
)
)
).one()
return "recent" if recent else "stale" if ever else "never"
def _now() -> datetime: def _now() -> datetime:
return datetime.now(timezone.utc) return datetime.now(timezone.utc)
async def _alert(db: AsyncSession, source: str, code: str, messages: list[str]) -> None: async def _alert(
db: AsyncSession,
source: str,
code: str,
messages: list[str],
*,
severity: str = "error",
dedup_hours: int = 24,
) -> None:
try: try:
await system_event_service.log_event( await system_event_service.log_event(
db, db,
severity="error", severity=severity,
source="data_import", source="data_import",
code=f"{source}_{code}", code=f"{source}_{code}",
message=(("; ".join(messages)) or code)[:_MAX_ERROR_LEN], message=(("; ".join(messages)) or code)[:_MAX_ERROR_LEN],
dedup_key=f"data_import:{source}:{code}", dedup_key=f"data_import:{source}:{code}",
dedup_hours=dedup_hours,
) )
except Exception: # noqa: BLE001 — alerting must never mask the real outcome except Exception: # noqa: BLE001 — alerting must never mask the real outcome
logger.exception("Failed to emit data_import alert %s/%s", source, code) logger.exception("Failed to emit data_import alert %s/%s", source, code)
@@ -149,7 +183,8 @@ async def run_import(
) -> DataImportRun | None: ) -> DataImportRun | None:
"""Run one import for ``importer``. """Run one import for ``importer``.
Returns the recorded ``DataImportRun`` (promoted / no_op / failed), or None Returns the recorded ``DataImportRun`` (promoted / no_op / deferred /
failed), or None
when the per-source advisory lock is already held (another run is active). when the per-source advisory lock is already held (another run is active).
``force`` runs even when the revision is unchanged. The revision tracks the ``force`` runs even when the revision is unchanged. The revision tracks the
@@ -208,9 +243,46 @@ async def run_import(
run.validation_json = json.dumps(result.summary, default=str) run.validation_json = json.dumps(result.summary, default=str)
if not result.ok: if not result.ok:
run.status = STATUS_FAILED
run.error_details = ("; ".join(result.messages))[:_MAX_ERROR_LEN] run.error_details = ("; ".join(result.messages))[:_MAX_ERROR_LEN]
run.completed_at = _now() run.completed_at = _now()
if result.retryable:
run.status = STATUS_DEFERRED
await session.commit()
alert_days = result.deferred_alert_after_days
if alert_days is not None:
alert_days = max(1, alert_days)
cutoff = run.started_at - timedelta(days=alert_days)
promotion_state = await _promotion_state_since(
session, source, cutoff
)
if promotion_state != "recent":
history = (
f"{source} import has never promoted successfully"
if promotion_state == "never"
else f"{source} import has not promoted successfully "
f"within {alert_days} day(s)"
)
await _alert(
session,
source,
"deferred_stale",
[
f"{history}; import remains deferred",
*result.deferred_alert_messages,
f"Current deferral: "
f"{run.error_details or 'validation deferred'}",
],
severity="warning",
dedup_hours=alert_days * 24,
)
logger.info(
"data_import %s: deferred for retry: %s",
source,
result.messages,
)
return run
run.status = STATUS_FAILED
await session.commit() await session.commit()
await _alert(session, source, "validation_failed", result.messages) await _alert(session, source, "validation_failed", result.messages)
logger.warning( logger.warning(
+20
View File
@@ -362,6 +362,7 @@ class SecFundamentalsImporter:
# the filings: "which ones" has to be in the alert itself, not merely # the filings: "which ones" has to be in the alert itself, not merely
# reconstructible by re-walking the index. # reconstructible by re-walking the index.
blocking = _within_retry_window(staged.missing_xbrl) blocking = _within_retry_window(staged.missing_xbrl)
aged_out = _past_retry_window(staged.missing_xbrl)
if blocking: if blocking:
messages.append( messages.append(
f"{len(blocking)} tracked XBRL filing(s) unresolved within the " f"{len(blocking)} tracked XBRL filing(s) unresolved within the "
@@ -413,6 +414,25 @@ class SecFundamentalsImporter:
summary=summary, summary=summary,
source_max_date=self._latest_index_date, source_max_date=self._latest_index_date,
messages=messages, 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") == "not_in_companyfacts" 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 and risk being "
f"promoted around without automatic retry: "
f"{_missing_detail(aged_out)}"
]
if aged_out
else []
),
) )
async def promote(self, db, staged: StagedFundamentals, run_id: int) -> dict[str, int]: async def promote(self, db, staged: StagedFundamentals, run_id: int) -> dict[str, int]:
@@ -25,7 +25,7 @@ function formatAgo(iso: string | null | undefined): string {
function lastRunColor(status: string | null | undefined): string { function lastRunColor(status: string | null | undefined): string {
if (status === 'error') return 'text-red-300'; if (status === 'error') return 'text-red-300';
if (status === 'rate_limited') return 'text-amber-300'; if (status === 'rate_limited' || status === 'deferred') return 'text-amber-300';
return 'text-gray-500'; return 'text-gray-500';
} }
@@ -127,7 +127,7 @@ export function JobControls() {
className={`text-[11px] font-medium ${ className={`text-[11px] font-medium ${
job.running job.running
? 'text-blue-300' ? 'text-blue-300'
: job.runtime_status === 'rate_limited' : job.runtime_status === 'rate_limited' || job.runtime_status === 'deferred'
? 'text-amber-300' ? 'text-amber-300'
: job.runtime_status === 'error' : job.runtime_status === 'error'
? 'text-red-300' ? 'text-red-300'
@@ -140,6 +140,8 @@ export function JobControls() {
? 'Running' ? 'Running'
: job.runtime_status === 'rate_limited' : job.runtime_status === 'rate_limited'
? 'Paused (rate-limited)' ? 'Paused (rate-limited)'
: job.runtime_status === 'deferred'
? 'Deferred (retrying)'
: job.runtime_status === 'error' : job.runtime_status === 'error'
? 'Last run error' ? 'Last run error'
: job.enabled : job.enabled
+67 -2
View File
@@ -17,7 +17,7 @@ from __future__ import annotations
import asyncio import asyncio
import os import os
import tempfile import tempfile
from datetime import date, datetime, timezone from datetime import date, datetime, timedelta, timezone
import pytest import pytest
from sqlalchemy import func, select from sqlalchemy import func, select
@@ -29,6 +29,7 @@ from app.models.data_import_run import DataImportRun
from app.models.fundamental_snapshot import FundamentalSnapshot from app.models.fundamental_snapshot import FundamentalSnapshot
from app.models.system_event import SystemEvent from app.models.system_event import SystemEvent
from app.services.data_import import ( from app.services.data_import import (
STATUS_DEFERRED,
STATUS_FAILED, STATUS_FAILED,
STATUS_NO_OP, STATUS_NO_OP,
STATUS_PROMOTED, STATUS_PROMOTED,
@@ -67,9 +68,16 @@ class FakeImporter:
source = "sec_facts" source = "sec_facts"
def __init__(self, revision, *, ok=True, n_rows=3, raise_in="none"): def __init__(
self, revision, *, ok=True, retryable=False, alert_days=None,
n_rows=3, raise_in="none",
alert_messages=None,
):
self.revision = revision self.revision = revision
self.ok = ok self.ok = ok
self.retryable = retryable
self.alert_days = alert_days
self.alert_messages = alert_messages or []
self.n_rows = n_rows self.n_rows = n_rows
self.raise_in = raise_in self.raise_in = raise_in
self.staged_called = False self.staged_called = False
@@ -94,6 +102,9 @@ class FakeImporter:
summary={"staged_rows": len(staged)}, summary={"staged_rows": len(staged)},
source_max_date=date(2026, 7, 21), source_max_date=date(2026, 7, 21),
messages=[] if self.ok else ["coverage below threshold"], messages=[] if self.ok else ["coverage below threshold"],
retryable=self.retryable,
deferred_alert_after_days=self.alert_days,
deferred_alert_messages=self.alert_messages,
) )
async def promote(self, db, staged, run_id): async def promote(self, db, staged, run_id):
@@ -200,6 +211,60 @@ async def test_failed_validation_leaves_data_untouched(engine):
assert await _count(factory, SystemEvent) == 1 # alerted assert await _count(factory, SystemEvent) == 1 # alerted
async def test_retryable_validation_defers_without_alerting(engine):
factory = _factory(engine)
await run_import(FakeImporter("rev1", n_rows=3), engine=engine)
run = await run_import(
FakeImporter("rev2", ok=False, retryable=True, n_rows=5), engine=engine
)
assert run is not None and run.status == STATUS_DEFERRED
assert "coverage" in (run.error_details or "")
assert await _count(factory, FundamentalSnapshot) == 3 # untouched
assert await _count(factory, SystemEvent) == 0 # expected retry does not alert
async def test_stale_deferred_validation_emits_deduplicated_warning(engine):
factory = _factory(engine)
promoted = await run_import(FakeImporter("rev1", n_rows=3), engine=engine)
async with factory() as s:
promoted.started_at = datetime.now(timezone.utc) - timedelta(days=4)
await s.merge(promoted)
await s.commit()
importer = FakeImporter(
"rev2", ok=False, retryable=True, alert_days=3,
alert_messages=["source detail names OLD-ACCESSION"],
n_rows=5,
)
first = await run_import(importer, engine=engine)
second = await run_import(importer, engine=engine)
assert first is not None and first.status == STATUS_DEFERRED
assert second is not None and second.status == STATUS_DEFERRED
async with factory() as s:
events = (await s.execute(select(SystemEvent))).scalars().all()
assert len(events) == 1
assert events[0].severity == "warning"
assert events[0].code == "sec_facts_deferred_stale"
assert "OLD-ACCESSION" in events[0].message
assert "aged-out" not in events[0].message
async def test_never_promoted_deferred_warning_says_never(engine):
factory = _factory(engine)
run = await run_import(
FakeImporter("rev1", ok=False, retryable=True, alert_days=3),
engine=engine,
)
assert run is not None and run.status == STATUS_DEFERRED
async with factory() as s:
event = (await s.execute(select(SystemEvent))).scalar_one()
assert "has never promoted successfully" in event.message
async def test_exception_in_promote_rolls_back(engine): async def test_exception_in_promote_rolls_back(engine):
factory = _factory(engine) factory = _factory(engine)
await run_import(FakeImporter("rev1", n_rows=3), engine=engine) # baseline await run_import(FakeImporter("rev1", n_rows=3), engine=engine) # baseline
+23
View File
@@ -20,6 +20,7 @@ from app.scheduler import (
queue_backtest_target_model, queue_backtest_target_model,
scheduler, scheduler,
) )
from app.services.data_import import STATUS_DEFERRED
def test_manual_backtest_target_model_is_one_shot(): def test_manual_backtest_target_model_is_one_shot():
@@ -250,6 +251,28 @@ class TestShadowImportJobs:
assert runtime["processed"] == 0 assert runtime["processed"] == 0
assert runtime["message"] == "validation failed" assert runtime["message"] == "validation failed"
async def test_deferred_run_is_visible_without_error_status(self, monkeypatch):
async def enabled(db, job_name):
return True
async def imported(importer):
return SimpleNamespace(
status=STATUS_DEFERRED,
revision="abcdef1234567890",
error_details="Company Facts publication lag; retrying",
)
monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory)
monkeypatch.setattr("app.scheduler._is_job_enabled", enabled)
monkeypatch.setattr("app.scheduler.run_import", imported)
await _run_shadow_import("sec_fundamentals_import", object())
runtime = get_job_runtime_snapshot("sec_fundamentals_import")
assert runtime["status"] == STATUS_DEFERRED
assert runtime["processed"] == 0
assert runtime["message"] == "Company Facts publication lag; retrying"
async def test_source_lock_surfaces_skipped(self, monkeypatch): async def test_source_lock_surfaces_skipped(self, monkeypatch):
async def enabled(db, job_name): async def enabled(db, job_name):
return True return True
+75 -6
View File
@@ -17,8 +17,17 @@ import app.models # noqa: F401
from app.models.fundamental_snapshot import FundamentalSnapshot from app.models.fundamental_snapshot import FundamentalSnapshot
from app.models.system_event import SystemEvent from app.models.system_event import SystemEvent
from app.models.ticker import Ticker from app.models.ticker import Ticker
from app.services.data_import import STATUS_FAILED, STATUS_PROMOTED, run_import from app.services.data_import import (
from app.services.sec_fundamentals_importer import SecFundamentalsImporter STATUS_DEFERRED,
STATUS_FAILED,
STATUS_PROMOTED,
run_import,
)
from app.services.sec_fundamentals_importer import (
SecFundamentalsImporter,
StagedFundamentals,
)
from app.services.sec_universe import ResolvedUniverse
@pytest.fixture @pytest.fixture
@@ -191,7 +200,7 @@ async def test_incremental_adds_only_new_filing(engine):
assert q2.fiscal_period == "Q2" and q2.revenue == 254940 assert q2.fiscal_period == "Q2" and q2.revenue == 254940
async def test_consistency_gate_fails_when_facts_lag_index(engine): async def test_consistency_gate_defers_without_alert_when_facts_lag_index(engine):
factory = _factory(engine) factory = _factory(engine)
await _seed(factory, ["AAPL"]) await _seed(factory, ["AAPL"])
backfill = FakeSecClient( backfill = FakeSecClient(
@@ -213,9 +222,9 @@ async def test_consistency_gate_fails_when_facts_lag_index(engine):
) )
run = await run_import(_importer(incr, today=date(2026, 5, 3)), engine=engine) run = await run_import(_importer(incr, today=date(2026, 5, 3)), engine=engine)
assert run.status == STATUS_FAILED assert run.status == STATUS_DEFERRED
# The gate blocks every later run until it clears, so the alert itself has to # The gate blocks every later run until it clears, so run history still has
# name the filing and say why it could not be resolved. # to name the filing and say why it could not be resolved.
details = run.error_details or "" details = run.error_details or ""
assert "GHOST" in details and "not_in_companyfacts" in details assert "GHOST" in details and "not_in_companyfacts" in details
assert "2026-05-01" in details # index date the filing was seen on assert "2026-05-01" in details # index date the filing was seen on
@@ -224,6 +233,66 @@ async def test_consistency_gate_fails_when_facts_lag_index(engine):
assert summary["missing_xbrl"][0]["accession"] == "GHOST" assert summary["missing_xbrl"][0]["accession"] == "GHOST"
assert summary["missing_xbrl"][0]["form"] == "10-Q" assert summary["missing_xbrl"][0]["form"] == "10-Q"
assert await _count(factory, FundamentalSnapshot) == 2 # nothing new written assert await _count(factory, FundamentalSnapshot) == 2 # nothing new written
assert await _count(factory, SystemEvent) == 0 # expected SEC lag does not alert
async def test_companyfacts_lag_does_not_mask_second_validation_failure():
importer = SecFundamentalsImporter(today=date(2026, 5, 3))
importer._latest_index_date = date(2026, 5, 2)
staged = StagedFundamentals(
resolved=ResolvedUniverse(),
missing_xbrl=[{
"cik": "0000320193",
"accession": "GHOST",
"form": "10-Q",
"index_date": date(2026, 5, 1),
"age_days": 2,
"reason": "not_in_companyfacts",
}],
invalid_payloads=[{
"cik": "0000789019",
"reason": "missing facts structure",
}],
)
result = await importer.validate(None, staged)
assert not result.ok
assert not result.retryable
assert len(result.messages) == 2
async def test_deferred_alert_names_aged_out_accessions_separately():
importer = SecFundamentalsImporter(today=date(2026, 5, 6))
importer._latest_index_date = date(2026, 5, 5)
staged = StagedFundamentals(
resolved=ResolvedUniverse(),
missing_xbrl=[
{
"cik": "0000320193",
"accession": "YOUNG",
"form": "10-Q",
"index_date": date(2026, 5, 5),
"age_days": 1,
"reason": "not_in_companyfacts",
},
{
"cik": "0000789019",
"accession": "AGED-OUT",
"form": "10-Q",
"index_date": date(2026, 5, 1),
"age_days": 5,
"reason": "not_in_companyfacts",
},
],
)
result = await importer.validate(None, staged)
assert result.retryable
assert len(result.messages) == 1 and "YOUNG" in result.messages[0]
assert len(result.deferred_alert_messages) == 1
assert "AGED-OUT" in result.deferred_alert_messages[0]
async def test_gate_separates_missing_submissions_from_missing_facts(engine): async def test_gate_separates_missing_submissions_from_missing_facts(engine):