--- a/app/workers/celery_app.py +++ b/app/workers/celery_app.py @@ -26,8 +26,14 @@ task_acks_late=settings.CELERY_TASK_ACKS_LATE, task_time_limit=settings.CELERY_TASK_TIME_LIMIT_SECONDS, task_soft_time_limit=settings.CELERY_TASK_SOFT_TIME_LIMIT_SECONDS, + imports=( + "app.workers.discovery", + "app.workers.ingest", + "app.workers.peer_groups", + "app.workers.recompute", + "app.workers.token_cleanup", + ), ) -celery_app.autodiscover_tasks(["app.workers"]) celery_app.conf.beat_schedule = { "daily-universe-ingestion": { --- a/app/workers/ingest.py +++ b/app/workers/ingest.py @@ -25,6 +25,7 @@ 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 @@ -480,7 +481,15 @@ 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) + 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) @@ -488,7 +497,15 @@ _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) + 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) @@ -496,7 +513,15 @@ _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) + 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) @@ -504,7 +529,15 @@ _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) + 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) --- a/app/adapters/resolver.py +++ b/app/adapters/resolver.py @@ -99,8 +99,35 @@ result = resolve_conflict(candidates) if result.selected_value is not None: resolved[field] = result.selected_value - if result.is_conflicting: - conflicts.append(FieldConflict(field, val_a, val_b, result.selected_value, result.selection_reason)) + + # Secondary-only provenance must be observable. A NULL in the + # primary provider followed by a value from the secondary provider + # is not a numeric conflict, but it is still a cross-provider merge. + if ( + secondary is not None + and secondary_tier is not None + and val_a is None + and val_b is not None + ): + conflicts.append( + FieldConflict( + field, + val_a, + val_b, + result.selected_value, + "Primary field missing; value supplied by secondary source", + ) + ) + elif result.is_conflicting: + conflicts.append( + FieldConflict( + field, + val_a, + val_b, + result.selected_value, + result.selection_reason, + ) + ) return resolved, conflicts --- a/app/adapters/validation.py +++ b/app/adapters/validation.py @@ -232,7 +232,47 @@ out[name] = number # Cross-field checks: each is an identity that must hold, not a heuristic. - revenue, gross_profit = out.get("revenue"), out.get("gross_profit") + revenue = out.get("revenue") + cogs = out.get("cogs") + gross_profit = out.get("gross_profit") + + # Accounting identity: + # gross_profit = revenue - cogs + # + # The provider resolver already treats <=0.5% financial differences + # as non-conflicting. Use the same tolerance here so ordinary provider + # rounding does not create a false integrity failure. + if revenue is not None and cogs is not None and gross_profit is not None: + expected_gross_profit = revenue - cogs + + if expected_gross_profit == 0: + inconsistent = abs(gross_profit) > 0.01 + else: + inconsistent = ( + abs(gross_profit - expected_gross_profit) + / abs(expected_gross_profit) + > 0.005 + ) + + if inconsistent: + report.add( + ValidationIssue( + "gross_profit", + gross_profit, + ( + f"does not reconcile with revenue ({revenue}) " + f"- cogs ({cogs}) = {expected_gross_profit}" + ), + SEVERITY_WARNING, + row_key, + ) + ) + + # Do not allow an internally inconsistent dependent value + # into the persisted merged statement. + del out["gross_profit"] + gross_profit = None + if revenue is not None and gross_profit is not None and gross_profit > revenue: report.add(ValidationIssue("gross_profit", gross_profit, f"exceeds revenue ({revenue})", SEVERITY_WARNING, row_key)) --- a/app/schemas/common.py +++ b/app/schemas/common.py @@ -96,7 +96,7 @@ class ScreenFilter(BaseModel): - metric: str + metric: str = Field(max_length=48) op: str # gt|gte|lt|lte|eq|between value: float value2: Optional[float] = None # for "between" @@ -115,7 +115,7 @@ class ScreenRequest(BaseModel): universe: ScreenUniverse = ScreenUniverse() logic: str = "AND" - filters: list[ScreenFilter] = [] + filters: list[ScreenFilter] = Field(default=[], max_length=20) sort_by: str = "overall_score" sort_direction: str = "desc" # AUDIT FIX (StockLab overhaul, performance audit, docs/AUDIT_PERFORMANCE.md finding #1): no