diff --git a/app/services/rr_scanner_service.py b/app/services/rr_scanner_service.py index a814e43..1884777 100644 --- a/app/services/rr_scanner_service.py +++ b/app/services/rr_scanner_service.py @@ -13,7 +13,7 @@ import logging from collections.abc import Callable from datetime import date, datetime, timedelta, timezone -from sqlalchemy import and_, func, select +from sqlalchemy import and_, func, select, update from sqlalchemy.ext.asyncio import AsyncSession from app.exceptions import NotFoundError @@ -58,6 +58,28 @@ async def _get_ticker(db: AsyncSession, symbol: str) -> Ticker: return ticker +async def _mark_ticker_scores_stale(db: AsyncSession, symbol: str) -> None: + """Prevent a failed refresh from being presented as a current signal.""" + result = await db.execute( + select(Ticker.id).where(Ticker.symbol == symbol.strip().upper()) + ) + ticker_id = result.scalar_one_or_none() + if ticker_id is None: + raise NotFoundError(f"Ticker not found: {symbol.strip().upper()}") + + await db.execute( + update(DimensionScore) + .where(DimensionScore.ticker_id == ticker_id) + .values(is_stale=True) + ) + await db.execute( + update(CompositeScore) + .where(CompositeScore.ticker_id == ticker_id) + .values(is_stale=True) + ) + await db.commit() + + def _compute_quality_score( rr: float, strength: int, @@ -116,8 +138,11 @@ async def _apply_live_recommendation_context( select(DimensionScore).where(DimensionScore.ticker_id.in_(ticker_ids)) ) dims_by_ticker: dict[int, dict[str, float]] = {} + stale_score_ticker_ids: set[int] = set() for ds in dim_result.scalars().all(): dims_by_ticker.setdefault(ds.ticker_id, {})[ds.dimension] = float(ds.score) + if ds.is_stale: + stale_score_ticker_ids.add(ds.ticker_id) comp_result = await db.execute( select(CompositeScore) @@ -153,6 +178,18 @@ async def _apply_live_recommendation_context( live_row["composite_score"] = float(comp.score) live_row["context_as_of"]["score_computed_at"] = comp.computed_at + if ( + comp is None + or comp.is_stale + or ticker_id in stale_score_ticker_ids + ): + live_row["confidence_score"] = None + live_row["recommended_action"] = "NEUTRAL" + live_row["reasoning"] = "Score refresh pending; recommendation withheld." + live_row["risk_level"] = "High" + live_rows.append(live_row) + continue + dimension_scores = dims_by_ticker.get(ticker_id) sentiment = sentiments.get(ticker_id) if sentiment is not None: @@ -591,6 +628,12 @@ async def scan_all_tickers( except Exception: await db.rollback() logger.exception("Error refreshing scores for %s", symbol) + try: + await _mark_ticker_scores_stale(db, symbol) + except Exception: + await db.rollback() + logger.exception("Could not mark scores stale for %s", symbol) + continue try: setups = await scan_ticker( diff --git a/tests/unit/test_rr_scanner_preservation.py b/tests/unit/test_rr_scanner_preservation.py index d7e45f4..58b57b7 100644 --- a/tests/unit/test_rr_scanner_preservation.py +++ b/tests/unit/test_rr_scanner_preservation.py @@ -689,6 +689,31 @@ async def test_live_recommendation_payload_uses_current_score_and_sentiment( assert persisted.reasoning == old_reasoning +@pytest.mark.asyncio +async def test_live_recommendation_withholds_stale_scores( + db_session: AsyncSession, +): + stale_setup = await _seed_stale_setup_with_current_scores(db_session) + composite = ( + await db_session.execute( + select(CompositeScore).where(CompositeScore.ticker_id == stale_setup.ticker_id) + ) + ).scalar_one() + composite.is_stale = True + await db_session.flush() + + rows = await get_trade_setups( + db_session, + symbol="TTWO", + live_recommendation=True, + ) + + assert len(rows) == 1 + assert rows[0]["confidence_score"] is None + assert rows[0]["recommended_action"] == "NEUTRAL" + assert rows[0]["reasoning"] == "Score refresh pending; recommendation withheld." + + @pytest.mark.asyncio async def test_live_recommendation_filters_apply_to_live_values( db_session: AsyncSession, diff --git a/tests/unit/test_rr_scanner_scan_all.py b/tests/unit/test_rr_scanner_scan_all.py index 4f84d5a..89a2997 100644 --- a/tests/unit/test_rr_scanner_scan_all.py +++ b/tests/unit/test_rr_scanner_scan_all.py @@ -3,7 +3,10 @@ from __future__ import annotations import pytest +from datetime import datetime, timezone +from sqlalchemy import select +from app.models.score import CompositeScore, DimensionScore from app.models.ticker import Ticker from app.services import rr_scanner_service, scoring_service from tests.conftest import _test_session_factory # type: ignore @@ -42,6 +45,49 @@ async def test_scan_proceeds_when_score_refresh_fails(session, monkeypatch): assert setups == [] +async def test_scan_marks_scores_stale_when_refresh_fails(session, monkeypatch): + ticker = Ticker(symbol="AAA") + session.add(ticker) + await session.flush() + now = datetime.now(timezone.utc) + session.add_all([ + DimensionScore( + ticker_id=ticker.id, + dimension="technical", + score=70.0, + is_stale=False, + computed_at=now, + ), + CompositeScore( + ticker_id=ticker.id, + score=70.0, + is_stale=False, + weights_json="{}", + computed_at=now, + ), + ]) + await session.commit() + + async def _boom(db, symbol): + raise RuntimeError("scoring unavailable") + + async def _fake_scan_ticker(db, symbol, *args, **kwargs): + comp = ( + await db.execute(select(CompositeScore).where(CompositeScore.ticker_id == ticker.id)) + ).scalar_one() + dimensions = ( + await db.execute(select(DimensionScore).where(DimensionScore.ticker_id == ticker.id)) + ).scalars().all() + assert comp.is_stale is True + assert all(score.is_stale for score in dimensions) + return [] + + monkeypatch.setattr(scoring_service, "compute_all_dimensions", _boom) + monkeypatch.setattr(rr_scanner_service, "scan_ticker", _fake_scan_ticker) + + await rr_scanner_service.scan_all_tickers(session) + + async def test_scan_error_does_not_stop_later_tickers(session, monkeypatch): session.add_all([Ticker(symbol="AAA"), Ticker(symbol="BBB")]) await session.commit()