From c059f838bbbc45c440bde45e498a94c1f42a365f Mon Sep 17 00:00:00 2001 From: Dorian Date: Fri, 18 Sep 2026 16:58:03 +0200 Subject: [PATCH] fix(backend): fiabilise le tri des lectures/predictions et la detection de redemarrage a zero --- apps/backend/app/repositories/prediction.py | 8 ++- apps/backend/app/repositories/reading.py | 5 +- apps/backend/app/services/alert.py | 60 +++++++++++----- .../tests/repositories/test_prediction.py | 24 +++++++ .../tests/repositories/test_reading.py | 23 +++++++ apps/backend/tests/services/test_alert.py | 69 ++++++++++++++++++- apps/backend/tests/test_internal_alerts.py | 36 ++++++++-- docs/architecture/20-backend.md | 13 +++- 8 files changed, 208 insertions(+), 30 deletions(-) diff --git a/apps/backend/app/repositories/prediction.py b/apps/backend/app/repositories/prediction.py index dd3dc28..5939899 100644 --- a/apps/backend/app/repositories/prediction.py +++ b/apps/backend/app/repositories/prediction.py @@ -16,10 +16,16 @@ class PredictionRepository: ) -> Sequence[Prediction]: # Restreint à `available` : une prévision `insufficient_data`/`error` n'a pas de # `predicted_value` à comparer à une lecture réelle (détection d'anomalie). + # Piège : `prediction` n'a pas d'unicité sur `(site_id, target_at)` (cf. + # `enervision_ml.score`, qui insère toujours une nouvelle ligne plutôt que d'écraser la + # précédente). `prediction_id` en dernier départage donc les égalités de `target_at` par + # ordre croissant : `_detect_anomaly` construit un dict qui garde le dernier rencontré, + # c'est-à-dire le run le plus récent plutôt qu'une ligne choisie au hasard par le plan + # d'exécution. requete = ( select(Prediction) .where(Prediction.target_at >= since, Prediction.status == "available") - .order_by(Prediction.site_id, Prediction.target_at) + .order_by(Prediction.site_id, Prediction.target_at, Prediction.prediction_id) ) if site_id is not None: requete = requete.where(Prediction.site_id == site_id) diff --git a/apps/backend/app/repositories/reading.py b/apps/backend/app/repositories/reading.py index 3e9cec5..82a8565 100644 --- a/apps/backend/app/repositories/reading.py +++ b/apps/backend/app/repositories/reading.py @@ -37,10 +37,13 @@ class ReadingRepository: async def list_since(self, *, since: datetime, site_id: str | None = None) -> Sequence[Reading]: # Trié par site puis par heure croissante : la détection d'alertes (spike) a besoin de # comparer chaque lecture à celle qui la précède immédiatement pour le même site. + # `reading_id` en dernier départage : `uq_reading_source` autorise deux lignes au même + # `site_id`+`timestamp` quand la `source` diffère (même piège que `latest_for_site`), sans + # quoi l'ordre entre elles ne serait pas garanti d'un appel à l'autre. requete = ( select(Reading) .where(Reading.timestamp >= since) - .order_by(Reading.site_id, Reading.timestamp) + .order_by(Reading.site_id, Reading.timestamp, Reading.reading_id) ) if site_id is not None: requete = requete.where(Reading.site_id == site_id) diff --git a/apps/backend/app/services/alert.py b/apps/backend/app/services/alert.py index acae8cf..44a1db7 100644 --- a/apps/backend/app/services/alert.py +++ b/apps/backend/app/services/alert.py @@ -142,43 +142,65 @@ def _detect_threshold(lectures: Sequence[Reading], sites_par_id: dict[str, Site] def _detect_spike(lectures: Sequence[Reading]) -> list[Alert]: - # `lectures` est triée par site puis par heure (cf. `ReadingRepository.list_since`) : deux - # lignes consécutives du même site sont donc deux mesures consécutives dans le temps. + # `lectures` est triée par site, heure puis `reading_id` (cf. `ReadingRepository.list_since`) : + # deux lignes consécutives du même site sont donc deux mesures consécutives dans le temps, + # sauf lorsqu'elles partagent le même horodatage (deux `source` différentes pour le même + # instant, permises par `uq_reading_source`) : ce n'est alors pas une variation réelle, on + # l'ignore plutôt que de générer une fausse alerte figée par son `source_alert_id`. alertes = [] precedente: Reading | None = None for lecture in lectures: - if precedente is None or precedente.site_id != lecture.site_id: + if ( + precedente is None + or precedente.site_id != lecture.site_id + or precedente.timestamp == lecture.timestamp + ): precedente = lecture continue avant, apres = precedente.consumption_kw, lecture.consumption_kw precedente = lecture - if avant is None or apres is None or avant == 0: + if avant is None or apres is None: + continue + if avant == 0: + # Une variation relative n'a pas de sens depuis zéro, mais un redémarrage direct à + # une consommation positive reste le signal le plus alarmant du lot : `critical` + # plutôt qu'un ratio indéfini. + if apres > 0: + alertes.append(_spike_alert(lecture, avant, apres, severity="critical")) continue variation = abs(apres - avant) / abs(avant) if variation < SPIKE_RELATIVE_THRESHOLD: continue alertes.append( - Alert( - source_alert_id=f"spike:{THRESHOLD_METRIC}:{lecture.timestamp.isoformat()}", - site_id=lecture.site_id, - source="enervision", - timestamp=lecture.timestamp, - type="spike", + _spike_alert( + lecture, + avant, + apres, severity=_severity_from_ratio(variation / SPIKE_RELATIVE_THRESHOLD), - message=( - f"Variation brutale de {variation * 100:.0f}% entre deux lectures " - f"consécutives ({avant:.1f} kW -> {apres:.1f} kW)" - ), - value=apres, - threshold=avant, - metric=THRESHOLD_METRIC, - prediction_id=None, - raw_data={}, ) ) return alertes +def _spike_alert(lecture: Reading, avant: float, apres: float, *, severity: str) -> Alert: + return Alert( + source_alert_id=f"spike:{THRESHOLD_METRIC}:{lecture.timestamp.isoformat()}", + site_id=lecture.site_id, + source="enervision", + timestamp=lecture.timestamp, + type="spike", + severity=severity, + message=( + f"Variation brutale entre deux lectures consécutives ({avant:.1f} kW -> {apres:.1f} kW)" + ), + value=apres, + threshold=avant, + metric=THRESHOLD_METRIC, + prediction_id=None, + raw_data={}, + ) + + def _detect_anomaly(lectures: Sequence[Reading], predictions: Sequence[Prediction]) -> list[Alert]: # Alignement strict (site_id, target_at == timestamp) : `enervision_ml.score` produit une # cible à l'heure pile suivant la dernière lecture, sur la même grille horaire que `reading`. diff --git a/apps/backend/tests/repositories/test_prediction.py b/apps/backend/tests/repositories/test_prediction.py index 708a469..da71aea 100644 --- a/apps/backend/tests/repositories/test_prediction.py +++ b/apps/backend/tests/repositories/test_prediction.py @@ -68,6 +68,30 @@ async def test_list_since_excludes_predictions_that_are_not_available( assert list(resultats) == [] +async def test_list_since_breaks_a_target_at_tie_by_ascending_prediction_id( + session: AsyncSession, +) -> None: + # `prediction` n'a pas d'unicité sur `(site_id, target_at)` : deux runs de scoring sans + # nouvelle lecture entre-temps produisent deux lignes `available` à la même cible. Sans ce + # départage, `_detect_anomaly` retiendrait une ligne au hasard plutôt que le run le plus + # récent. + site = await creer_site(session) + depot = PredictionRepository(session) + cible = datetime(2026, 9, 16, tzinfo=UTC) + premier_run = await creer_prediction( + session, site_id=site.site_id, target_at=cible, predicted_value=10.0 + ) + second_run = await creer_prediction( + session, site_id=site.site_id, target_at=cible, predicted_value=20.0 + ) + + resultats = await depot.list_since(since=datetime(2026, 9, 1, tzinfo=UTC), site_id=site.site_id) + identifiants = [p.prediction_id for p in resultats] + await session.rollback() + + assert identifiants == [premier_run.prediction_id, second_run.prediction_id] + + async def test_list_since_filters_by_site_id(session: AsyncSession) -> None: premier = await creer_site(session) second = await creer_site(session) diff --git a/apps/backend/tests/repositories/test_reading.py b/apps/backend/tests/repositories/test_reading.py index aaff856..ac3f854 100644 --- a/apps/backend/tests/repositories/test_reading.py +++ b/apps/backend/tests/repositories/test_reading.py @@ -189,6 +189,29 @@ async def test_list_since_excludes_readings_before_the_cutoff(session: AsyncSess assert identifiants == [dedans.reading_id] +async def test_list_since_breaks_a_timestamp_tie_by_ascending_reading_id( + session: AsyncSession, +) -> None: + # `uq_reading_source` autorise deux lignes au même `site_id`+`timestamp` quand la `source` + # diffère (même piège que `latest_for_site`). Sans ce départage, `_detect_spike` traiterait + # cette paire comme une variation réelle selon un ordre non garanti par le plan d'exécution. + site = await creer_site(session) + depot = ReadingRepository(session) + horodatage = datetime(2026, 9, 16, tzinfo=UTC) + premiere = await creer_lecture( + session, site_id=site.site_id, timestamp=horodatage, source="api_history", consumption_kw=10 + ) + seconde = await creer_lecture( + session, site_id=site.site_id, timestamp=horodatage, source="api_current", consumption_kw=42 + ) + + resultats = await depot.list_since(since=datetime(2026, 9, 1, tzinfo=UTC), site_id=site.site_id) + identifiants = [r.reading_id for r in resultats] + await session.rollback() + + assert identifiants == [premiere.reading_id, seconde.reading_id] + + async def test_list_since_filters_by_site_id(session: AsyncSession) -> None: premier = await creer_site(session) second = await creer_site(session) diff --git a/apps/backend/tests/services/test_alert.py b/apps/backend/tests/services/test_alert.py index c99fd6f..97b2b0a 100644 --- a/apps/backend/tests/services/test_alert.py +++ b/apps/backend/tests/services/test_alert.py @@ -262,6 +262,28 @@ async def test_detect_ignores_a_prediction_whose_target_at_does_not_match_the_re assert [a for a in depot.crees if a.type == "anomaly"] == [] +async def test_detect_keeps_the_most_recent_run_when_two_predictions_share_the_same_target() -> ( + None +): + # `PredictionRepository.list_since` départage les égalités de `target_at` par `prediction_id` + # croissant : le repository fait donc déjà passer le run le plus récent en dernier dans la + # liste, et c'est ce dernier que le dict de `_detect_anomaly` doit retenir. + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, consumption_kwh=100.0)], + predictions=[ + FauxPrediction("A", target_at=NOW, predicted_value=100.0, prediction_id=1), + FauxPrediction("A", target_at=NOW, predicted_value=70.0, prediction_id=2), + ], + ) + + await svc.detect(now=NOW) + + (candidate,) = [a for a in depot.crees if a.type == "anomaly"] + assert candidate.threshold == 70.0 + assert candidate.prediction_id == 2 + + async def test_detect_raises_an_outage_alert_past_the_threshold() -> None: derniere = NOW - OUTAGE_THRESHOLD - timedelta(minutes=1) svc, depot = service( @@ -345,7 +367,35 @@ async def test_detect_returns_early_when_there_is_no_site() -> None: assert depot.crees == [] -async def test_detect_ignores_a_spike_when_the_previous_reading_is_zero() -> None: +async def test_detect_ignores_a_spike_pair_with_a_missing_measurement() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[ + FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=None), + FauxLecture("A", NOW, consumption_kw=160.0), + ], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "spike"] == [] + + +async def test_detect_ignores_a_reading_still_at_zero_after_a_previous_zero() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[ + FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=0.0), + FauxLecture("A", NOW, consumption_kw=0.0), + ], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "spike"] == [] + + +async def test_detect_raises_a_critical_spike_when_a_site_restarts_from_zero() -> None: svc, depot = service( sites=[FauxSite("A")], lectures=[ @@ -356,6 +406,23 @@ async def test_detect_ignores_a_spike_when_the_previous_reading_is_zero() -> Non await svc.detect(now=NOW) + (candidate,) = [a for a in depot.crees if a.type == "spike"] + assert candidate.severity == "critical" + assert candidate.value == 50.0 + assert candidate.threshold == 0.0 + + +async def test_detect_ignores_a_spike_pair_sharing_the_same_timestamp() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[ + FauxLecture("A", NOW, consumption_kw=100.0), + FauxLecture("A", NOW, consumption_kw=160.0), + ], + ) + + await svc.detect(now=NOW) + assert [a for a in depot.crees if a.type == "spike"] == [] diff --git a/apps/backend/tests/test_internal_alerts.py b/apps/backend/tests/test_internal_alerts.py index 4740150..690ac00 100644 --- a/apps/backend/tests/test_internal_alerts.py +++ b/apps/backend/tests/test_internal_alerts.py @@ -1,8 +1,10 @@ from datetime import UTC, datetime import pytest +from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession +from app.db.session import get_session_factory from app.detection import internal_alerts from app.repositories.alert import AlertRepository from tests.repositories.test_reading import creer_lecture @@ -50,16 +52,36 @@ def test_main_prints_how_many_alerts_were_recorded( @pytest.mark.integration async def test_run_detection_writes_a_threshold_alert_end_to_end(session: AsyncSession) -> None: + # `run_detection` ouvre sa propre session et commite : `session.rollback()` seul ne défait + # rien ici (contrairement au reste de la suite), d'où le nettoyage explicite ci-dessous, sur + # le modèle de `tests/api/test_matrice_acces.py`. site = await creer_site(session, capacity_kw=100.0) + site_id = site.site_id instant = datetime(2026, 9, 16, 12, tzinfo=UTC) - await creer_lecture(session, site_id=site.site_id, timestamp=instant, consumption_kw=150.0) + await creer_lecture(session, site_id=site_id, timestamp=instant, consumption_kw=150.0) await session.commit() - nombre = await internal_alerts.run_detection(now=instant, site_id=site.site_id) + try: + nombre = await internal_alerts.run_detection(now=instant, site_id=site_id) - alertes = await AlertRepository(session).list_all(site_id=site.site_id) - types = [a.type for a in alertes] - await session.rollback() + alertes = await AlertRepository(session).list_all(site_id=site_id) + types = [a.type for a in alertes] + await session.rollback() - assert nombre == 1 - assert types == ["threshold"] + assert nombre == 1 + assert types == ["threshold"] + finally: + # `site.site_id` n'est plus sûr après `session.rollback()` : le rollback expire tous les + # objets de la session (indépendamment d'`expire_on_commit`), et y accéder ici relance une + # requête hors contexte async. D'où `site_id`, capturé avant. + async with get_session_factory()() as nettoyage: + await nettoyage.execute( + text("delete from alert where site_id = :site_id"), {"site_id": site_id} + ) + await nettoyage.execute( + text("delete from reading where site_id = :site_id"), {"site_id": site_id} + ) + await nettoyage.execute( + text("delete from site where site_id = :site_id"), {"site_id": site_id} + ) + await nettoyage.commit() diff --git a/docs/architecture/20-backend.md b/docs/architecture/20-backend.md index 8c3641d..731ff21 100644 --- a/docs/architecture/20-backend.md +++ b/docs/architecture/20-backend.md @@ -217,7 +217,7 @@ leur sens dans chaque règle plutôt que d'être laissés à `null` par commodit | `type` | Règle | `value` / `threshold` | |---|---|---| | `threshold` | `reading.consumption_kw` dépasse `site.capacity_kw` (site sans capacité déclarée : ignoré) | mesure / capacité du site | -| `spike` | Variation relative ≥ 50% (`SPIKE_RELATIVE_THRESHOLD`) entre deux lectures consécutives du même site | mesure actuelle / mesure précédente | +| `spike` | Variation relative ≥ 50% (`SPIKE_RELATIVE_THRESHOLD`) entre deux lectures consécutives du même site, ou redémarrage direct à une valeur positive depuis zéro (`critical`) | mesure actuelle / mesure précédente | | `anomaly` | Écart relatif ≥ 30% (`ANOMALY_RELATIVE_THRESHOLD`) entre `reading.consumption_kwh` et la `prediction` du même site dont `target_at == timestamp` | mesure réelle / valeur prédite | | `outage` | Aucune lecture depuis plus de 3h (`OUTAGE_THRESHOLD`, 3x la cadence horaire nominale), ou site jamais lu | `null` / `null` | | `sensor` | `reading.data_quality` ∈ `partial`/`degraded`/`critical` | `null` / `null` | @@ -229,6 +229,17 @@ au-delà. `AlertRepository.create_many()` insère par lot avec `ON CONFLICT DO N horodatage) : rejouer la détection sur une fenêtre déjà analysée ne duplique donc jamais une alerte. +**Pièges de tri corrigés en revue** : `reading`/`prediction` n'ont pas d'unicité sur leur couple +métier (`uq_reading_source` autorise deux `source` différentes au même `site_id`+`timestamp`, +`prediction` n'a aucune contrainte sur `(site_id, target_at)`, chaque run de scoring gardant sa +propre ligne). `ReadingRepository.list_since()`/`PredictionRepository.list_since()` départagent +donc les égalités par `reading_id`/`prediction_id` croissant, comme le font déjà +`latest_by_site()`/`latest_for_site()` sur les mêmes tables ; sans ce départage, l'ordre entre +lignes à égalité n'est pas garanti d'un appel à l'autre, et `_detect_spike`/`_detect_anomaly` +auraient pu comparer des lectures/choisir une prévision au hasard. `_detect_spike` ignore en plus +explicitement les paires de lectures qui partagent le même horodatage (deux `source` pour un seul +instant réel, pas une variation). + Comme `enervision_ml.score`, la détection est un script lancé à la main, pas encore ordonnancé par Airflow : `uv run python -m app.detection.internal_alerts [--site-id ...] [--now ...]`, dans `apps/backend` puisque les règles s'appuient sur les repositories ORM de l'API plutôt que sur une