--- 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/api/v1/rankings.py +++ b/app/api/v1/rankings.py @@ -11,7 +11,7 @@ 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 @@ -46,8 +46,38 @@ # 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()) --- 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 --- 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/workers/recompute.py +++ b/app/workers/recompute.py @@ -120,7 +120,12 @@ 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() @@ -347,7 +352,7 @@ 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 @@ -358,7 +363,7 @@ 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) --- a/tests/test_provider_resilience.py +++ b/tests/test_provider_resilience.py @@ -108,8 +108,13 @@ 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) --- --- /dev/null +++ b/.pytest_cache/CACHEDIR.TAG @@ -0,0 +1,4 @@ +Signature: 8a477f597d28d172789f06886806bc55 +# This file is a cache directory tag created by pytest. +# For information about cache directory tags, see: +# https://bford.info/cachedir/spec.html --- /dev/null +++ b/.pytest_cache/.gitignore @@ -0,0 +1,2 @@ +# Created by pytest automatically. +* --- /dev/null +++ b/.pytest_cache/README.md @@ -0,0 +1,8 @@ +# pytest cache directory # + +This directory contains data from the pytest's cache plugin, +which provides the `--lf` and `--ff` options, as well as the `cache` fixture. + +**Do not** commit this to version control. + +See [the docs](https://docs.pytest.org/en/stable/how-to/cache.html) for more information. --- /dev/null +++ b/.pytest_cache/v/cache/nodeids @@ -0,0 +1,51 @@ +[ + "tests/test_provider_resilience.py::test_merge_combines_all_three_statements", + "tests/test_provider_resilience.py::test_merge_drops_none_valued_keys_entirely", + "tests/test_provider_resilience.py::test_merge_fills_a_field_the_income_statement_left_none", + "tests/test_provider_resilience.py::test_merge_income_statement_wins_on_a_shared_field", + "tests/test_provider_resilience.py::test_merge_later_none_never_erases_an_earlier_real_value", + "tests/test_provider_resilience.py::test_merge_of_all_empty_inputs_is_empty", + "tests/test_provider_resilience.py::test_merge_of_three_identical_dicts_is_idempotent", + "tests/test_provider_resilience.py::test_merge_with_only_an_income_statement_is_identity", + "tests/test_provider_resilience.py::test_resolve_conflict_flags_real_disagreement", + "tests/test_provider_resilience.py::test_resolve_conflict_handles_zero_value_without_dividing_by_zero", + "tests/test_provider_resilience.py::test_resolve_conflict_no_sources_available", + "tests/test_provider_resilience.py::test_resolve_conflict_picks_highest_priority_source", + "tests/test_provider_resilience.py::test_resolve_conflict_tolerates_small_rounding_differences", + "tests/test_provider_resilience.py::test_resolve_conflict_unknown_tier_sorts_last", + "tests/test_provider_resilience.py::test_resolve_line_items_agreeing_sources_no_conflicts", + "tests/test_provider_resilience.py::test_resolve_line_items_field_absent_from_both_sources_is_omitted", + "tests/test_provider_resilience.py::test_resolve_line_items_flags_real_disagreement_on_one_field_only", + "tests/test_provider_resilience.py::test_resolve_line_items_no_secondary_passes_through_primary_unchanged", + "tests/test_provider_resilience.py::test_resolve_line_items_official_filing_secondary_overrides_primary_tier", + "tests/test_provider_resilience.py::test_resolve_line_items_one_source_missing_a_field_other_has_it", + "tests/test_validation.py::test_a_bad_optional_field_does_not_reject_the_bar", + "tests/test_validation.py::test_a_warning_only_report_is_still_ok", + "tests/test_validation.py::test_an_empty_report_has_a_zero_rejection_rate_not_a_division_error", + "tests/test_validation.py::test_bar_with_a_non_positive_close_is_rejected", + "tests/test_validation.py::test_bar_without_a_close_is_rejected", + "tests/test_validation.py::test_bar_without_a_date_is_rejected", + "tests/test_validation.py::test_close_above_high_is_warned_but_kept", + "tests/test_validation.py::test_close_below_low_is_warned_but_kept", + "tests/test_validation.py::test_coerce_accepts_a_numeric_string", + "tests/test_validation.py::test_coerce_accepts_int_and_float", + "tests/test_validation.py::test_coerce_records_an_issue_for_an_absent_required_field", + "tests/test_validation.py::test_coerce_rejects_a_boolean", + "tests/test_validation.py::test_coerce_rejects_a_list", + "tests/test_validation.py::test_coerce_rejects_a_non_numeric_string", + "tests/test_validation.py::test_coerce_rejects_nan_and_infinity", + "tests/test_validation.py::test_coerce_treats_none_and_empty_string_as_absent_without_an_issue", + "tests/test_validation.py::test_coerce_without_a_report_still_returns_the_right_value", + "tests/test_validation.py::test_consistent_line_items_produce_no_issues", + "tests/test_validation.py::test_current_assets_above_total_assets_is_warned", + "tests/test_validation.py::test_diluted_below_outstanding_is_warned", + "tests/test_validation.py::test_gross_profit_above_revenue_is_warned", + "tests/test_validation.py::test_high_below_low_drops_both_and_warns", + "tests/test_validation.py::test_issues_carry_the_row_key_so_they_can_be_traced", + "tests/test_validation.py::test_legitimately_negative_fields_are_kept", + "tests/test_validation.py::test_line_items_coerces_numeric_strings", + "tests/test_validation.py::test_line_items_drops_a_bad_field_without_losing_the_others", + "tests/test_validation.py::test_negative_revenue_is_dropped_as_impossible", + "tests/test_validation.py::test_report_counts_and_rejection_rate", + "tests/test_validation.py::test_valid_bar_is_accepted_with_every_field" +] \ No newline at end of file --- /dev/null +++ b/.pytest_cache/v/cache/stepwise @@ -0,0 +1 @@ +[] \ No newline at end of file