From edd5e82d29633cc010e9a2c7dbd785f32aee1ff7 Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Fri, 18 Sep 2026 16:04:10 +0200 Subject: [PATCH 1/3] fix --- sonar-project.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sonar-project.properties b/sonar-project.properties index 43c28fd..49c6abe 100644 --- a/sonar-project.properties +++ b/sonar-project.properties @@ -9,7 +9,7 @@ sonar.tests=apps/frontend/src,apps/backend/tests sonar.test.inclusions=**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py # Liste des fichiers et dossiers à exclure de l'analyse -sonar.exclusions=.pytest_cache,.venv,alembic,tests,**/*/node_modules/**,**/*/dist/**,**/*/build/**,**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py +sonar.exclusions=.pytest_cache,.venv,alembic,tests,**/*/node_modules/**,**/*/dist/**,**/*/build/**,**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py,**/*.spec.ts # Chemin vers le rapport de couverture de code # Fichier généré par Pytest From a9e124a97d55445c72a1a8632f05515c19533873 Mon Sep 17 00:00:00 2001 From: Dorian Date: Fri, 18 Sep 2026 16:10:06 +0200 Subject: [PATCH 2/3] feat(backend): detecte les alertes internes a partir des lectures et previsions --- apps/backend/app/api/deps.py | 7 +- apps/backend/app/detection/__init__.py | 0 apps/backend/app/detection/internal_alerts.py | 68 ++++ apps/backend/app/repositories/alert.py | 34 ++ apps/backend/app/repositories/prediction.py | 15 + apps/backend/app/repositories/reading.py | 12 + apps/backend/app/services/alert.py | 291 ++++++++++++++- apps/backend/tests/repositories/test_alert.py | 57 +++ .../tests/repositories/test_prediction.py | 55 +++ .../tests/repositories/test_reading.py | 50 +++ apps/backend/tests/services/test_alert.py | 347 +++++++++++++++++- apps/backend/tests/test_internal_alerts.py | 65 ++++ docs/architecture/20-backend.md | 29 ++ 13 files changed, 1021 insertions(+), 9 deletions(-) create mode 100644 apps/backend/app/detection/__init__.py create mode 100644 apps/backend/app/detection/internal_alerts.py create mode 100644 apps/backend/tests/test_internal_alerts.py diff --git a/apps/backend/app/api/deps.py b/apps/backend/app/api/deps.py index c247f4c..6098403 100644 --- a/apps/backend/app/api/deps.py +++ b/apps/backend/app/api/deps.py @@ -180,7 +180,12 @@ SiteServiceDep = Annotated[SiteService, Depends(get_site_service)] def get_alert_service(session: SessionDep) -> AlertService: - return AlertService(alerts=AlertRepository(session)) + return AlertService( + alerts=AlertRepository(session), + readings=ReadingRepository(session), + predictions=PredictionRepository(session), + sites=SiteRepository(session), + ) AlertServiceDep = Annotated[AlertService, Depends(get_alert_service)] diff --git a/apps/backend/app/detection/__init__.py b/apps/backend/app/detection/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/backend/app/detection/internal_alerts.py b/apps/backend/app/detection/internal_alerts.py new file mode 100644 index 0000000..d7e0785 --- /dev/null +++ b/apps/backend/app/detection/internal_alerts.py @@ -0,0 +1,68 @@ +# Détection d'alertes internes EnerVision (issue #104) : script lancé à la main pour l'instant, +# comme `enervision_ml.score` côté ML, sans automatisation Airflow pour l'ordonnancer. + +from __future__ import annotations + +import argparse +import asyncio +import sys +from datetime import UTC, datetime + +from app.core.config import get_settings +from app.db.session import get_session_factory +from app.repositories.alert import AlertRepository +from app.repositories.prediction import PredictionRepository +from app.repositories.reading import ReadingRepository +from app.repositories.site import SiteRepository +from app.services.alert import AlertService + + +async def run_detection(*, now: datetime | None = None, site_id: str | None = None) -> int: + """Exécute les cinq règles de détection et enregistre les nouvelles alertes. Rend le nombre de + lignes effectivement insérées (les doublons de `source_alert_id` sont silencieusement + ignorés).""" + async with get_session_factory()() as session: + service = AlertService( + alerts=AlertRepository(session), + readings=ReadingRepository(session), + predictions=PredictionRepository(session), + sites=SiteRepository(session), + ) + nouvelles = await service.detect(now=now, site_id=site_id) + await session.commit() + return len(nouvelles) + + +def _parse_instant(valeur: str) -> datetime: + instant = datetime.fromisoformat(valeur) + return instant if instant.tzinfo is not None else instant.replace(tzinfo=UTC) + + +def parse_args(argv: list[str] | None = None) -> argparse.Namespace: + parser = argparse.ArgumentParser( + prog="python -m app.detection.internal_alerts", + description="Détection d'alertes internes EnerVision", + ) + parser.add_argument("--site-id", default=None, help="Limite la détection à un seul site.") + parser.add_argument( + "--now", + type=_parse_instant, + default=None, + help=( + "Instant de référence (ISO 8601, UTC si le fuseau est omis). Défaut : l'heure courante." + ), + ) + return parser.parse_args(argv) + + +def main(argv: list[str] | None = None) -> int: + args = parse_args(argv) + # Échoue tôt si `APP_SECRET_KEY`/`DATABASE_URL` manquent, avant toute requête à la base. + get_settings() + nombre = asyncio.run(run_detection(now=args.now, site_id=args.site_id)) + print(f"{nombre} nouvelle(s) alerte(s) enregistrée(s).") + return 0 + + +if __name__ == "__main__": # pragma: no cover + sys.exit(main()) diff --git a/apps/backend/app/repositories/alert.py b/apps/backend/app/repositories/alert.py index 4b0766f..f495a3b 100644 --- a/apps/backend/app/repositories/alert.py +++ b/apps/backend/app/repositories/alert.py @@ -1,6 +1,7 @@ from collections.abc import Sequence from sqlalchemy import select +from sqlalchemy.dialects.postgresql import insert from sqlalchemy.ext.asyncio import AsyncSession from app.models.energy import Alert @@ -19,3 +20,36 @@ class AlertRepository: if severity is not None: requete = requete.where(Alert.severity == severity) return (await self._session.scalars(requete)).all() + + async def create_many(self, alerts: Sequence[Alert]) -> Sequence[Alert]: + # `ON CONFLICT DO NOTHING` sur `uq_alert_source_reference` : rejouer la détection sur une + # fenêtre qui recouvre une exécution précédente ne doit pas dupliquer une alerte déjà + # enregistrée. `RETURNING` ne renvoie donc que les lignes effectivement insérées. + if not alerts: + return [] + valeurs = [ + { + "source_alert_id": alerte.source_alert_id, + "site_id": alerte.site_id, + "source": alerte.source, + "timestamp": alerte.timestamp, + "type": alerte.type, + "severity": alerte.severity, + "message": alerte.message, + "value": alerte.value, + "threshold": alerte.threshold, + "metric": alerte.metric, + "prediction_id": alerte.prediction_id, + "raw_data": alerte.raw_data, + } + for alerte in alerts + ] + requete = ( + insert(Alert) + .values(valeurs) + .on_conflict_do_nothing(constraint="uq_alert_source_reference") + .returning(Alert) + ) + resultat = await self._session.execute(requete) + await self._session.flush() + return resultat.scalars().all() diff --git a/apps/backend/app/repositories/prediction.py b/apps/backend/app/repositories/prediction.py index f79311a..dd3dc28 100644 --- a/apps/backend/app/repositories/prediction.py +++ b/apps/backend/app/repositories/prediction.py @@ -1,4 +1,5 @@ from collections.abc import Sequence +from datetime import datetime from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession @@ -10,6 +11,20 @@ class PredictionRepository: def __init__(self, session: AsyncSession) -> None: self._session = session + async def list_since( + self, *, since: datetime, site_id: str | None = None + ) -> 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). + requete = ( + select(Prediction) + .where(Prediction.target_at >= since, Prediction.status == "available") + .order_by(Prediction.site_id, Prediction.target_at) + ) + if site_id is not None: + requete = requete.where(Prediction.site_id == site_id) + return (await self._session.scalars(requete)).all() + async def latest_by_site(self) -> Sequence[Prediction]: # `.distinct(site_id)` compile en `DISTINCT ON (site_id)` sous PostgreSQL : une seule # ligne par site, la plus récente grâce à l'ordre composite qui suit. Même mécanisme que diff --git a/apps/backend/app/repositories/reading.py b/apps/backend/app/repositories/reading.py index d005d16..3e9cec5 100644 --- a/apps/backend/app/repositories/reading.py +++ b/apps/backend/app/repositories/reading.py @@ -34,6 +34,18 @@ class ReadingRepository: lecture: Reading | None = await self._session.scalar(requete) return lecture + 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. + requete = ( + select(Reading) + .where(Reading.timestamp >= since) + .order_by(Reading.site_id, Reading.timestamp) + ) + if site_id is not None: + requete = requete.where(Reading.site_id == site_id) + return (await self._session.scalars(requete)).all() + async def list_history( self, *, diff --git a/apps/backend/app/services/alert.py b/apps/backend/app/services/alert.py index a3ad16e..acae8cf 100644 --- a/apps/backend/app/services/alert.py +++ b/apps/backend/app/services/alert.py @@ -1,14 +1,301 @@ from collections.abc import Sequence +from datetime import UTC, datetime, timedelta -from app.models.energy import Alert +from app.models.energy import Alert, Prediction, Reading, Site from app.repositories.alert import AlertRepository +from app.repositories.prediction import PredictionRepository +from app.repositories.reading import ReadingRepository +from app.repositories.site import SiteRepository + +# Fenêtre de lectures/prédictions analysée à chaque exécution : assez large pour couvrir une paire +# de lectures consécutives (spike) et une coupure prolongée (outage), sans réanalyser tout +# l'historique à chaque lancement manuel du script de détection. +LOOKBACK = timedelta(hours=48) + +# Cadence nominale d'une lecture : le CSV historique comme l'API Mock livrent un pas horaire. +EXPECTED_INTERVAL = timedelta(hours=1) +# Au-delà de trois pas manqués, on parle de coupure plutôt que d'un simple retard d'ingestion. +OUTAGE_THRESHOLD = EXPECTED_INTERVAL * 3 + +# +/-50% entre deux lectures consécutives du même site. +SPIKE_RELATIVE_THRESHOLD = 0.5 +# 30% d'écart entre la consommation réelle et la prévision du même site/instant. +ANOMALY_RELATIVE_THRESHOLD = 0.3 +# Une prévision quasi nulle rend l'écart relatif ininterprétable ; on l'ignore plutôt. +ANOMALY_MINIMUM_PREDICTED_VALUE = 1e-6 + +THRESHOLD_METRIC = "consumption_kw" +ANOMALY_METRIC = "consumption_kwh" +# `data_quality` -> sévérité du capteur défaillant. `good` est volontairement absent : il ne +# déclenche jamais d'alerte. +QUALITE_VERS_SEVERITE: dict[str, str] = { + "partial": "low", + "degraded": "medium", + "critical": "critical", +} class AlertService: - def __init__(self, *, alerts: AlertRepository) -> None: + def __init__( + self, + *, + alerts: AlertRepository, + readings: ReadingRepository, + predictions: PredictionRepository, + sites: SiteRepository, + ) -> None: self._alerts = alerts + self._readings = readings + self._predictions = predictions + self._sites = sites async def list_all( self, *, site_id: str | None = None, severity: str | None = None ) -> Sequence[Alert]: return await self._alerts.list_all(site_id=site_id, severity=severity) + + async def detect( + self, *, now: datetime | None = None, site_id: str | None = None + ) -> Sequence[Alert]: + """Compare les lectures/prévisions récentes aux cinq règles internes et enregistre les + alertes déclenchées (`source='enervision'`). Idempotent grâce à `source_alert_id` : + rejouer sur une fenêtre déjà analysée ne recrée pas les mêmes lignes.""" + instant = now or datetime.now(UTC) + depuis = instant - LOOKBACK + + sites = await self._sites.list_all() + if site_id is not None: + sites = [site for site in sites if site.site_id == site_id] + sites_par_id = {site.site_id: site for site in sites} + if not sites_par_id: + return [] + + lectures = [ + lecture + for lecture in await self._readings.list_since(since=depuis, site_id=site_id) + if lecture.site_id in sites_par_id + ] + predictions = [ + prediction + for prediction in await self._predictions.list_since(since=depuis, site_id=site_id) + if prediction.site_id in sites_par_id + ] + dernieres_lectures = { + lecture.site_id: lecture + for lecture in await self._readings.latest_by_site() + if lecture.site_id in sites_par_id + } + + candidates = [ + *_detect_threshold(lectures, sites_par_id), + *_detect_spike(lectures), + *_detect_anomaly(lectures, predictions), + *_detect_outage(sites, dernieres_lectures, instant), + *_detect_sensor(lectures), + ] + if not candidates: + return [] + return await self._alerts.create_many(candidates) + + +def _severity_from_ratio(ratio: float) -> str: + if ratio >= 2.0: + return "critical" + if ratio >= 1.5: + return "high" + if ratio >= 1.2: + return "medium" + return "low" + + +def _detect_threshold(lectures: Sequence[Reading], sites_par_id: dict[str, Site]) -> list[Alert]: + # Seuil fixe = la capacité déclarée du site : dépasser `capacity_kw` est un dépassement + # matériel, pas une simple variation, et évite un seuil arbitraire non fourni par le domaine. + alertes = [] + for lecture in lectures: + site = sites_par_id[lecture.site_id] + valeur = lecture.consumption_kw + if site.capacity_kw is None or site.capacity_kw <= 0 or valeur is None: + continue + if valeur <= site.capacity_kw: + continue + alertes.append( + Alert( + source_alert_id=f"threshold:{THRESHOLD_METRIC}:{lecture.timestamp.isoformat()}", + site_id=lecture.site_id, + source="enervision", + timestamp=lecture.timestamp, + type="threshold", + severity=_severity_from_ratio(valeur / site.capacity_kw), + message=( + f"Puissance appelée {valeur:.1f} kW au-dessus de la capacité du site " + f"({site.capacity_kw:.1f} kW)" + ), + value=valeur, + threshold=site.capacity_kw, + metric=THRESHOLD_METRIC, + prediction_id=None, + raw_data={}, + ) + ) + return alertes + + +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. + alertes = [] + precedente: Reading | None = None + for lecture in lectures: + if precedente is None or precedente.site_id != lecture.site_id: + precedente = lecture + continue + avant, apres = precedente.consumption_kw, lecture.consumption_kw + precedente = lecture + if avant is None or apres is None or avant == 0: + 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", + 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 _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`. + predictions_par_cle = { + (prediction.site_id, prediction.target_at): prediction + for prediction in predictions + if prediction.target_metric == ANOMALY_METRIC + } + alertes = [] + for lecture in lectures: + prediction = predictions_par_cle.get((lecture.site_id, lecture.timestamp)) + reel = lecture.consumption_kwh + if prediction is None or reel is None or prediction.predicted_value is None: + continue + predite = prediction.predicted_value + if abs(predite) < ANOMALY_MINIMUM_PREDICTED_VALUE: + continue + ecart = abs(reel - predite) / abs(predite) + if ecart < ANOMALY_RELATIVE_THRESHOLD: + continue + alertes.append( + Alert( + source_alert_id=f"anomaly:{ANOMALY_METRIC}:{lecture.timestamp.isoformat()}", + site_id=lecture.site_id, + source="enervision", + timestamp=lecture.timestamp, + type="anomaly", + severity=_severity_from_ratio(ecart / ANOMALY_RELATIVE_THRESHOLD), + message=( + f"Écart de {ecart * 100:.0f}% entre la consommation mesurée ({reel:.1f} kWh) " + f"et la prévision ({predite:.1f} kWh)" + ), + value=reel, + threshold=predite, + metric=ANOMALY_METRIC, + prediction_id=prediction.prediction_id, + raw_data={}, + ) + ) + return alertes + + +def _detect_outage( + sites: Sequence[Site], dernieres_lectures: dict[str, Reading], now: datetime +) -> list[Alert]: + alertes = [] + for site in sites: + derniere = dernieres_lectures.get(site.site_id) + if derniere is None: + alertes.append( + _outage_alert( + site.site_id, + now, + reference=None, + message="Aucune lecture n'a jamais été reçue pour ce site", + severity="critical", + ) + ) + continue + absence = now - derniere.timestamp + if absence < OUTAGE_THRESHOLD: + continue + alertes.append( + _outage_alert( + site.site_id, + now, + reference=derniere.timestamp, + message=( + f"Aucune lecture depuis {absence} (dernière lecture : " + f"{derniere.timestamp.isoformat()})" + ), + severity=_severity_from_ratio(absence / OUTAGE_THRESHOLD), + ) + ) + return alertes + + +def _outage_alert( + site_id: str, now: datetime, *, reference: datetime | None, message: str, severity: str +) -> Alert: + return Alert( + source_alert_id=f"outage:{reference.isoformat() if reference is not None else 'jamais'}", + site_id=site_id, + source="enervision", + timestamp=now, + type="outage", + severity=severity, + message=message, + value=None, + threshold=None, + metric=None, + prediction_id=None, + raw_data={}, + ) + + +def _detect_sensor(lectures: Sequence[Reading]) -> list[Alert]: + alertes = [] + for lecture in lectures: + severite = QUALITE_VERS_SEVERITE.get(lecture.data_quality or "") + if severite is None: + continue + raisons = ", ".join(lecture.null_reasons or []) or "raison non précisée" + alertes.append( + Alert( + source_alert_id=f"sensor:{lecture.timestamp.isoformat()}", + site_id=lecture.site_id, + source="enervision", + timestamp=lecture.timestamp, + type="sensor", + severity=severite, + message=f"Qualité de mesure {lecture.data_quality} ({raisons})", + value=None, + threshold=None, + metric=None, + prediction_id=None, + raw_data={}, + ) + ) + return alertes diff --git a/apps/backend/tests/repositories/test_alert.py b/apps/backend/tests/repositories/test_alert.py index d2a78d0..16c5a9a 100644 --- a/apps/backend/tests/repositories/test_alert.py +++ b/apps/backend/tests/repositories/test_alert.py @@ -89,3 +89,60 @@ async def test_list_all_returns_an_empty_list_when_there_is_nothing( alertes = await depot.list_all(site_id=identifiant_site()) assert list(alertes) == [] + + +def _alerte_a_inserer(*, site_id: str, source_alert_id: str) -> Alert: + return Alert( + source_alert_id=source_alert_id, + site_id=site_id, + source="enervision", + timestamp=datetime(2026, 9, 16, tzinfo=UTC), + type="threshold", + severity="high", + message="Dépassement du seuil configuré", + value=812.5, + threshold=720.0, + metric="consumption_kw", + prediction_id=None, + raw_data={}, + ) + + +async def test_create_many_inserts_every_alert(session: AsyncSession) -> None: + site = await creer_site(session) + depot = AlertRepository(session) + + creees = await depot.create_many( + [ + _alerte_a_inserer(site_id=site.site_id, source_alert_id="threshold:a"), + _alerte_a_inserer(site_id=site.site_id, source_alert_id="threshold:b"), + ] + ) + identifiants = [a.alert_id for a in creees] + await session.rollback() + + assert len(identifiants) == 2 + assert all(identifiant is not None for identifiant in identifiants) + + +async def test_create_many_skips_a_duplicate_source_alert_id(session: AsyncSession) -> None: + site = await creer_site(session) + depot = AlertRepository(session) + await depot.create_many( + [_alerte_a_inserer(site_id=site.site_id, source_alert_id="threshold:rejouee")] + ) + + rejouees = await depot.create_many( + [_alerte_a_inserer(site_id=site.site_id, source_alert_id="threshold:rejouee")] + ) + await session.rollback() + + assert rejouees == [] + + +async def test_create_many_does_nothing_for_an_empty_list(session: AsyncSession) -> None: + depot = AlertRepository(session) + + creees = await depot.create_many([]) + + assert creees == [] diff --git a/apps/backend/tests/repositories/test_prediction.py b/apps/backend/tests/repositories/test_prediction.py index 2be7ddd..708a469 100644 --- a/apps/backend/tests/repositories/test_prediction.py +++ b/apps/backend/tests/repositories/test_prediction.py @@ -29,6 +29,61 @@ async def creer_prediction( return prediction +async def test_list_since_excludes_predictions_before_the_cutoff(session: AsyncSession) -> None: + site = await creer_site(session) + depot = PredictionRepository(session) + dedans = await creer_prediction( + session, site_id=site.site_id, target_at=datetime(2026, 9, 16, tzinfo=UTC) + ) + await creer_prediction( + session, site_id=site.site_id, target_at=datetime(2026, 9, 1, tzinfo=UTC) + ) + + resultats = await depot.list_since( + since=datetime(2026, 9, 10, tzinfo=UTC), site_id=site.site_id + ) + identifiants = [p.prediction_id for p in resultats] + await session.rollback() + + assert identifiants == [dedans.prediction_id] + + +async def test_list_since_excludes_predictions_that_are_not_available( + session: AsyncSession, +) -> None: + site = await creer_site(session) + depot = PredictionRepository(session) + await creer_prediction( + session, + site_id=site.site_id, + target_at=datetime(2026, 9, 16, tzinfo=UTC), + status="insufficient_data", + predicted_value=None, + failure_reason="pas assez d'historique", + ) + + resultats = await depot.list_since(since=datetime(2026, 9, 1, tzinfo=UTC), site_id=site.site_id) + await session.rollback() + + assert list(resultats) == [] + + +async def test_list_since_filters_by_site_id(session: AsyncSession) -> None: + premier = await creer_site(session) + second = await creer_site(session) + depot = PredictionRepository(session) + voulue = await creer_prediction(session, site_id=premier.site_id) + await creer_prediction(session, site_id=second.site_id) + + resultats = await depot.list_since( + since=datetime(2026, 8, 1, tzinfo=UTC), site_id=premier.site_id + ) + identifiants = [p.prediction_id for p in resultats] + await session.rollback() + + assert identifiants == [voulue.prediction_id] + + async def test_latest_by_site_keeps_only_the_most_recent_target(session: AsyncSession) -> None: site = await creer_site(session) depot = PredictionRepository(session) diff --git a/apps/backend/tests/repositories/test_reading.py b/apps/backend/tests/repositories/test_reading.py index 4f12df0..aaff856 100644 --- a/apps/backend/tests/repositories/test_reading.py +++ b/apps/backend/tests/repositories/test_reading.py @@ -155,6 +155,56 @@ async def test_latest_for_site_ignores_the_readings_of_the_other_sites( assert trouvee is None +async def test_list_since_orders_by_site_then_by_time_ascending(session: AsyncSession) -> None: + site = await creer_site(session) + depot = ReadingRepository(session) + plus_recente = await creer_lecture( + session, site_id=site.site_id, timestamp=datetime(2026, 9, 16, tzinfo=UTC) + ) + plus_ancienne = await creer_lecture( + session, site_id=site.site_id, timestamp=datetime(2026, 9, 15, tzinfo=UTC) + ) + + 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 == [plus_ancienne.reading_id, plus_recente.reading_id] + + +async def test_list_since_excludes_readings_before_the_cutoff(session: AsyncSession) -> None: + site = await creer_site(session) + depot = ReadingRepository(session) + dedans = await creer_lecture( + session, site_id=site.site_id, timestamp=datetime(2026, 9, 16, tzinfo=UTC) + ) + await creer_lecture(session, site_id=site.site_id, timestamp=datetime(2026, 9, 1, tzinfo=UTC)) + + resultats = await depot.list_since( + since=datetime(2026, 9, 10, tzinfo=UTC), site_id=site.site_id + ) + identifiants = [r.reading_id for r in resultats] + await session.rollback() + + assert identifiants == [dedans.reading_id] + + +async def test_list_since_filters_by_site_id(session: AsyncSession) -> None: + premier = await creer_site(session) + second = await creer_site(session) + depot = ReadingRepository(session) + voulue = await creer_lecture(session, site_id=premier.site_id) + await creer_lecture(session, site_id=second.site_id) + + resultats = await depot.list_since( + since=datetime(2026, 8, 1, tzinfo=UTC), site_id=premier.site_id + ) + identifiants = [r.reading_id for r in resultats] + await session.rollback() + + assert identifiants == [voulue.reading_id] + + async def test_list_history_orders_the_readings_by_timestamp_descending( session: AsyncSession, ) -> None: diff --git a/apps/backend/tests/services/test_alert.py b/apps/backend/tests/services/test_alert.py index 4a88802..c99fd6f 100644 --- a/apps/backend/tests/services/test_alert.py +++ b/apps/backend/tests/services/test_alert.py @@ -1,7 +1,10 @@ -from datetime import UTC, datetime +from dataclasses import dataclass, field +from datetime import UTC, datetime, timedelta from app.models.energy import Alert -from app.services.alert import AlertService +from app.services.alert import OUTAGE_THRESHOLD, AlertService, _severity_from_ratio + +NOW = datetime(2026, 9, 16, 12, 0, tzinfo=UTC) def alert( @@ -26,10 +29,36 @@ def alert( ) +@dataclass +class FauxSite: + site_id: str + capacity_kw: float | None = None + + +@dataclass +class FauxLecture: + site_id: str + timestamp: datetime + consumption_kw: float | None = None + consumption_kwh: float | None = None + data_quality: str | None = None + null_reasons: list[str] | None = None + + +@dataclass +class FauxPrediction: + site_id: str + target_at: datetime + predicted_value: float | None + target_metric: str = "consumption_kwh" + prediction_id: int = 1 + + class FakeRepository: def __init__(self, alerts: list[Alert]) -> None: self._alerts = alerts self.appels: list[tuple[str | None, str | None]] = [] + self.crees: list[Alert] = [] async def list_all( self, *, site_id: str | None = None, severity: str | None = None @@ -37,19 +66,325 @@ class FakeRepository: self.appels.append((site_id, severity)) return self._alerts + async def create_many(self, alerts: list[Alert]) -> list[Alert]: + self.crees = list(alerts) + return self.crees + + +@dataclass +class FauxDepotLectures: + depuis: list[FauxLecture] = field(default_factory=list) + dernieres: list[FauxLecture] = field(default_factory=list) + + async def list_since(self, *, since: datetime, site_id: str | None = None) -> list[FauxLecture]: + return [lecture for lecture in self.depuis if site_id is None or lecture.site_id == site_id] + + async def latest_by_site(self) -> list[FauxLecture]: + return self.dernieres + + +@dataclass +class FauxDepotPredictions: + predictions: list[FauxPrediction] = field(default_factory=list) + + async def list_since( + self, *, since: datetime, site_id: str | None = None + ) -> list[FauxPrediction]: + return [p for p in self.predictions if site_id is None or p.site_id == site_id] + + +@dataclass +class FauxDepotSites: + sites: list[FauxSite] + + async def list_all(self) -> list[FauxSite]: + return self.sites + + +def service( + *, + sites: list[FauxSite], + lectures: list[FauxLecture] | None = None, + dernieres: list[FauxLecture] | None = None, + predictions: list[FauxPrediction] | None = None, + alerts: FakeRepository | None = None, +) -> tuple[AlertService, FakeRepository]: + depot_alertes = alerts or FakeRepository([]) + dernieres_lectures = dernieres if dernieres is not None else (lectures or []) + return ( + AlertService( + alerts=depot_alertes, # type: ignore[arg-type] + readings=FauxDepotLectures(depuis=lectures or [], dernieres=dernieres_lectures), # type: ignore[arg-type] + predictions=FauxDepotPredictions(predictions or []), # type: ignore[arg-type] + sites=FauxDepotSites(sites), # type: ignore[arg-type] + ), + depot_alertes, + ) + async def test_list_all_returns_the_repository_alerts() -> None: - service = AlertService(alerts=FakeRepository([alert(1), alert(2)])) + svc, _ = service(sites=[], alerts=FakeRepository([alert(1), alert(2)])) - alertes = await service.list_all() + alertes = await svc.list_all() assert [a.alert_id for a in alertes] == [1, 2] async def test_list_all_relays_the_filters_to_the_repository() -> None: depot = FakeRepository([]) - service = AlertService(alerts=depot) + svc, _ = service(sites=[], alerts=depot) - await service.list_all(site_id="site-1", severity="critical") + await svc.list_all(site_id="site-1", severity="critical") assert depot.appels == [("site-1", "critical")] + + +async def test_detect_raises_a_threshold_alert_above_site_capacity() -> None: + svc, depot = service( + sites=[FauxSite("A", capacity_kw=100.0)], + lectures=[FauxLecture("A", NOW, consumption_kw=150.0)], + ) + + await svc.detect(now=NOW) + + (candidate,) = depot.crees + assert candidate.type == "threshold" + assert candidate.severity == "high" + assert candidate.value == 150.0 + assert candidate.threshold == 100.0 + assert candidate.metric == "consumption_kw" + + +async def test_detect_ignores_a_reading_within_capacity() -> None: + svc, depot = service( + sites=[FauxSite("A", capacity_kw=100.0)], + lectures=[FauxLecture("A", NOW, consumption_kw=80.0)], + ) + + await svc.detect(now=NOW) + + assert depot.crees == [] + + +async def test_detect_ignores_threshold_when_the_site_has_no_declared_capacity() -> None: + svc, depot = service( + sites=[FauxSite("A", capacity_kw=None)], + lectures=[FauxLecture("A", NOW, consumption_kw=9999.0)], + ) + + await svc.detect(now=NOW) + + assert depot.crees == [] + + +async def test_detect_raises_a_spike_alert_on_a_brutal_consecutive_variation() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[ + FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=100.0), + FauxLecture("A", NOW, consumption_kw=160.0), + ], + ) + + await svc.detect(now=NOW) + + (candidate,) = [a for a in depot.crees if a.type == "spike"] + assert candidate.value == 160.0 + assert candidate.threshold == 100.0 + assert candidate.timestamp == NOW + + +async def test_detect_ignores_a_moderate_consecutive_variation() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[ + FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=100.0), + FauxLecture("A", NOW, consumption_kw=110.0), + ], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "spike"] == [] + + +async def test_detect_never_compares_consecutive_readings_across_two_sites() -> None: + svc, depot = service( + sites=[FauxSite("A"), FauxSite("B")], + lectures=[ + FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=10.0), + FauxLecture("B", NOW, consumption_kw=1000.0), + ], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "spike"] == [] + + +async def test_detect_raises_an_anomaly_alert_far_from_the_matching_prediction() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, consumption_kwh=100.0)], + predictions=[FauxPrediction("A", target_at=NOW, predicted_value=70.0)], + ) + + await svc.detect(now=NOW) + + (candidate,) = [a for a in depot.crees if a.type == "anomaly"] + assert candidate.value == 100.0 + assert candidate.threshold == 70.0 + assert candidate.metric == "consumption_kwh" + assert candidate.prediction_id == 1 + + +async def test_detect_ignores_a_reading_close_to_its_prediction() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, consumption_kwh=100.0)], + predictions=[FauxPrediction("A", target_at=NOW, predicted_value=95.0)], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "anomaly"] == [] + + +async def test_detect_ignores_a_prediction_whose_target_at_does_not_match_the_reading() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, consumption_kwh=100.0)], + predictions=[FauxPrediction("A", target_at=NOW - timedelta(hours=1), predicted_value=1.0)], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "anomaly"] == [] + + +async def test_detect_raises_an_outage_alert_past_the_threshold() -> None: + derniere = NOW - OUTAGE_THRESHOLD - timedelta(minutes=1) + svc, depot = service( + sites=[FauxSite("A")], + lectures=[], + dernieres=[FauxLecture("A", derniere)], + ) + + await svc.detect(now=NOW) + + (candidate,) = [a for a in depot.crees if a.type == "outage"] + assert candidate.severity in {"low", "medium", "high", "critical"} + + +async def test_detect_ignores_a_site_still_within_the_outage_threshold() -> None: + derniere = NOW - OUTAGE_THRESHOLD + timedelta(minutes=1) + svc, depot = service( + sites=[FauxSite("A")], + lectures=[], + dernieres=[FauxLecture("A", derniere)], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "outage"] == [] + + +async def test_detect_raises_a_critical_outage_alert_for_a_site_never_read() -> None: + svc, depot = service(sites=[FauxSite("A")], lectures=[], dernieres=[]) + + await svc.detect(now=NOW) + + (candidate,) = [a for a in depot.crees if a.type == "outage"] + assert candidate.severity == "critical" + assert candidate.source_alert_id == "outage:jamais" + + +async def test_detect_raises_a_sensor_alert_on_a_degraded_reading() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, data_quality="critical", null_reasons=["missing:x"])], + ) + + await svc.detect(now=NOW) + + (candidate,) = [a for a in depot.crees if a.type == "sensor"] + assert candidate.severity == "critical" + + +async def test_detect_ignores_a_good_quality_reading_for_the_sensor_rule() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, data_quality="good")], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "sensor"] == [] + + +async def test_detect_scopes_to_a_single_site_when_asked() -> None: + svc, depot = service( + sites=[FauxSite("A", capacity_kw=100.0), FauxSite("B", capacity_kw=100.0)], + lectures=[ + FauxLecture("A", NOW, consumption_kw=150.0), + FauxLecture("B", NOW, consumption_kw=150.0), + ], + ) + + await svc.detect(now=NOW, site_id="A") + + assert {a.site_id for a in depot.crees} == {"A"} + + +async def test_detect_returns_early_when_there_is_no_site() -> None: + svc, depot = service(sites=[]) + + resultat = await svc.detect(now=NOW) + + assert resultat == [] + assert depot.crees == [] + + +async def test_detect_ignores_a_spike_when_the_previous_reading_is_zero() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[ + FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=0.0), + FauxLecture("A", NOW, consumption_kw=50.0), + ], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "spike"] == [] + + +async def test_detect_ignores_an_anomaly_when_the_prediction_is_near_zero() -> None: + svc, depot = service( + sites=[FauxSite("A")], + lectures=[FauxLecture("A", NOW, consumption_kwh=5.0)], + predictions=[FauxPrediction("A", target_at=NOW, predicted_value=0.0)], + ) + + await svc.detect(now=NOW) + + assert [a for a in depot.crees if a.type == "anomaly"] == [] + + +def test_severity_from_ratio_covers_every_band() -> None: + assert _severity_from_ratio(1.0) == "low" + assert _severity_from_ratio(1.2) == "medium" + assert _severity_from_ratio(1.5) == "high" + assert _severity_from_ratio(2.0) == "critical" + + +async def test_detect_does_not_call_create_many_when_nothing_triggers() -> None: + svc, depot = service( + sites=[FauxSite("A", capacity_kw=100.0)], + lectures=[FauxLecture("A", NOW, consumption_kw=10.0, data_quality="good")], + ) + + resultat = await svc.detect(now=NOW) + + assert resultat == [] + assert depot.crees == [] diff --git a/apps/backend/tests/test_internal_alerts.py b/apps/backend/tests/test_internal_alerts.py new file mode 100644 index 0000000..4740150 --- /dev/null +++ b/apps/backend/tests/test_internal_alerts.py @@ -0,0 +1,65 @@ +from datetime import UTC, datetime + +import pytest +from sqlalchemy.ext.asyncio import AsyncSession + +from app.detection import internal_alerts +from app.repositories.alert import AlertRepository +from tests.repositories.test_reading import creer_lecture +from tests.repositories.test_site import creer as creer_site + + +def test_parse_args_defaults_to_no_site_and_no_instant() -> None: + arguments = internal_alerts.parse_args([]) + + assert arguments.site_id is None + assert arguments.now is None + + +def test_parse_args_reads_the_site_id() -> None: + arguments = internal_alerts.parse_args(["--site-id", "site-1"]) + + assert arguments.site_id == "site-1" + + +def test_parse_args_parses_the_instant_option() -> None: + arguments = internal_alerts.parse_args(["--now", "2026-09-16T12:00:00+00:00"]) + + assert arguments.now == datetime(2026, 9, 16, 12, tzinfo=UTC) + + +def test_parse_instant_treats_a_naive_datetime_as_utc() -> None: + assert internal_alerts._parse_instant("2026-09-16T12:00:00") == datetime( + 2026, 9, 16, 12, tzinfo=UTC + ) + + +def test_main_prints_how_many_alerts_were_recorded( + monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str] +) -> None: + async def fausse_execution(*, now: datetime | None, site_id: str | None) -> int: + return 3 + + monkeypatch.setattr(internal_alerts, "run_detection", fausse_execution) + + code = internal_alerts.main([]) + + assert code == 0 + assert "3 nouvelle" in capsys.readouterr().out + + +@pytest.mark.integration +async def test_run_detection_writes_a_threshold_alert_end_to_end(session: AsyncSession) -> None: + site = await creer_site(session, capacity_kw=100.0) + instant = datetime(2026, 9, 16, 12, tzinfo=UTC) + await creer_lecture(session, site_id=site.site_id, timestamp=instant, consumption_kw=150.0) + await session.commit() + + nombre = await internal_alerts.run_detection(now=instant, site_id=site.site_id) + + alertes = await AlertRepository(session).list_all(site_id=site.site_id) + types = [a.type for a in alertes] + await session.rollback() + + assert nombre == 1 + assert types == ["threshold"] diff --git a/docs/architecture/20-backend.md b/docs/architecture/20-backend.md index 09d3b3d..8c3641d 100644 --- a/docs/architecture/20-backend.md +++ b/docs/architecture/20-backend.md @@ -207,6 +207,35 @@ par exemple `limit` hors bornes). Un datetime sans fuseau dans `start`/`end` est l'UTC plutôt que rejeté : le comparer tel quel à `reading.timestamp` (`timestamptz`) échouerait côté pilote, en `500` plutôt qu'un refus propre. +### Détection d'alertes internes + +`AlertService` n'est plus lecture seule : `AlertService.detect()` compare les `reading` (et, pour +le type `anomaly`, les `prediction`) des dernières 48h (`LOOKBACK`) à cinq règles et enregistre une +ligne `alert` par déclenchement, avec `source="enervision"`. `metric`/`value`/`threshold` gardent +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 | +| `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` | + +La sévérité de chaque alerte (hors `sensor`, dérivée directement de `data_quality`) suit le même +barème par ratio observé/seuil : `low` sous 1.2, `medium` sous 1.5, `high` sous 2.0, `critical` +au-delà. `AlertRepository.create_many()` insère par lot avec `ON CONFLICT DO NOTHING` sur +`uq_alert_source_reference`, et `source_alert_id` est construit de façon déterministe (règle + +horodatage) : rejouer la détection sur une fenêtre déjà analysée ne duplique donc jamais une +alerte. + +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 +connexion SQL directe (contrairement à `app/etl/historical_import.py`). Cette issue (#104) +débloquait #38 (moteur de règles pour recommandations), dont la FK `alert_id` `NOT NULL` n'avait +jusqu'ici rien à référencer côté `source="enervision"`. + ### `/health/ready` Cette sonde porte une garde décrite dans l'[ADR 0001](../adr/0001-postgresql-timescaledb.md) : un From c059f838bbbc45c440bde45e498a94c1f42a365f Mon Sep 17 00:00:00 2001 From: Dorian Date: Fri, 18 Sep 2026 16:58:03 +0200 Subject: [PATCH 3/3] 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