Withhold stale-score trade recommendations

This commit is contained in:
2026-07-11 16:56:52 +02:00
parent 292b9934b1
commit 727b147c81
3 changed files with 115 additions and 1 deletions
+44 -1
View File
@@ -13,7 +13,7 @@ import logging
from collections.abc import Callable from collections.abc import Callable
from datetime import date, datetime, timedelta, timezone 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 sqlalchemy.ext.asyncio import AsyncSession
from app.exceptions import NotFoundError from app.exceptions import NotFoundError
@@ -58,6 +58,28 @@ async def _get_ticker(db: AsyncSession, symbol: str) -> Ticker:
return 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( def _compute_quality_score(
rr: float, rr: float,
strength: int, strength: int,
@@ -116,8 +138,11 @@ async def _apply_live_recommendation_context(
select(DimensionScore).where(DimensionScore.ticker_id.in_(ticker_ids)) select(DimensionScore).where(DimensionScore.ticker_id.in_(ticker_ids))
) )
dims_by_ticker: dict[int, dict[str, float]] = {} dims_by_ticker: dict[int, dict[str, float]] = {}
stale_score_ticker_ids: set[int] = set()
for ds in dim_result.scalars().all(): for ds in dim_result.scalars().all():
dims_by_ticker.setdefault(ds.ticker_id, {})[ds.dimension] = float(ds.score) 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( comp_result = await db.execute(
select(CompositeScore) select(CompositeScore)
@@ -153,6 +178,18 @@ async def _apply_live_recommendation_context(
live_row["composite_score"] = float(comp.score) live_row["composite_score"] = float(comp.score)
live_row["context_as_of"]["score_computed_at"] = comp.computed_at 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) dimension_scores = dims_by_ticker.get(ticker_id)
sentiment = sentiments.get(ticker_id) sentiment = sentiments.get(ticker_id)
if sentiment is not None: if sentiment is not None:
@@ -591,6 +628,12 @@ async def scan_all_tickers(
except Exception: except Exception:
await db.rollback() await db.rollback()
logger.exception("Error refreshing scores for %s", symbol) 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: try:
setups = await scan_ticker( setups = await scan_ticker(
@@ -689,6 +689,31 @@ async def test_live_recommendation_payload_uses_current_score_and_sentiment(
assert persisted.reasoning == old_reasoning 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 @pytest.mark.asyncio
async def test_live_recommendation_filters_apply_to_live_values( async def test_live_recommendation_filters_apply_to_live_values(
db_session: AsyncSession, db_session: AsyncSession,
+46
View File
@@ -3,7 +3,10 @@
from __future__ import annotations from __future__ import annotations
import pytest 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.models.ticker import Ticker
from app.services import rr_scanner_service, scoring_service from app.services import rr_scanner_service, scoring_service
from tests.conftest import _test_session_factory # type: ignore 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 == [] 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): async def test_scan_error_does_not_stop_later_tickers(session, monkeypatch):
session.add_all([Ticker(symbol="AAA"), Ticker(symbol="BBB")]) session.add_all([Ticker(symbol="AAA"), Ticker(symbol="BBB")])
await session.commit() await session.commit()