--- a/app/adapters/resolver.py +++ b/app/adapters/resolver.py @@ -97,12 +97,39 @@ if secondary is not None and secondary_tier is not None: candidates.append((secondary_tier, val_b)) result = resolve_conflict(candidates) if result.selected_value is not None: resolved[field] = result.selected_value - 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 def merge_statement_line_items( income: dict, --- a/app/adapters/validation.py +++ b/app/adapters/validation.py @@ -230,11 +230,51 @@ SEVERITY_WARNING, row_key)) continue 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)) total_assets, current_assets = out.get("total_assets"), out.get("current_assets") if total_assets is not None and current_assets is not None and current_assets > total_assets: --- a/app/api/v1/rankings.py +++ b/app/api/v1/rankings.py @@ -9,11 +9,11 @@ """ from typing import Optional from fastapi import APIRouter, Depends, Query, Request -from sqlalchemy import select +from sqlalchemy import func, select, and_ from sqlalchemy.orm import Session from app.api.v1.deps import limiter from app.api.v1.serializers import company_summary, is_demo_by_security, security_eager_load_options from app.core.config import get_settings @@ -44,12 +44,42 @@ # AUDIT FIX (StockLab overhaul, Part A2, docs/AUDIT_PERFORMANCE.md's remaining # company_summary() N+1 finding): same fix as screeners.py -- batches the # Security->Company->{country,sector,industry} chain instead of 3 lazy-load queries per row. .options(*security_eager_load_options()) .join(Company, Security.company_id == Company.id) - .join(Score, Score.security_id == Security.id) - .outerjoin(Valuation, Valuation.security_id == Security.id) + .join( + Score, + and_( + Score.security_id == Security.id, + Score.id + == select(Score.id) + .distinct(Score.security_id) + .order_by( + Score.security_id, + Score.calculation_date.desc(), + ) + .limit(1) + .correlate(Security) + .scalar_subquery(), + ), + ) + .outerjoin( + Valuation, + and_( + Valuation.security_id == Security.id, + Valuation.id + == select(Valuation.id) + .distinct(Valuation.security_id) + .order_by( + Valuation.security_id, + Valuation.calculation_date.desc(), + ) + .limit(1) + .correlate(Security) + .scalar_subquery(), + ), + ) ) if country: query = query.join(Country, Company.country_id == Country.id).where(Country.iso2 == country.upper()) if sector: query = query.where(Company.sector.has(code=sector.upper())) --- a/app/engines/screening/executor.py +++ b/app/engines/screening/executor.py @@ -12,10 +12,42 @@ _SCORE_FIELDS = { "overall_score", "quality_score", "financial_health_score", "growth_score", "competitive_advantage_score", "valuation_score", "risk_score", } + +_PERCENTAGE_METRICS = { + "roic", + "roic_minus_wacc", + "revenue_growth_yoy", + "revenue_cagr_5y", + "eps_growth_yoy", + "fcf_growth_yoy", + "gross_margin", + "operating_margin", + "net_margin", + "fcf_margin", + "fcf_yield", + "dividend_yield", + "buyback_yield", + "shareholder_yield", + "roe", + "fcf_payout_ratio", +} + + +def _normalize_metric_threshold( + metric: str, + value: float | None, +) -> float | None: + """Convert public percentage-point thresholds to stored ratio values.""" + if value is None: + return None + if metric in _PERCENTAGE_METRICS: + return value / 100.0 + return value + _OP_MAP = { "gt": lambda col, v, v2: col > v, "gte": lambda col, v, v2: col >= v, "lt": lambda col, v, v2: col < v, @@ -39,11 +71,13 @@ raise InvalidScreenFilter( f"relative='{f.relative}' filters require percentile columns not yet materialized in " f"this build โ€” see docs/SPEC_COVERAGE.md. Use relative='absolute' for now." ) m = Metric.__table__.alias(f"metric_{f.metric}") - condition = _OP_MAP[f.op](m.c.value, f.value, f.value2) + value = _normalize_metric_threshold(f.metric, f.value) + value2 = _normalize_metric_threshold(f.metric, f.value2) + condition = _OP_MAP[f.op](m.c.value, value, value2) return exists( select(1).select_from(m).where(m.c.security_id == Security.id, m.c.metric_key == f.metric, condition) ) --- a/app/schemas/common.py +++ b/app/schemas/common.py @@ -94,11 +94,11 @@ metrics: list[MetricOut] = [] as_of: Optional[date] = None 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" relative: str = "absolute" # absolute|industry_percentile|historical_percentile|peer_percentile @@ -113,11 +113,11 @@ 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 # upper bound existed here before this pass -- a client (or an unrate-limited abusive caller, # see docs/AUDIT_SECURITY.md finding #3) could request limit=100000 and force --- a/app/workers/celery_app.py +++ b/app/workers/celery_app.py @@ -24,12 +24,18 @@ worker_concurrency=settings.CELERY_WORKER_CONCURRENCY, worker_prefetch_multiplier=settings.CELERY_WORKER_PREFETCH_MULTIPLIER, 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": { "task": "app.workers.ingest.ingest_universe_task", "schedule": crontab(hour=2, minute=0), # spec ยง52: universe refresh minimum daily --- a/app/workers/ingest.py +++ b/app/workers/ingest.py @@ -23,10 +23,11 @@ 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 @@ -478,35 +479,67 @@ 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) + 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) + 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) + 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) + 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) --- a/app/workers/recompute.py +++ b/app/workers/recompute.py @@ -6,11 +6,11 @@ """ from __future__ import annotations from datetime import date -from sqlalchemy import select +from sqlalchemy import or_, select from sqlalchemy.orm import Session from app.core.config import get_settings from app.core.db import SessionLocal from app.core.logging import get_logger @@ -94,10 +94,14 @@ select(FinancialPeriod) .where( FinancialPeriod.security_id == security.id, FinancialPeriod.period_type.in_(("Q1", "Q2", "Q3", "Q4")), FinancialPeriod.period_end <= as_of, + or_( + FinancialPeriod.filing_date.is_(None), + FinancialPeriod.filing_date <= as_of, + ), ) .order_by(FinancialPeriod.period_end.desc()) .limit(8) ).scalars().all() ) @@ -118,11 +122,16 @@ def build_snapshot_from_db(db: Session, security: Security, as_of: date) -> FinancialSnapshot | None: periods = ( db.execute( select(FinancialPeriod) - .where(FinancialPeriod.security_id == security.id, FinancialPeriod.period_type == "FY") + .where( + FinancialPeriod.security_id == security.id, + FinancialPeriod.period_type == "FY", + FinancialPeriod.period_end <= as_of, + FinancialPeriod.filing_date <= as_of, + ) .order_by(FinancialPeriod.period_end.desc()) .limit(11) ).scalars().all() ) if not periods: @@ -337,30 +346,35 @@ empty/partial `source_tiers` list without fabricating a score for missing entries.""" rows = db.execute( select(Source.provider_tier) .select_from(FinancialPeriod) .join(Source, FinancialPeriod.source_id == Source.id) - .where(FinancialPeriod.security_id == security_id, FinancialPeriod.period_end <= as_of) + .where(FinancialPeriod.security_id == security_id, + FinancialPeriod.period_end <= as_of, + or_( + FinancialPeriod.filing_date.is_(None), + FinancialPeriod.filing_date <= as_of, + )) .order_by(FinancialPeriod.period_end.desc()) .limit(limit) ).all() return [tier for (tier,) in rows] def recompute_security(db: Session, security_id: str, peer_metric_values: dict | None = None, - industry_medians: dict | None = None) -> None: + industry_medians: dict | None = None, as_of: date | None = None) -> None: """`peer_metric_values` (industry-percentile universe) is supplied by the caller โ€” computing it requires a cross-security aggregate query, kept out of this function to keep it unit-testable with a hand-built peer set; see app/workers/peer_groups.py (IMPLEMENTED in Part B2 โ€” see SPEC_COVERAGE.md) for the production aggregate-query version.""" settings = get_settings() security = db.get(Security, security_id) if security is None: logger.warning("recompute.security_not_found", security_id=security_id) return - as_of = date.today() + as_of = as_of or date.today() snapshot = build_snapshot_from_db(db, security, as_of) if snapshot is None: logger.warning("recompute.no_financial_data", security_id=security_id) return --- a/tests/test_provider_resilience.py +++ b/tests/test_provider_resilience.py @@ -106,12 +106,17 @@ def test_resolve_line_items_one_source_missing_a_field_other_has_it(): primary = {"revenue": 100.0} # net_income missing from primary entirely secondary = {"revenue": 100.0, "net_income": 10.0} resolved, conflicts = resolve_line_items(primary, secondary, _FIELDS, "PRIMARY", "SECONDARY") - assert resolved["net_income"] == 10.0 # only source with a value wins, not treated as a conflict - assert conflicts == [] + assert resolved["net_income"] == 10.0 + assert len(conflicts) == 1 + assert conflicts[0].field == "net_income" + assert conflicts[0].value_a is None + assert conflicts[0].value_b == 10.0 + assert conflicts[0].selected_value == 10.0 + assert conflicts[0].selection_reason == "Primary field missing; value supplied by secondary source" # --- Part A9: merge_statement_line_items (the critical ingestion bug fix) ---