feat(backend): detecte les alertes internes a partir des lectures et previsions
This commit is contained in:
@@ -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)]
|
||||
|
||||
@@ -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())
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
*,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user