Compare commits

...
Author SHA1 Message Date
Dorian c059f838bb fix(backend): fiabilise le tri des lectures/predictions et la detection de redemarrage a zero
Backend / Tests exigeant une base (push) Failing after 38s
Backend / Lint, typage et tests (push) Successful in 1m29s
Backend / Audit des dépendances (push) Successful in 1m3s
SonarQube / build-back (push) Successful in 1m8s
SonarQube / build-front (push) Successful in 9m36s
SonarQube / test-back (push) Failing after 1m5s
SonarQube / test-front (push) Failing after 5m15s
SonarQube / SonarQube (push) Skipped
2026-09-18 16:58:03 +02:00
Dorian a9e124a97d feat(backend): detecte les alertes internes a partir des lectures et previsions 2026-09-18 16:10:06 +02:00
Johan LEROYandGitHub 619024f547 Merge pull request #105 from ineszang/feat/service-de-scoring
feat(ml,backend): implemente le service de scoring et GET /predictions
2026-09-18 15:33:12 +02:00
Dorian 460d6c1b3e Merge remote-tracking branch 'origin/dev' into feat/service-de-scoring 2026-09-18 15:06:41 +02:00
Dorian eb4291b10a fix(ml,backend,frontend): borne la peremption des predictions et isole les erreurs par flux 2026-09-18 14:58:39 +02:00
Johan LEROYandGitHub 9a72bb3a8f Merge pull request #108 from ineszang/test/matrice-acces-roles
test(backend): matrice d'accès par rôle et CI d'intégration
2026-09-18 14:25:40 +02:00
Dorian db81290026 Merge remote-tracking branch 'origin/dev' into feat/service-de-scoring 2026-09-18 12:05:38 +02:00
Dorian d9229e5a93 feat(frontend): affiche les prévisions de consommation sur le dashboard 2026-09-18 11:51:02 +02:00
Dorian e9376a98bf feat(ml,backend): implemente le service de scoring et GET /predictions (#37) 2026-09-18 11:06:04 +02:00
40 changed files with 2994 additions and 62 deletions
+4 -1
View File
@@ -6,7 +6,7 @@ ML := ml
.PHONY: help install install-backend install-frontend install-ml dev dev-backend dev-frontend \ .PHONY: help install install-backend install-frontend install-ml dev dev-backend dev-frontend \
lint format typecheck test test-cov test-integration check \ lint format typecheck test test-cov test-integration check \
openapi docker-build db-up db-down db-reset db-logs db-psql migrate bootstrap-admin \ openapi docker-build db-up db-down db-reset db-logs db-psql migrate bootstrap-admin \
ml-lint ml-typecheck ml-test ml-check ml-train ml-lint ml-typecheck ml-test ml-check ml-train ml-score
help: ## Liste les cibles disponibles help: ## Liste les cibles disponibles
@grep -E '^[a-zA-Z_-]+:.*?## .*$$' $(MAKEFILE_LIST) | awk 'BEGIN {FS = ":.*?## "}; {printf " \033[36m%-16s\033[0m %s\n", $$1, $$2}' @grep -E '^[a-zA-Z_-]+:.*?## .*$$' $(MAKEFILE_LIST) | awk 'BEGIN {FS = ":.*?## "}; {printf " \033[36m%-16s\033[0m %s\n", $$1, $$2}'
@@ -74,6 +74,9 @@ ml-check: ml-lint ml-typecheck ml-test ## Chaîne de vérification complète du
ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DATABASE_URL ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DATABASE_URL
cd $(ML) && uv run python -m enervision_ml.train $(if $(CSV),--csv $(CSV),) cd $(ML) && uv run python -m enervision_ml.train $(if $(CSV),--csv $(CSV),)
ml-score: ## Score le prochain pas horaire et l'ecrit dans `prediction`. CSV=chemin optionnel
cd $(ML) && uv run python -m enervision_ml.score $(if $(CSV),--csv $(CSV),)
docker-build: ## Construit l'image du backend docker-build: ## Construit l'image du backend
docker build -t enervision-backend:local $(BACKEND) docker build -t enervision-backend:local $(BACKEND)
+17 -1
View File
@@ -27,6 +27,7 @@ from app.repositories.audit_log import AuditLogRepository
from app.repositories.login_attempt import LoginAttemptRepository from app.repositories.login_attempt import LoginAttemptRepository
from app.repositories.password_reset_attempt import PasswordResetAttemptRepository from app.repositories.password_reset_attempt import PasswordResetAttemptRepository
from app.repositories.password_reset_token import PasswordResetTokenRepository from app.repositories.password_reset_token import PasswordResetTokenRepository
from app.repositories.prediction import PredictionRepository
from app.repositories.reading import ReadingRepository from app.repositories.reading import ReadingRepository
from app.repositories.recommendation import RecommendationRepository from app.repositories.recommendation import RecommendationRepository
from app.repositories.refresh_token import RefreshTokenRepository from app.repositories.refresh_token import RefreshTokenRepository
@@ -34,6 +35,7 @@ from app.repositories.site import SiteRepository
from app.repositories.user import UserRepository from app.repositories.user import UserRepository
from app.services.alert import AlertService from app.services.alert import AlertService
from app.services.auth import AuthService, LoginPolicy, PasswordResetPolicy from app.services.auth import AuthService, LoginPolicy, PasswordResetPolicy
from app.services.prediction import PredictionService
from app.services.reading import ReadingService from app.services.reading import ReadingService
from app.services.recommendation import RecommendationService from app.services.recommendation import RecommendationService
from app.services.sensor import SensorService from app.services.sensor import SensorService
@@ -178,7 +180,12 @@ SiteServiceDep = Annotated[SiteService, Depends(get_site_service)]
def get_alert_service(session: SessionDep) -> AlertService: 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)] AlertServiceDep = Annotated[AlertService, Depends(get_alert_service)]
@@ -212,6 +219,15 @@ def get_sensor_service(session: SessionDep) -> SensorService:
SensorServiceDep = Annotated[SensorService, Depends(get_sensor_service)] SensorServiceDep = Annotated[SensorService, Depends(get_sensor_service)]
def get_prediction_service(session: SessionDep) -> PredictionService:
return PredictionService(
sites=SiteRepository(session), predictions=PredictionRepository(session)
)
PredictionServiceDep = Annotated[PredictionService, Depends(get_prediction_service)]
async def get_current_principal( async def get_current_principal(
credentials: CredentialsDep, credentials: CredentialsDep,
session: SessionDep, session: SessionDep,
+7
View File
@@ -83,6 +83,13 @@ TAGS: Final[list[dict[str, Any]]] = [
"name": "sensors", "name": "sensors",
"description": "État de santé des capteurs par site. Réservé au rôle `admin`.", "description": "État de santé des capteurs par site. Réservé au rôle `admin`.",
}, },
{
"name": "predictions",
"description": (
"Dernière prévision de consommation par site, calculée hors ligne par le pipeline "
"de scoring (`ml/`) et simplement lue ici. Accessible à partir du rôle `lecteur`."
),
},
] ]
cookie_de_rafraichissement = APIKeyCookie( cookie_de_rafraichissement = APIKeyCookie(
@@ -0,0 +1,18 @@
from fastapi import APIRouter
from app.api.deps import LecteurDep, PredictionServiceDep
from app.schemas.prediction import PredictionSummaryResponse
router = APIRouter()
@router.get(
"",
response_model=PredictionSummaryResponse,
summary="Dernière prédiction de consommation par site",
)
async def get_predictions(
_: LecteurDep, service: PredictionServiceDep
) -> PredictionSummaryResponse:
resume = await service.summary()
return PredictionSummaryResponse.model_validate(resume)
+4
View File
@@ -5,6 +5,7 @@ from app.api.v1.endpoints import (
alerts, alerts,
auth, auth,
health, health,
predictions,
readings, readings,
recommendations, recommendations,
sensors, sensors,
@@ -34,3 +35,6 @@ api_router.include_router(
api_router.include_router( api_router.include_router(
sensors.router, prefix="/sensors", tags=["sensors"], responses=REPONSES_ADMIN sensors.router, prefix="/sensors", tags=["sensors"], responses=REPONSES_ADMIN
) )
api_router.include_router(
predictions.router, prefix="/predictions", tags=["predictions"], responses=REPONSES_LECTEUR
)
@@ -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())
+34
View File
@@ -1,6 +1,7 @@
from collections.abc import Sequence from collections.abc import Sequence
from sqlalchemy import select from sqlalchemy import select
from sqlalchemy.dialects.postgresql import insert
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.models.energy import Alert from app.models.energy import Alert
@@ -19,3 +20,36 @@ class AlertRepository:
if severity is not None: if severity is not None:
requete = requete.where(Alert.severity == severity) requete = requete.where(Alert.severity == severity)
return (await self._session.scalars(requete)).all() 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()
@@ -0,0 +1,49 @@
from collections.abc import Sequence
from datetime import datetime
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.energy import Prediction
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).
# 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, Prediction.prediction_id)
)
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
# `ReadingRepository.latest_by_site`. Trié sur `target_at` (couvert par
# `ix_prediction_site_target`) plutôt que `created_at` : c'est la prévision la plus
# récente qui compte pour un tableau de bord, pas forcément le dernier run de scoring.
requete = (
select(Prediction)
.distinct(Prediction.site_id)
.order_by(
Prediction.site_id,
Prediction.target_at.desc(),
Prediction.prediction_id.desc(),
)
)
return (await self._session.scalars(requete)).all()
+15
View File
@@ -34,6 +34,21 @@ class ReadingRepository:
lecture: Reading | None = await self._session.scalar(requete) lecture: Reading | None = await self._session.scalar(requete)
return lecture 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.
# `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, Reading.reading_id)
)
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( async def list_history(
self, self,
*, *,
+43
View File
@@ -0,0 +1,43 @@
from datetime import datetime
from enum import StrEnum
from pydantic import BaseModel, ConfigDict
class PredictionTargetMetric(StrEnum):
CONSUMPTION_KWH = "consumption_kwh"
CONSUMPTION_KW = "consumption_kw"
class PredictionStatus(StrEnum):
AVAILABLE = "available"
INSUFFICIENT_DATA = "insufficient_data"
ERROR = "error"
class SitePredictionResponse(BaseModel):
model_config = ConfigDict(from_attributes=True)
target_at: datetime
target_metric: PredictionTargetMetric
period_minutes: int | None
predicted_value: float | None
status: PredictionStatus
failure_reason: str | None
model_reference: str
created_at: datetime
class SitePredictionSummaryResponse(BaseModel):
model_config = ConfigDict(from_attributes=True)
site_id: str
site_name: str
prediction: SitePredictionResponse | None
class PredictionSummaryResponse(BaseModel):
model_config = ConfigDict(from_attributes=True)
timestamp: datetime
sites: list[SitePredictionSummaryResponse]
+311 -2
View File
@@ -1,14 +1,323 @@
from collections.abc import Sequence 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.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: class AlertService:
def __init__(self, *, alerts: AlertRepository) -> None: def __init__(
self,
*,
alerts: AlertRepository,
readings: ReadingRepository,
predictions: PredictionRepository,
sites: SiteRepository,
) -> None:
self._alerts = alerts self._alerts = alerts
self._readings = readings
self._predictions = predictions
self._sites = sites
async def list_all( async def list_all(
self, *, site_id: str | None = None, severity: str | None = None self, *, site_id: str | None = None, severity: str | None = None
) -> Sequence[Alert]: ) -> Sequence[Alert]:
return await self._alerts.list_all(site_id=site_id, severity=severity) 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, 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
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:
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(
_spike_alert(
lecture,
avant,
apres,
severity=_severity_from_ratio(variation / SPIKE_RELATIVE_THRESHOLD),
)
)
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`.
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
+68
View File
@@ -0,0 +1,68 @@
from dataclasses import dataclass
from datetime import UTC, datetime
from app.models.energy import Prediction, Site
from app.repositories.prediction import PredictionRepository
from app.repositories.site import SiteRepository
@dataclass(frozen=True, slots=True)
class SitePrediction:
target_at: datetime
target_metric: str
period_minutes: int | None
predicted_value: float | None
status: str
failure_reason: str | None
model_reference: str
created_at: datetime
@dataclass(frozen=True, slots=True)
class SitePredictionSummary:
site_id: str
site_name: str
prediction: SitePrediction | None
@dataclass(frozen=True, slots=True)
class PredictionSummary:
timestamp: datetime
sites: list[SitePredictionSummary]
class PredictionService:
def __init__(self, sites: SiteRepository, predictions: PredictionRepository) -> None:
self._sites = sites
self._predictions = predictions
async def summary(self) -> PredictionSummary:
sites = await self._sites.list_all()
dernieres = {p.site_id: p for p in await self._predictions.latest_by_site()}
return PredictionSummary(
timestamp=datetime.now(UTC),
sites=[_resume_site(site, dernieres.get(site.site_id)) for site in sites],
)
def _resume_site(site: Site, derniere: Prediction | None) -> SitePredictionSummary:
# Piège : l'absence de ligne signifie « jamais scoré », pas une valeur pseudo-statut, qui
# n'existe pas dans la contrainte de la table. `prediction` reste `None` plutôt que de
# fabriquer un statut absent du domaine `available`/`insufficient_data`/`error`.
prediction = None
if derniere is not None:
prediction = SitePrediction(
target_at=derniere.target_at,
target_metric=derniere.target_metric,
period_minutes=derniere.period_minutes,
predicted_value=derniere.predicted_value,
status=derniere.status,
failure_reason=derniere.failure_reason,
model_reference=derniere.model_reference,
created_at=derniere.created_at,
)
return SitePredictionSummary(
site_id=site.site_id, site_name=site.site_name, prediction=prediction
)
+197
View File
@@ -1719,6 +1719,62 @@
} }
] ]
} }
},
"/api/v1/predictions": {
"get": {
"tags": [
"predictions"
],
"summary": "Dernière prédiction de consommation par site",
"operationId": "get_predictions_api_v1_predictions_get",
"responses": {
"200": {
"description": "Successful Response",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/PredictionSummaryResponse"
}
}
}
},
"500": {
"description": "Erreur interne. `correlation` identifie la trace côté serveur, qui n'est pas renvoyée au client.",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/InternalErrorResponse"
}
}
}
},
"401": {
"description": "Jeton absent, illisible, périmé, ou rendu caduc par un changement de rôle ou une désactivation. L'en-tête `WWW-Authenticate` porte la cause dans `error=`.",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/ErrorResponse"
}
}
}
},
"403": {
"description": "Mot de passe provisoire à changer (`detail` vaut `password_change_required`).",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/ErrorResponse"
}
}
}
}
},
"security": [
{
"Jeton d'accès": []
}
]
}
} }
}, },
"components": { "components": {
@@ -1972,6 +2028,45 @@
], ],
"title": "PasswordChangeRequest" "title": "PasswordChangeRequest"
}, },
"PredictionStatus": {
"type": "string",
"enum": [
"available",
"insufficient_data",
"error"
],
"title": "PredictionStatus"
},
"PredictionSummaryResponse": {
"properties": {
"timestamp": {
"type": "string",
"format": "date-time",
"title": "Timestamp"
},
"sites": {
"items": {
"$ref": "#/components/schemas/SitePredictionSummaryResponse"
},
"type": "array",
"title": "Sites"
}
},
"type": "object",
"required": [
"timestamp",
"sites"
],
"title": "PredictionSummaryResponse"
},
"PredictionTargetMetric": {
"type": "string",
"enum": [
"consumption_kwh",
"consumption_kw"
],
"title": "PredictionTargetMetric"
},
"PrincipalResponse": { "PrincipalResponse": {
"properties": { "properties": {
"id": { "id": {
@@ -2518,6 +2613,104 @@
], ],
"title": "SiteCurrentResponse" "title": "SiteCurrentResponse"
}, },
"SitePredictionResponse": {
"properties": {
"target_at": {
"type": "string",
"format": "date-time",
"title": "Target At"
},
"target_metric": {
"$ref": "#/components/schemas/PredictionTargetMetric"
},
"period_minutes": {
"anyOf": [
{
"type": "integer"
},
{
"type": "null"
}
],
"title": "Period Minutes"
},
"predicted_value": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Predicted Value"
},
"status": {
"$ref": "#/components/schemas/PredictionStatus"
},
"failure_reason": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"title": "Failure Reason"
},
"model_reference": {
"type": "string",
"title": "Model Reference"
},
"created_at": {
"type": "string",
"format": "date-time",
"title": "Created At"
}
},
"type": "object",
"required": [
"target_at",
"target_metric",
"period_minutes",
"predicted_value",
"status",
"failure_reason",
"model_reference",
"created_at"
],
"title": "SitePredictionResponse"
},
"SitePredictionSummaryResponse": {
"properties": {
"site_id": {
"type": "string",
"title": "Site Id"
},
"site_name": {
"type": "string",
"title": "Site Name"
},
"prediction": {
"anyOf": [
{
"$ref": "#/components/schemas/SitePredictionResponse"
},
{
"type": "null"
}
]
}
},
"type": "object",
"required": [
"site_id",
"site_name",
"prediction"
],
"title": "SitePredictionSummaryResponse"
},
"SiteResponse": { "SiteResponse": {
"properties": { "properties": {
"site_id": { "site_id": {
@@ -2973,6 +3166,10 @@
{ {
"name": "sensors", "name": "sensors",
"description": "État de santé des capteurs par site. Réservé au rôle `admin`." "description": "État de santé des capteurs par site. Réservé au rôle `admin`."
},
{
"name": "predictions",
"description": "Dernière prévision de consommation par site, calculée hors ligne par le pipeline de scoring (`ml/`) et simplement lue ici. Accessible à partir du rôle `lecteur`."
} }
] ]
} }
+1
View File
@@ -53,6 +53,7 @@ ROLE_MINIMUM: Final[dict[Route, Role]] = {
("GET", "/api/v1/recommendations/{recommendation_id}"): Role.LECTEUR, ("GET", "/api/v1/recommendations/{recommendation_id}"): Role.LECTEUR,
("GET", "/api/v1/stats/summary"): Role.LECTEUR, ("GET", "/api/v1/stats/summary"): Role.LECTEUR,
("GET", "/api/v1/readings"): Role.LECTEUR, ("GET", "/api/v1/readings"): Role.LECTEUR,
("GET", "/api/v1/predictions"): Role.LECTEUR,
("GET", "/api/v1/sensors/status"): Role.ADMIN, ("GET", "/api/v1/sensors/status"): Role.ADMIN,
("GET", "/api/v1/users"): Role.ADMIN, ("GET", "/api/v1/users"): Role.ADMIN,
("POST", "/api/v1/users"): Role.ADMIN, ("POST", "/api/v1/users"): Role.ADMIN,
@@ -0,0 +1,85 @@
from collections.abc import Callable, Iterator
from datetime import UTC, datetime
from uuid import uuid4
import pytest
from fastapi import FastAPI
from httpx import AsyncClient
from app.api.deps import get_current_principal, get_prediction_service
from app.core.principal import Principal
from app.core.roles import AccountKind, Role
from app.services.prediction import PredictionSummary, SitePrediction, SitePredictionSummary
TARGET_AT = datetime(2026, 9, 16, 13, 0, tzinfo=UTC)
CREATED_AT = datetime(2026, 9, 16, 12, 0, tzinfo=UTC)
def lecteur() -> Principal:
# Le garde-fou de rôle (`lecteur` minimum) est déjà couvert par l'ensemble `ROUTES_A_ROLE`
# de `tests/api/test_openapi.py` : pas besoin ici d'un paramètre de rôle jamais appelé avec
# autre chose que sa valeur par défaut.
return Principal(
id=uuid4(),
email="lecteur@enervision.fr",
role=Role.LECTEUR,
kind=AccountKind.HUMAIN,
must_change_password=False,
)
class FauxService:
def __init__(self) -> None:
self.resume = PredictionSummary(
timestamp=datetime.now(UTC),
sites=[
SitePredictionSummary(
site_id="SITE001",
site_name="Bureau Paris La Défense",
prediction=SitePrediction(
target_at=TARGET_AT,
target_metric="consumption_kwh",
period_minutes=60,
predicted_value=812.5,
status="available",
failure_reason=None,
model_reference="lightgbm-abc123",
created_at=CREATED_AT,
),
),
SitePredictionSummary(site_id="SITE002", site_name="Usine Lyon", prediction=None),
],
)
async def summary(self) -> PredictionSummary:
return self.resume
@pytest.fixture
def servi(app: FastAPI) -> Iterator[Callable[[], FauxService]]:
def installe() -> FauxService:
service = FauxService()
app.dependency_overrides[get_prediction_service] = lambda: service
app.dependency_overrides[get_current_principal] = lambda: lecteur()
return service
yield installe
app.dependency_overrides.pop(get_prediction_service, None)
app.dependency_overrides.pop(get_current_principal, None)
async def test_get_predictions_returns_the_service_result(
servi: Callable[[], FauxService], client: AsyncClient
) -> None:
servi()
response = await client.get("/api/v1/predictions")
assert response.status_code == 200
corps = response.json()
premier, second = corps["sites"]
assert premier["site_id"] == "SITE001"
assert premier["prediction"]["predicted_value"] == 812.5
assert premier["prediction"]["status"] == "available"
assert second["site_id"] == "SITE002"
assert second["prediction"] is None
@@ -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()) alertes = await depot.list_all(site_id=identifiant_site())
assert list(alertes) == [] 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 == []
@@ -0,0 +1,171 @@
from datetime import UTC, datetime
import pytest
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.energy import Prediction
from app.repositories.prediction import PredictionRepository
from tests.repositories.test_site import creer as creer_site
from tests.repositories.test_site import identifiant as identifiant_site
pytestmark = pytest.mark.integration
async def creer_prediction(
session: AsyncSession, *, site_id: str, **overrides: object
) -> Prediction:
prediction = Prediction(
site_id=site_id,
target_at=overrides.get("target_at", datetime(2026, 9, 16, tzinfo=UTC)),
target_metric=overrides.get("target_metric", "consumption_kwh"),
period_minutes=overrides.get("period_minutes", 60),
predicted_value=overrides.get("predicted_value", 42.0),
model_reference=overrides.get("model_reference", "lightgbm-test"),
status=overrides.get("status", "available"),
failure_reason=overrides.get("failure_reason"),
)
session.add(prediction)
await session.flush()
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_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)
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)
ancienne = await creer_prediction(
session, site_id=site.site_id, target_at=datetime(2026, 9, 1, tzinfo=UTC)
)
recente = await creer_prediction(
session, site_id=site.site_id, target_at=datetime(2026, 9, 15, tzinfo=UTC)
)
resultats = await depot.latest_by_site()
identifiants = [
p.prediction_id
for p in resultats
if p.prediction_id in (ancienne.prediction_id, recente.prediction_id)
]
await session.rollback()
assert identifiants == [recente.prediction_id]
async def test_latest_by_site_returns_one_row_per_site(session: AsyncSession) -> None:
premier = await creer_site(session)
second = await creer_site(session)
depot = PredictionRepository(session)
voulue_premier = await creer_prediction(session, site_id=premier.site_id)
voulue_second = await creer_prediction(session, site_id=second.site_id)
resultats = await depot.latest_by_site()
identifiants = {p.site_id for p in resultats if p.site_id in (premier.site_id, second.site_id)}
await session.rollback()
assert identifiants == {voulue_premier.site_id, voulue_second.site_id}
async def test_latest_by_site_keeps_an_insufficient_data_prediction(session: AsyncSession) -> None:
site = await creer_site(session)
depot = PredictionRepository(session)
voulue = await creer_prediction(
session,
site_id=site.site_id,
status="insufficient_data",
predicted_value=None,
failure_reason="pas assez d'historique",
)
resultats = await depot.latest_by_site()
identifiants = [p.prediction_id for p in resultats if p.site_id == site.site_id]
await session.rollback()
assert identifiants == [voulue.prediction_id]
async def test_latest_by_site_returns_an_empty_list_when_there_is_nothing(
session: AsyncSession,
) -> None:
depot = PredictionRepository(session)
resultats = [p for p in await depot.latest_by_site() if p.site_id == identifiant_site()]
assert resultats == []
@@ -155,6 +155,79 @@ async def test_latest_for_site_ignores_the_readings_of_the_other_sites(
assert trouvee is None 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_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)
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( async def test_list_history_orders_the_readings_by_timestamp_descending(
session: AsyncSession, session: AsyncSession,
) -> None: ) -> None:
+408 -6
View File
@@ -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.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( 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: class FakeRepository:
def __init__(self, alerts: list[Alert]) -> None: def __init__(self, alerts: list[Alert]) -> None:
self._alerts = alerts self._alerts = alerts
self.appels: list[tuple[str | None, str | None]] = [] self.appels: list[tuple[str | None, str | None]] = []
self.crees: list[Alert] = []
async def list_all( async def list_all(
self, *, site_id: str | None = None, severity: str | None = None self, *, site_id: str | None = None, severity: str | None = None
@@ -37,19 +66,392 @@ class FakeRepository:
self.appels.append((site_id, severity)) self.appels.append((site_id, severity))
return self._alerts 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: 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] assert [a.alert_id for a in alertes] == [1, 2]
async def test_list_all_relays_the_filters_to_the_repository() -> None: async def test_list_all_relays_the_filters_to_the_repository() -> None:
depot = FakeRepository([]) 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")] 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_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(
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_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=[
FauxLecture("A", NOW - timedelta(hours=1), consumption_kw=0.0),
FauxLecture("A", NOW, consumption_kw=50.0),
],
)
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"] == []
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 == []
@@ -0,0 +1,121 @@
from dataclasses import dataclass
from datetime import UTC, datetime
from app.services.prediction import PredictionService
TARGET_AT = datetime(2026, 9, 16, 13, 0, tzinfo=UTC)
CREATED_AT = datetime(2026, 9, 16, 12, 0, tzinfo=UTC)
@dataclass
class FauxSite:
site_id: str
site_name: str
@dataclass
class FauxPrediction:
site_id: str
target_at: datetime
target_metric: str
period_minutes: int | None
predicted_value: float | None
status: str
failure_reason: str | None
model_reference: str
created_at: datetime
class FauxDepotSites:
def __init__(self, sites: list[FauxSite]) -> None:
self._sites = sites
async def list_all(self) -> list[FauxSite]:
return self._sites
class FauxDepotPredictions:
def __init__(self, predictions: list[FauxPrediction]) -> None:
self._predictions = predictions
async def latest_by_site(self) -> list[FauxPrediction]:
return self._predictions
def prediction_disponible(site_id: str = "A") -> FauxPrediction:
return FauxPrediction(
site_id=site_id,
target_at=TARGET_AT,
target_metric="consumption_kwh",
period_minutes=60,
predicted_value=812.5,
status="available",
failure_reason=None,
model_reference="lightgbm-abc123",
created_at=CREATED_AT,
)
async def test_summary_attaches_the_latest_prediction_to_its_site() -> None:
service = PredictionService(
sites=FauxDepotSites([FauxSite("A", "Site A")]), # type: ignore[arg-type]
predictions=FauxDepotPredictions([prediction_disponible("A")]), # type: ignore[arg-type]
)
resume = await service.summary()
site = resume.sites[0]
assert site.site_id == "A"
assert site.prediction is not None
assert site.prediction.predicted_value == 812.5
assert site.prediction.status == "available"
async def test_summary_leaves_prediction_none_for_a_site_never_scored() -> None:
service = PredictionService(
sites=FauxDepotSites([FauxSite("A", "Site A")]), # type: ignore[arg-type]
predictions=FauxDepotPredictions([]), # type: ignore[arg-type]
)
resume = await service.summary()
assert resume.sites[0].prediction is None
async def test_summary_carries_an_insufficient_data_prediction_without_a_value() -> None:
insuffisante = FauxPrediction(
site_id="A",
target_at=TARGET_AT,
target_metric="consumption_kwh",
period_minutes=60,
predicted_value=None,
status="insufficient_data",
failure_reason="pas assez d'historique",
model_reference="lightgbm-abc123",
created_at=CREATED_AT,
)
service = PredictionService(
sites=FauxDepotSites([FauxSite("A", "Site A")]), # type: ignore[arg-type]
predictions=FauxDepotPredictions([insuffisante]), # type: ignore[arg-type]
)
resume = await service.summary()
site = resume.sites[0]
assert site.prediction is not None
assert site.prediction.status == "insufficient_data"
assert site.prediction.predicted_value is None
assert site.prediction.failure_reason == "pas assez d'historique"
async def test_summary_covers_every_site_even_with_a_single_prediction_in_the_repository() -> None:
service = PredictionService(
sites=FauxDepotSites([FauxSite("A", "Site A"), FauxSite("B", "Site B")]), # type: ignore[arg-type]
predictions=FauxDepotPredictions([prediction_disponible("A")]), # type: ignore[arg-type]
)
resume = await service.summary()
par_site = {site.site_id: site for site in resume.sites}
assert par_site["A"].prediction is not None
assert par_site["B"].prediction is None
@@ -0,0 +1,87 @@
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
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:
# `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_id, timestamp=instant, consumption_kw=150.0)
await session.commit()
try:
nombre = await internal_alerts.run_detection(now=instant, site_id=site_id)
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"]
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()
@@ -64,4 +64,13 @@ describe('mockApiInterceptor', () => {
httpMock.expectNone(`${environment.apiUrl}/alerts`); httpMock.expectNone(`${environment.apiUrl}/alerts`);
expect((result as unknown[]).length).toBeGreaterThan(0); expect((result as unknown[]).length).toBeGreaterThan(0);
}); });
it('laisse toujours passer /predictions vers le réseau, même avec useMockFixtures activé', () => {
environment.useMockFixtures = true;
http.get(`${environment.apiUrl}/predictions`).subscribe();
const req = httpMock.expectOne(`${environment.apiUrl}/predictions`);
req.flush({ timestamp: '2026-09-18T09:00:00Z', sites: [] });
});
}); });
@@ -26,5 +26,7 @@ export const mockApiInterceptor: HttpInterceptorFn = (req, next) => {
if (req.url.endsWith(`${environment.apiUrl}/alerts`)) { if (req.url.endsWith(`${environment.apiUrl}/alerts`)) {
return of(new HttpResponse({ status: 200, body: ALERTS_FIXTURE })); return of(new HttpResponse({ status: 200, body: ALERTS_FIXTURE }));
} }
// Volontairement jamais mocké, contrairement à `stats`/`alerts` : les prévisions sont servies
// par l'API réelle dès maintenant (au même titre que `/auth/*`, déjà toujours réel).
return next(req); return next(req);
}; };
@@ -0,0 +1,35 @@
import { TestBed } from '@angular/core/testing';
import { provideHttpClient } from '@angular/common/http';
import { provideHttpClientTesting, HttpTestingController } from '@angular/common/http/testing';
import { PredictionsService } from './predictions.service';
import { environment } from '../../../environments/environment';
describe('PredictionsService', () => {
let service: PredictionsService;
let httpMock: HttpTestingController;
beforeEach(() => {
TestBed.configureTestingModule({
providers: [provideHttpClient(), provideHttpClientTesting()],
});
service = TestBed.inject(PredictionsService);
httpMock = TestBed.inject(HttpTestingController);
});
afterEach(() => httpMock.verify());
it('appelle le bon endpoint et retourne un résumé de prévisions', () => {
let result: unknown;
service.getPredictions().subscribe((r) => (result = r));
const req = httpMock.expectOne(`${environment.apiUrl}/predictions`);
expect(req.request.method).toBe('GET');
req.flush({
timestamp: '2026-09-18T09:00:00Z',
sites: [{ site_id: 'SITE001', site_name: 'Test', prediction: null }],
});
expect((result as { sites: unknown[] }).sites.length).toBe(1);
});
});
@@ -0,0 +1,13 @@
import { Service, inject } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { environment } from '../../../environments/environment';
import { PredictionSummary } from '../../shared/models/prediction.model';
@Service()
export class PredictionsService {
private http = inject(HttpClient);
getPredictions() {
return this.http.get<PredictionSummary>(`${environment.apiUrl}/predictions`);
}
}
@@ -21,7 +21,13 @@
</div> </div>
</header> </header>
@if (error(); as message) { @if (statsError(); as message) {
<ev-alert severity="danger" class="banner-error">{{ message }}</ev-alert>
}
@if (alertsError(); as message) {
<ev-alert severity="danger" class="banner-error">{{ message }}</ev-alert>
}
@if (predictionsError(); as message) {
<ev-alert severity="danger" class="banner-error">{{ message }}</ev-alert> <ev-alert severity="danger" class="banner-error">{{ message }}</ev-alert>
} }
@@ -72,4 +78,33 @@
</ul> </ul>
</section> </section>
} }
@if (predictions().length > 0) {
<section class="predictions-section">
<h2>Prévisions de consommation</h2>
<ul class="predictions-list">
@for (site of predictions(); track site.site_id) {
<li class="prediction-item">
<span class="prediction-item__site">{{ site.site_name }}</span>
@if (site.prediction; as prediction) {
@if (prediction.status === 'available') {
<span class="prediction-item__value">
{{ prediction.predicted_value | number: '1.0-1' }} kWh
<span class="prediction-item__target"
>{{ prediction.target_at | date: "dd/MM 'à' HH:mm" }}</span
>
</span>
} @else {
<ev-badge [tone]="badgeToneForPredictionStatus(prediction.status)">{{
prediction.status === 'insufficient_data' ? 'Historique insuffisant' : 'Erreur'
}}</ev-badge>
}
} @else {
<ev-badge tone="neutral">Pas encore de prévision</ev-badge>
}
</li>
}
</ul>
</section>
}
</div> </div>
@@ -121,3 +121,40 @@ h2 {
.alert-item__message { .alert-item__message {
font-size: 0.9rem; font-size: 0.9rem;
} }
.predictions-list {
list-style: none;
margin: 0;
padding: 0;
display: flex;
flex-direction: column;
gap: 0.5rem;
}
.prediction-item {
display: flex;
align-items: center;
justify-content: space-between;
gap: 0.75rem;
padding: 0.7rem 1rem;
border-radius: var(--radius-md);
background: var(--color-surface);
border: 1px solid var(--color-border-light);
}
.prediction-item__site {
font-size: 0.9rem;
font-weight: 600;
}
.prediction-item__value {
font-size: 0.9rem;
font-weight: 600;
}
.prediction-item__target {
margin-left: 0.35rem;
font-size: 0.8rem;
font-weight: 400;
color: var(--color-text-muted);
}
@@ -4,6 +4,7 @@ import { of, throwError } from 'rxjs';
import { Dashboard } from './dashboard'; import { Dashboard } from './dashboard';
import { StatsService } from '../../core/services/stats.service'; import { StatsService } from '../../core/services/stats.service';
import { AlertsService } from '../../core/services/alerts.service'; import { AlertsService } from '../../core/services/alerts.service';
import { PredictionsService } from '../../core/services/predictions.service';
import {AuthService} from '../../core/services/auth.service'; import {AuthService} from '../../core/services/auth.service';
import {Router, provideRouter} from '@angular/router'; import {Router, provideRouter} from '@angular/router';
@@ -17,18 +18,24 @@ vi.mock('chart.js', () => {
return { Chart: ChartMock, registerables: [] }; return { Chart: ChartMock, registerables: [] };
}); });
function predictionsMock(sites: unknown[] = []) {
return { getPredictions: vi.fn().mockReturnValue(of({ timestamp: '2026-09-18T09:00:00Z', sites })) };
}
describe('Dashboard', () => { describe('Dashboard', () => {
afterEach(() => vi.useRealTimers()); afterEach(() => vi.useRealTimers());
it('charge les stats et les alertes au démarrage', async () => { it('charge les stats, les alertes et les prévisions au démarrage', async () => {
const statsMock = { getSummary: vi.fn().mockReturnValue(of({ total_sites: 7, sites: [] })) }; const statsMock = { getSummary: vi.fn().mockReturnValue(of({ total_sites: 7, sites: [] })) };
const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([{ alert_id: 'A1' }])) }; const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([{ alert_id: 'A1' }])) };
const predictions = predictionsMock([{ site_id: 'SITE001', site_name: 'Test', prediction: null }]);
TestBed.configureTestingModule({ TestBed.configureTestingModule({
imports: [Dashboard], imports: [Dashboard],
providers: [ providers: [
{ provide: StatsService, useValue: statsMock }, { provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock }, { provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictions },
provideRouter([]), provideRouter([]),
], ],
}); });
@@ -42,8 +49,12 @@ describe('Dashboard', () => {
expect(statsMock.getSummary).toHaveBeenCalled(); expect(statsMock.getSummary).toHaveBeenCalled();
expect(alertsMock.getAlerts).toHaveBeenCalled(); expect(alertsMock.getAlerts).toHaveBeenCalled();
expect(predictions.getPredictions).toHaveBeenCalled();
expect(fixture.componentInstance.alerts().length).toBe(1); expect(fixture.componentInstance.alerts().length).toBe(1);
expect(fixture.componentInstance.error()).toBeNull(); expect(fixture.componentInstance.predictions().length).toBe(1);
expect(fixture.componentInstance.statsError()).toBeNull();
expect(fixture.componentInstance.alertsError()).toBeNull();
expect(fixture.componentInstance.predictionsError()).toBeNull();
}); });
it("signale l'indisponibilité puis repart au rafraîchissement suivant", () => { it("signale l'indisponibilité puis repart au rafraîchissement suivant", () => {
@@ -61,6 +72,7 @@ describe('Dashboard', () => {
providers: [ providers: [
{ provide: StatsService, useValue: statsMock }, { provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock }, { provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictionsMock() },
provideRouter([]), provideRouter([]),
], ],
}); });
@@ -70,13 +82,13 @@ describe('Dashboard', () => {
vi.advanceTimersByTime(1); vi.advanceTimersByTime(1);
expect(statsMock.getSummary).toHaveBeenCalledTimes(1); expect(statsMock.getSummary).toHaveBeenCalledTimes(1);
expect(fixture.componentInstance.error()).not.toBeNull(); expect(fixture.componentInstance.statsError()).not.toBeNull();
expect(fixture.componentInstance.stats()).toBeNull(); expect(fixture.componentInstance.stats()).toBeNull();
vi.advanceTimersByTime(10000); vi.advanceTimersByTime(10000);
expect(statsMock.getSummary).toHaveBeenCalledTimes(2); expect(statsMock.getSummary).toHaveBeenCalledTimes(2);
expect(fixture.componentInstance.stats()).not.toBeNull(); expect(fixture.componentInstance.stats()).not.toBeNull();
expect(fixture.componentInstance.error()).toBeNull(); expect(fixture.componentInstance.statsError()).toBeNull();
}); });
it("n'interrompt pas la page quand le chargement des alertes échoue", () => { it("n'interrompt pas la page quand le chargement des alertes échoue", () => {
@@ -88,6 +100,7 @@ describe('Dashboard', () => {
providers: [ providers: [
{ provide: StatsService, useValue: statsMock }, { provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock }, { provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictionsMock() },
provideRouter([]), provideRouter([]),
], ],
}); });
@@ -96,6 +109,62 @@ describe('Dashboard', () => {
fixture.detectChanges(); fixture.detectChanges();
expect(fixture.componentInstance.alerts().length).toBe(0); expect(fixture.componentInstance.alerts().length).toBe(0);
expect(fixture.componentInstance.alertsError()).not.toBeNull();
});
it("n'interrompt pas la page quand le chargement des prévisions échoue", () => {
const statsMock = { getSummary: vi.fn().mockReturnValue(of({ total_sites: 7, sites: [] })) };
const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([])) };
const predictions = {
getPredictions: vi.fn().mockReturnValue(throwError(() => new Error('nope'))),
};
TestBed.configureTestingModule({
imports: [Dashboard],
providers: [
{ provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictions },
provideRouter([]),
],
});
const fixture = TestBed.createComponent(Dashboard);
fixture.detectChanges();
expect(fixture.componentInstance.predictions().length).toBe(0);
expect(fixture.componentInstance.predictionsError()).not.toBeNull();
});
it("un rafraîchissement de stats n'efface pas une erreur de prévisions en attente", () => {
vi.useFakeTimers();
const statsMock = { getSummary: vi.fn().mockReturnValue(of({ total_sites: 7, sites: [] })) };
const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([])) };
const predictions = {
getPredictions: vi.fn().mockReturnValue(throwError(() => new Error('nope'))),
};
TestBed.configureTestingModule({
imports: [Dashboard],
providers: [
{ provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictions },
provideRouter([]),
],
});
const fixture = TestBed.createComponent(Dashboard);
fixture.detectChanges();
expect(fixture.componentInstance.predictionsError()).not.toBeNull();
// Plusieurs cycles de `timer(0, 10_000)` (stats) plus tard, l'erreur des prévisions doit
// toujours être visible : rien ne vient la rafraîchir tant que la section n'est pas rechargée.
vi.advanceTimersByTime(30000);
expect(fixture.componentInstance.predictionsError()).not.toBeNull();
expect(fixture.componentInstance.statsError()).toBeNull();
}); });
it('appelle logout et redirige vers /login au clic sur le bouton de déconnexion', () => { it('appelle logout et redirige vers /login au clic sur le bouton de déconnexion', () => {
@@ -108,6 +177,7 @@ describe('Dashboard', () => {
providers: [ providers: [
{ provide: StatsService, useValue: statsMock }, { provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock }, { provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictionsMock() },
{ provide: AuthService, useValue: authMock }, { provide: AuthService, useValue: authMock },
provideRouter([]), provideRouter([]),
], ],
@@ -137,6 +207,7 @@ describe('Dashboard', () => {
providers: [ providers: [
{ provide: StatsService, useValue: statsMock }, { provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock }, { provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictionsMock() },
{ provide: AuthService, useValue: authMock }, { provide: AuthService, useValue: authMock },
provideRouter([]), provideRouter([]),
], ],
@@ -164,6 +235,7 @@ describe('Dashboard', () => {
providers: [ providers: [
{ provide: StatsService, useValue: statsMock }, { provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock }, { provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictionsMock() },
provideRouter([]), provideRouter([]),
], ],
}); });
@@ -179,4 +251,26 @@ describe('Dashboard', () => {
dashboard.badgeToneForSeverity('critical'), dashboard.badgeToneForSeverity('critical'),
); );
}); });
it('distingue le ton des statuts de prévision', () => {
const statsMock = { getSummary: vi.fn().mockReturnValue(of({ total_sites: 7, sites: [] })) };
const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([])) };
TestBed.configureTestingModule({
imports: [Dashboard],
providers: [
{ provide: StatsService, useValue: statsMock },
{ provide: AlertsService, useValue: alertsMock },
{ provide: PredictionsService, useValue: predictionsMock() },
provideRouter([]),
],
});
const fixture = TestBed.createComponent(Dashboard);
const dashboard = fixture.componentInstance;
expect(dashboard.badgeToneForPredictionStatus('available')).toBe('success');
expect(dashboard.badgeToneForPredictionStatus('insufficient_data')).toBe('warning');
expect(dashboard.badgeToneForPredictionStatus('error')).toBe('danger');
});
}); });
@@ -1,15 +1,17 @@
import { Component, OnInit, inject, signal, DestroyRef } from '@angular/core'; import { Component, OnInit, inject, signal, DestroyRef, WritableSignal } from '@angular/core';
import { takeUntilDestroyed } from '@angular/core/rxjs-interop'; import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
import { timer, switchMap, catchError, EMPTY, Observable } from 'rxjs'; import { timer, switchMap, catchError, EMPTY, Observable } from 'rxjs';
import { DecimalPipe } from '@angular/common'; import { DecimalPipe, DatePipe } from '@angular/common';
import { Router, RouterLink } from '@angular/router'; import { Router, RouterLink } from '@angular/router';
import { StatsService } from '../../core/services/stats.service'; import { StatsService } from '../../core/services/stats.service';
import { ConsumptionGauge } from '../../shared/components/consumption-gauge/consumption-gauge'; import { ConsumptionGauge } from '../../shared/components/consumption-gauge/consumption-gauge';
import { SiteLoadChart } from '../../shared/components/site-load-chart/site-load-chart'; import { SiteLoadChart } from '../../shared/components/site-load-chart/site-load-chart';
import { AlertsService } from '../../core/services/alerts.service'; import { AlertsService } from '../../core/services/alerts.service';
import { PredictionsService } from '../../core/services/predictions.service';
import { AuthService } from '../../core/services/auth.service'; import { AuthService } from '../../core/services/auth.service';
import { StatsSummary } from '../../shared/models/stats.model'; import { StatsSummary } from '../../shared/models/stats.model';
import { Alert, AlertSeverity } from '../../shared/models/alert.model'; import { Alert, AlertSeverity } from '../../shared/models/alert.model';
import { PredictionStatus, SitePredictionSummary } from '../../shared/models/prediction.model';
import { Card } from '../../shared/components/ui/card/card'; import { Card } from '../../shared/components/ui/card/card';
import { Alert as EvAlert } from '../../shared/components/ui/alert/alert'; import { Alert as EvAlert } from '../../shared/components/ui/alert/alert';
import { Badge, BadgeTone } from '../../shared/components/ui/badge/badge'; import { Badge, BadgeTone } from '../../shared/components/ui/badge/badge';
@@ -27,11 +29,21 @@ const TON_PAR_SEVERITE: Record<AlertSeverity, BadgeTone> = {
critical: 'critical', critical: 'critical',
}; };
// `error` n'a pas de précédent dans les fixtures ou l'API à ce jour, mais figure dans le
// domaine du schéma backend (`ck_prediction_status`) : mieux vaut une couleur définie que
// tomber sur `undefined` si ce statut apparaît un jour.
const TON_PAR_STATUT_PREDICTION: Record<PredictionStatus, BadgeTone> = {
available: 'success',
insufficient_data: 'warning',
error: 'danger',
};
@Component({ @Component({
selector: 'app-dashboard', selector: 'app-dashboard',
standalone: true, standalone: true,
imports: [ imports: [
DecimalPipe, DecimalPipe,
DatePipe,
RouterLink, RouterLink,
ConsumptionGauge, ConsumptionGauge,
SiteLoadChart, SiteLoadChart,
@@ -47,31 +59,52 @@ const TON_PAR_SEVERITE: Record<AlertSeverity, BadgeTone> = {
export class Dashboard implements OnInit { export class Dashboard implements OnInit {
private statsService = inject(StatsService); private statsService = inject(StatsService);
private alertsService = inject(AlertsService); private alertsService = inject(AlertsService);
private predictionsService = inject(PredictionsService);
private auth = inject(AuthService); private auth = inject(AuthService);
private router = inject(Router); private router = inject(Router);
private destroyRef = inject(DestroyRef); private destroyRef = inject(DestroyRef);
stats = signal<StatsSummary | null>(null); stats = signal<StatsSummary | null>(null);
alerts = signal<Alert[]>([]); alerts = signal<Alert[]>([]);
error = signal<string | null>(null); predictions = signal<SitePredictionSummary[]>([]);
// Un signal par flux, pas un seul `error` partagé : sinon le tick suivant de `timer` (stats)
// efface silencieusement un message d'échec des prévisions ou des alertes après 10s au plus,
// sans retry ni indication pour l'utilisateur que la section correspondante est restée vide.
statsError = signal<string | null>(null);
alertsError = signal<string | null>(null);
predictionsError = signal<string | null>(null);
ngOnInit(): void { ngOnInit(): void {
this.alertsService this.alertsService
.getAlerts() .getAlerts()
.pipe(catchError(() => this.reportUnavailable())) .pipe(catchError(() => this.reportUnavailable(this.alertsError)))
.subscribe((alerts) => this.alerts.set(alerts)); .subscribe((alerts) => {
this.alertsError.set(null);
this.alerts.set(alerts);
});
// Les prévisions viennent d'un scoring hors ligne, pas d'un calcul à la demande : un seul
// chargement au démarrage suffit, pas besoin du rafraîchissement périodique de `stats`.
this.predictionsService
.getPredictions()
.pipe(catchError(() => this.reportUnavailable(this.predictionsError)))
.subscribe((summary) => {
this.predictionsError.set(null);
this.predictions.set(summary.sites);
});
// Piège : le catchError porte sur l'observable interne. Sur le flux externe il // Piège : le catchError porte sur l'observable interne. Sur le flux externe il
// terminerait le timer, et le rafraîchissement ne repartirait jamais. // terminerait le timer, et le rafraîchissement ne repartirait jamais.
timer(0, REFRESH_INTERVAL_MS) timer(0, REFRESH_INTERVAL_MS)
.pipe( .pipe(
switchMap(() => switchMap(() =>
this.statsService.getSummary().pipe(catchError(() => this.reportUnavailable())), this.statsService.getSummary().pipe(catchError(() => this.reportUnavailable(this.statsError))),
), ),
takeUntilDestroyed(this.destroyRef), takeUntilDestroyed(this.destroyRef),
) )
.subscribe((stats) => { .subscribe((stats) => {
this.error.set(null); this.statsError.set(null);
this.stats.set(stats); this.stats.set(stats);
}); });
} }
@@ -80,6 +113,10 @@ export class Dashboard implements OnInit {
return TON_PAR_SEVERITE[severity]; return TON_PAR_SEVERITE[severity];
} }
badgeToneForPredictionStatus(status: PredictionStatus): BadgeTone {
return TON_PAR_STATUT_PREDICTION[status];
}
onLogout(): void { onLogout(): void {
this.auth.logout().subscribe({ this.auth.logout().subscribe({
next: () => this.router.navigate(['/login']), next: () => this.router.navigate(['/login']),
@@ -91,8 +128,8 @@ export class Dashboard implements OnInit {
}); });
} }
private reportUnavailable(): Observable<never> { private reportUnavailable(target: WritableSignal<string | null>): Observable<never> {
this.error.set(UNAVAILABLE_MESSAGE); target.set(UNAVAILABLE_MESSAGE);
return EMPTY; return EMPTY;
} }
} }
@@ -0,0 +1,24 @@
export type PredictionStatus = 'available' | 'insufficient_data' | 'error';
export type PredictionTargetMetric = 'consumption_kwh' | 'consumption_kw';
export interface SitePrediction {
target_at: string;
target_metric: PredictionTargetMetric;
period_minutes: number | null;
predicted_value: number | null;
status: PredictionStatus;
failure_reason: string | null;
model_reference: string;
created_at: string;
}
export interface SitePredictionSummary {
site_id: string;
site_name: string;
prediction: SitePrediction | null;
}
export interface PredictionSummary {
timestamp: string;
sites: SitePredictionSummary[];
}
+3 -3
View File
@@ -74,10 +74,10 @@ collecteur ne vient le lire.
| Domaine | Technologie | Emplacement | Statut | Ce qui existe réellement | | Domaine | Technologie | Emplacement | Statut | Ce qui existe réellement |
|---|---|---|---|---| |---|---|---|---|---|
| Backend | FastAPI, Python 3.14 | `apps/backend` | `En cours` | Factory, configuration, journalisation, 2 sondes de santé, `/metrics`, contrat OpenAPI versionné, routes `sites`, `alerts`, `recommendations`, `stats/summary` et `readings` en lecture (endpoints → services → repositories → models) | | Backend | FastAPI, Python 3.14 | `apps/backend` | `En cours` | Factory, configuration, journalisation, 2 sondes de santé, `/metrics`, contrat OpenAPI versionné, routes `sites`, `alerts`, `recommendations`, `stats/summary`, `readings`, `sensors/status` et `predictions` en lecture (endpoints → services → repositories → models) |
| Frontend | Angular 22, Node 24 | `apps/frontend` | `En cours` | Tableau de bord sur route `/dashboard`, deux services HTTP, graphiques Chart.js, données servies par des fixtures | | Frontend | Angular 22, Node 24 | `apps/frontend` | `En cours` | Tableau de bord sur route `/dashboard`, authentification complète (garde de route, intercepteur de jeton), cinq services HTTP, graphiques Chart.js. `stats`/`alerts` sur fixtures, `predictions` branché sur l'API réelle |
| Base | PostgreSQL 17 + TimescaleDB | `db` | `Fait` | Bootstrap de l'extension, base de test, chaîne Alembic. Schéma applicatif créé (`site`, `dataset`, `reading` en hypertable, `prediction`, `alert`, `recommendation`) | | Base | PostgreSQL 17 + TimescaleDB | `db` | `Fait` | Bootstrap de l'extension, base de test, chaîne Alembic. Schéma applicatif créé (`site`, `dataset`, `reading` en hypertable, `prediction`, `alert`, `recommendation`) |
| ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement (features par lags/moyennes glissantes, baseline de persistance saisonnière, suivi MLflow local), voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../../ML-START.md). Scoring, endpoint et surveillance de dérive pas encore construits | | ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`. Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../../ML-START.md). Automatisation (Airflow) et surveillance de dérive (EC06, #44/#45) pas encore construites |
| Infra | Terraform, k3s single-node | `infra/terraform` | `En cours` | Module d'installation du cluster. Jamais appliqué, aucune ressource Kubernetes déclarée | | Infra | Terraform, k3s single-node | `infra/terraform` | `En cours` | Module d'installation du cluster. Jamais appliqué, aucune ressource Kubernetes déclarée |
| Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API | | Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API |
| ETL | Apache Airflow | `etl/airflow` | `Cible` | Rien | | ETL | Apache Airflow | `etl/airflow` | `Cible` | Rien |
+58 -5
View File
@@ -12,10 +12,10 @@ Les quatre couches existent désormais, portées par l'authentification.
```mermaid ```mermaid
flowchart TB flowchart TB
ep["endpoints<br/>health, auth, users, sites, alerts,<br/>recommendations, stats, sensors"] ep["endpoints<br/>health, auth, users, sites, alerts,<br/>recommendations, stats, readings, sensors, predictions"]
sc["schemas<br/>Pydantic"] sc["schemas<br/>Pydantic"]
sv["services<br/>AuthService, UserService,<br/>SiteService, AlertService, RecommendationService,<br/>StatsService, SensorService"] sv["services<br/>AuthService, UserService,<br/>SiteService, AlertService, RecommendationService,<br/>StatsService, ReadingService, SensorService, PredictionService"]
rp["repositories<br/>user, refresh_token,<br/>login_attempt, audit_log,<br/>site, alert, recommendation, reading"] rp["repositories<br/>user, refresh_token,<br/>login_attempt, audit_log,<br/>site, alert, recommendation, reading, prediction"]
md["models<br/>10 tables"] md["models<br/>10 tables"]
db[("PostgreSQL")] db[("PostgreSQL")]
@@ -149,6 +149,7 @@ Deux fichiers d'environnement, deux usages : `.env` à la racine alimente `docke
| GET | `/api/v1/stats/summary` | Résume la consommation instantanée du parc. `lecteur` | 401, 403, 500 | | GET | `/api/v1/stats/summary` | Résume la consommation instantanée du parc. `lecteur` | 401, 403, 500 |
| GET | `/api/v1/readings` | Historique des lectures, filtrable par `site_id`, fenêtre `start`/`end` (24h par défaut, 90 jours maximum) et paginé par `limit`/`offset`. `lecteur` | 400, 401, 403, 422, 500 | | GET | `/api/v1/readings` | Historique des lectures, filtrable par `site_id`, fenêtre `start`/`end` (24h par défaut, 90 jours maximum) et paginé par `limit`/`offset`. `lecteur` | 400, 401, 403, 422, 500 |
| GET | `/api/v1/sensors/status` | État de santé des capteurs par site, dérivé de la dernière lecture. `admin` | 401, 403, 500 | | GET | `/api/v1/sensors/status` | État de santé des capteurs par site, dérivé de la dernière lecture. `admin` | 401, 403, 500 |
| GET | `/api/v1/predictions` | Dernière prévision de consommation par site, calculée hors ligne par le pipeline de scoring (`ml/`). `lecteur` | 401, 403, 500 |
| GET | `/metrics` | Format Prometheus, hors du schéma. Jeton requis si `APP_METRICS_TOKEN` est posé | | | GET | `/metrics` | Format Prometheus, hors du schéma. Jeton requis si `APP_METRICS_TOKEN` est posé | |
| GET | `/docs`, `/redoc`, `/openapi.json` | Hors du schéma. Fermés en `staging` et en `prod` | | | GET | `/docs`, `/redoc`, `/openapi.json` | Hors du schéma. Fermés en `staging` et en `prod` | |
@@ -163,7 +164,7 @@ le jeton à usage unique plutôt que dans un `Principal`.
donc de modifier `ROUTES_PUBLIQUES` dans `tests/api/acces.py`. donc de modifier `ROUTES_PUBLIQUES` dans `tests/api/acces.py`.
`GET /sites` et `GET /sites/{site_id}` sont la première route métier, et le gabarit repris pour `GET /sites` et `GET /sites/{site_id}` sont la première route métier, et le gabarit repris pour
`GET /alerts` puis pour les suivantes (`dataset`, `prediction`) : les quatre couches `GET /alerts` puis pour les suivantes (`dataset`) : les quatre couches
`endpoints → services → repositories → models` y sont toutes présentes, sur des tables déjà créées `endpoints → services → repositories → models` y sont toutes présentes, sur des tables déjà créées
par la révision Alembic `e6d2026091501`. Elles n'exigent que le rôle `lecteur`, contrairement aux par la révision Alembic `e6d2026091501`. Elles n'exigent que le rôle `lecteur`, contrairement aux
routes d'administration qui exigent `admin`. `SiteRepository` lit par `AsyncSession.scalar()` (une routes d'administration qui exigent `admin`. `SiteRepository` lit par `AsyncSession.scalar()` (une
@@ -181,6 +182,18 @@ dernière `Reading` du site : un site connu sans lecture rend `200` avec tous le
détaillé pour le frontend est dans détaillé pour le frontend est dans
[31-contrat-authentification.md](31-contrat-authentification.md). [31-contrat-authentification.md](31-contrat-authentification.md).
`GET /predictions` reprend ce même sous-gabarit « dernière valeur par site » (`SiteRepository` +
`PredictionRepository`, un `SitePredictionSummaryResponse` par site plutôt qu'une table brute).
Différence avec `stats`/`sensors` : `prediction` est une vraie table accumulée par un processus
externe (`enervision_ml.score`, cf. `ml/README.md`), pas une valeur recalculée à la volée depuis
`reading` à chaque appel. `PredictionRepository.latest_by_site()` isole donc un `DISTINCT ON
(site_id)` ordonné par `target_at DESC` (couvert par l'index `ix_prediction_site_target`), le même
mécanisme que `ReadingRepository.latest_by_site()`. Un site jamais scoré rend `prediction: null`
plutôt qu'un statut inventé : le domaine `available`/`insufficient_data`/`error` de la contrainte
`ck_prediction_status` n'a pas de valeur pour « pas encore de ligne ». L'API ne lance jamais
LightGBM elle-même ; elle lit ce que le pipeline de scoring a déjà écrit, cf.
[ML-START.md](../../ML-START.md) section 3.
`GET /readings` reprend le même gabarit mais s'en écarte sur un point : `reading` est l'hypertable, `GET /readings` reprend le même gabarit mais s'en écarte sur un point : `reading` est l'hypertable,
donc la seule table métier pouvant porter des années d'historique, ce que `docs/architecture/ donc la seule table métier pouvant porter des années d'historique, ce que `docs/architecture/
owasp-traceabilite.md` documentait comme un risque ouvert (API4, aucune pagination plafonnée ni owasp-traceabilite.md` documentait comme un risque ouvert (API4, aucune pagination plafonnée ni
@@ -194,6 +207,46 @@ 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 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. 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, 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` |
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.
**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
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` ### `/health/ready`
Cette sonde porte une garde décrite dans l'[ADR 0001](../adr/0001-postgresql-timescaledb.md) : un Cette sonde porte une garde décrite dans l'[ADR 0001](../adr/0001-postgresql-timescaledb.md) : un
@@ -269,7 +322,7 @@ Les modèles de `app/schemas/errors.py` décrivent ce que les gestionnaires renv
### Ajouter une route métier ### Ajouter une route métier
Checklist pour toute nouvelle route sur le gabarit `sites`/`alerts`/`recommendations`/`stats`/ Checklist pour toute nouvelle route sur le gabarit `sites`/`alerts`/`recommendations`/`stats`/
`readings`/`sensors` (`dataset`, `prediction`) : `readings`/`sensors`/`predictions` (`dataset`) :
1. Composer ses `responses=` depuis `app/api/openapi.py` : `REPONSES_LECTEUR`/`REPONSES_ADMIN` 1. Composer ses `responses=` depuis `app/api/openapi.py` : `REPONSES_LECTEUR`/`REPONSES_ADMIN`
au niveau de l'`include_router()` du routeur, `REPONSE_VALIDATION` et les codes locaux au niveau de l'`include_router()` du routeur, `REPONSE_VALIDATION` et les codes locaux
+25 -19
View File
@@ -13,24 +13,29 @@ Ce qui est en place :
- `app.config.ts` fournit `provideBrowserGlobalErrorListeners()`, `provideRouter(routes)` et - `app.config.ts` fournit `provideBrowserGlobalErrorListeners()`, `provideRouter(routes)` et
`provideHttpClient(withInterceptors([mockApiInterceptor]))`. `provideHttpClient(withInterceptors([mockApiInterceptor]))`.
- Une route `/dashboard` en composant différé, et une redirection depuis la racine. - Une route `/dashboard` en composant différé, et une redirection depuis la racine.
- `core/services` porte `StatsService` et `AlertsService`, `core/interceptors` l'intercepteur de - `core/services` porte `StatsService`, `AlertsService`, `PredictionsService`, `SitesService` et
fixtures, `features/dashboard` la page, `shared/components` la jauge de consommation et le `AuthService`, `core/interceptors` l'intercepteur de fixtures et l'intercepteur d'authentification
(jeton porteur, rafraîchissement sur 401), `core/guards` la garde de route `authGuard`,
`features/dashboard` la page principale, `shared/components` la jauge de consommation et le
graphique de charge par site, tous deux construits sur Chart.js. graphique de charge par site, tous deux construits sur Chart.js.
- Une authentification complète côté interface : connexion, mot de passe oublié/réinitialisation,
changement de mot de passe, garde de route sur `/dashboard` et `/sites`. Détail :
[31-contrat-authentification.md](31-contrat-authentification.md).
- Un système de design partagé (`shared/components/ui/` : `ev-button`, `ev-card`, `ev-alert`, - Un système de design partagé (`shared/components/ui/` : `ev-button`, `ev-card`, `ev-alert`,
`ev-badge`, `ev-brand`, tokens CSS dans `styles/_tokens.scss`) que toute nouvelle page doit `ev-badge`, `ev-brand`, tokens CSS dans `styles/_tokens.scss`) que toute nouvelle page doit
réutiliser plutôt que redéfinir ses propres styles. Détail : réutiliser plutôt que redéfinir ses propres styles. Détail :
[32-design-systeme-frontend.md](32-design-systeme-frontend.md). [32-design-systeme-frontend.md](32-design-systeme-frontend.md).
- L'état vit dans des signaux, sans bibliothèque dédiée. - L'état vit dans des signaux, sans bibliothèque dédiée.
- Vitest via le builder `@angular/build:unit-test`, couverture activée, sept fichiers de test. - Vitest via le builder `@angular/build:unit-test`, couverture activée.
- Prettier configuré, parser `angular` pour les gabarits HTML. - Prettier configuré, parser `angular` pour les gabarits HTML.
Ce qui n'existe pas encore : Ce qui n'existe pas encore :
- **Aucun endpoint réel derrière l'écran.** `GET /api/v1/stats/summary` et `GET /api/v1/alerts` - **`stats`/`alerts` restent sur fixtures.** `GET /api/v1/stats/summary` et `GET /api/v1/alerts`
sont servis par l'intercepteur ; l'API expose `/health`, `/auth` et `/users`, rien d'autre. sont servis par l'intercepteur de fixtures ; l'API expose bien ces routes désormais, mais rien
- Aucune authentification côté interface : ni garde de route, ni intercepteur de jeton, alors que ne bascule `useMockFixtures` à `false` en développement pour les consommer réellement.
les routes métier de l'API en exigent un. Voir `GET /api/v1/predictions` fait exception : jamais mocké, branché sur l'API réelle depuis cette
[31-contrat-authentification.md](31-contrat-authentification.md). PR (voir plus bas).
- Aucun état de chargement : tant que la première réponse n'est pas arrivée, la page reste vide. - Aucun état de chargement : tant que la première réponse n'est pas arrivée, la page reste vide.
- Aucun lint : ESLint n'est pas installé. - Aucun lint : ESLint n'est pas installé.
@@ -84,19 +89,19 @@ sequenceDiagram
`mockApiInterceptor` n'intercepte que `/stats/summary` et `/alerts`, et seulement si `mockApiInterceptor` n'intercepte que `/stats/summary` et `/alerts`, et seulement si
`environment.useMockFixtures` est vrai. Le drapeau est à `true` en développement, à `false` en `environment.useMockFixtures` est vrai. Le drapeau est à `true` en développement, à `false` en
production : toute autre requête, et toutes les requêtes en production, suivent le chemin réel. production : toute autre requête, et toutes les requêtes en production, suivent le chemin réel.
`/predictions` est volontairement exclu de cette liste (contrairement à `stats`/`alerts`) : il
suit toujours le chemin réel, comme `/auth/*` - en développement, ça veut dire qu'un jeton valide
et un backend joignable sont nécessaires pour que la section prévisions du dashboard s'affiche.
En développement, `proxy.conf.json` redirige tout `/api` vers `http://localhost:8000`. C'est ce En développement, `proxy.conf.json` redirige tout `/api` vers `http://localhost:8000`. C'est ce
qui évite le CORS sur le poste, et c'est pourquoi `environment.development.ts` se contente d'un qui évite le CORS sur le poste, et c'est pourquoi `environment.development.ts` se contente d'un
`apiUrl` relatif, `/api/v1`. `apiUrl` relatif, `/api/v1`.
En production, il n'y a pas de proxy : `environment.ts` porte une URL absolue. Angular substitue En production, il n'y a pas de proxy, mais `environment.ts` porte lui aussi un `apiUrl` relatif
le fichier via `fileReplacements`, et la configuration `production` est celle par défaut. (`/api/v1`) plutôt qu'une URL absolue : la dette qui pointait en dur sur
`http://localhost:8000/api/v1` a été corrigée. Un build de production sert donc l'appel `/api/v1/...`
**Dette connue.** `src/environments/environment.ts`, qui est la configuration de production, sur son propre origin, ce qui suppose qu'un ingress ou un reverse proxy route `/api` vers le
pointe `http://localhost:8000/api/v1` en dur. La valeur est celle du poste de développement : backend une fois déployé — question toujours ouverte dans [10-infra.md](10-infra.md).
telle quelle, un build de production ne joindra jamais l'API. À corriger avant le premier
déploiement, en même temps que sera tranchée la question de l'ingress dans
[10-infra.md](10-infra.md).
## Exécution ## Exécution
@@ -124,9 +129,10 @@ avec un service statique, il reste à écrire.
## Sécurité ## Sécurité
- Le frontend ne détient aucun secret : `environment.ts` ne porte qu'une URL. - Le frontend ne détient aucun secret : `environment.ts` ne porte qu'une URL.
- L'authentification existe côté API mais pas côté interface : aucune garde de route, aucun - L'authentification existe des deux côtés désormais : `authGuard` protège `/dashboard` et
intercepteur de jeton. `core/guards` reste à créer, `core/interceptors` n'héberge aujourd'hui `/sites`, `authInterceptor` pose le jeton porteur sur les requêtes sortantes et déclenche le
que les fixtures. rafraîchissement sur 401. Détail complet dans
[31-contrat-authentification.md](31-contrat-authentification.md).
## Tests ## Tests
+51 -5
View File
@@ -57,6 +57,46 @@ validation. La coupure est **chronologique**, jamais un tirage aleatoire de lign
aleatoire laisserait des lignes de validation "voir" des lignes d'entrainement via leurs aleatoire laisserait des lignes de validation "voir" des lignes d'entrainement via leurs
lags/moyennes glissantes, une fuite qui masquerait un surapprentissage. lags/moyennes glissantes, une fuite qui masquerait un surapprentissage.
## Scoring
```bash
uv run python -m enervision_ml.score --csv data/all_sites_combined.csv
# ou, une fois la base peuplee et ML_DATABASE_URL positionnee :
uv run python -m enervision_ml.score
```
Calcule, pour chaque site (ou un seul avec `--site-id`), la consommation prevue de l'heure suivant
sa derniere lecture connue, et ecrit une ligne dans `prediction`. Etapes, cf. `ML-START.md`
section 2 :
1. Lit une fenetre recente de `reading`+`site` (21 jours par defaut, une marge au-dessus des 168h
necessaires au lag hebdomadaire) plutot que tout l'historique -- le meme piege que celui deja
corrige sur `GET /readings` (fenetre non plafonnee sur une hypertable).
2. Ajoute une ligne "future" par site (l'heure suivante) et calcule ses features avec
`enervision_ml.features.build_features`, **exactement** la meme fonction qu'a l'entrainement.
3. Si le lag de 168h est absent (moins d'une semaine d'historique pour ce site) : ecrit
`status="insufficient_data"` directement, sans jamais appeler LightGBM.
4. Sinon : appelle `booster.predict(...)` et ecrit `status="available"` avec la valeur predite.
`--model` pointe vers le fichier entraine (`models/lightgbm-consumption.txt` par defaut).
`model_reference` en base est le hache SHA-256 (tronque) du fichier modele, pas son nom de
fichier : `train.py` reecrit toujours le meme chemin a chaque entrainement, donc le nom seul ne
distinguerait pas deux versions du modele.
En mode `--csv`, rien n'est ecrit en base : c'est un instantane historique fige (l'heure "future"
calculee a partir de la fin du CSV n'existe dans aucune base reelle), utile pour valider le
pipeline sans base joignable.
**Limite assumee** : la feature `is_working_hours` de la ligne future est recopiee depuis la
derniere lecture reelle, pas recalculee -- il n'existe aucune regle horaire ouvrable dans ce
depot (elle vit dans le generateur du jeu de donnees d'origine). L'approximation n'est fausse
qu'aux heures de bascule ouverture/fermeture, sur une seule feature parmi une dizaine, pour une
prevision a un seul pas.
`prediction` n'a pas de contrainte d'unicite sur `(site_id, target_at)` : chaque run de scoring
insere une nouvelle ligne plutot que d'ecraser la precedente, pour garder une trace de chaque
prevision (utile plus tard pour comparer prevision et realise, surveillance de derive #44/#45).
## Commandes ## Commandes
```bash ```bash
@@ -81,8 +121,14 @@ environnement de developpement pour le moment.
## Piege a connaitre ## Piege a connaitre
`enervision_ml.features.build_features` est **le seul endroit** qui doit construire les features `enervision_ml.features.build_features` est **le seul endroit** qui doit construire les features
du modele, a l'entrainement comme au futur scoring (service #37, pas encore construit). Si les du modele, a l'entrainement comme au scoring (`enervision_ml.score`). Si les deux divergent meme
deux divergent meme legerement (une fenetre de moyenne glissante calculee differemment, par legerement (une fenetre de moyenne glissante calculee differemment, par exemple), le modele
exemple), le modele recoit en production des features qui ne ressemblent plus a ce qu'il a recoit en production des features qui ne ressemblent plus a ce qu'il a appris, et ses predictions
appris, et ses predictions deviennent silencieusement mauvaises sans qu'aucune erreur ne se deviennent silencieusement mauvaises sans qu'aucune erreur ne se declenche. Ne jamais reecrire
declenche. Ne jamais reecrire cette logique ailleurs : importer `enervision_ml.features`. cette logique ailleurs : importer `enervision_ml.features`.
## Et cote API ?
`GET /api/v1/predictions` (backend, `apps/backend`) lit ce que `enervision_ml.score` a ecrit dans
`prediction` -- la derniere prevision par site, jamais un recalcul a la volee. FastAPI ne fait
jamais tourner LightGBM lui-meme, cf. `ML-START.md` section 3.
+66 -3
View File
@@ -16,6 +16,7 @@ Deux chemins, qui doivent produire le meme schema de sortie (colonnes `site_id`,
colonne est renvoyee a `NaN`, que LightGBM gere nativement comme valeur manquante. colonne est renvoyee a `NaN`, que LightGBM gere nativement comme valeur manquante.
""" """
from datetime import datetime
from pathlib import Path from pathlib import Path
import pandas as pd import pandas as pd
@@ -34,6 +35,14 @@ OUTPUT_COLUMNS = [
"capacity_kw", "capacity_kw",
] ]
NUMERIC_COLUMNS = [
"consumption_kwh",
"temperature_celsius",
"humidity_percent",
"solar_irradiance_wm2",
"capacity_kw",
]
_READING_QUERY = text( _READING_QUERY = text(
""" """
SELECT SELECT
@@ -53,10 +62,43 @@ _READING_QUERY = text(
) )
_RECENT_READING_QUERY = text(
"""
SELECT
r.site_id,
r.timestamp,
r.consumption_kwh,
r.temperature_celsius,
r.humidity_percent,
r.solar_irradiance_wm2,
r.is_working_hours,
s.site_type,
s.capacity_kw
FROM reading r
JOIN site s ON s.site_id = r.site_id
WHERE r.timestamp >= :since
ORDER BY r.site_id, r.timestamp
"""
)
def load_from_database(connection: Connectable) -> pd.DataFrame: def load_from_database(connection: Connectable) -> pd.DataFrame:
"""Lit l'historique complet `reading` + `site` depuis PostgreSQL.""" """Lit l'historique complet `reading` + `site` depuis PostgreSQL. Entrainement seulement :
le scoring n'a besoin que d'une fenetre recente, cf. `load_recent_from_database`.
"""
frame = pd.read_sql(_READING_QUERY, connection) frame = pd.read_sql(_READING_QUERY, connection)
return frame[OUTPUT_COLUMNS] return _typer(frame[OUTPUT_COLUMNS])
def load_recent_from_database(connection: Connectable, *, since: datetime) -> pd.DataFrame:
"""Lit `reading` + `site` depuis `since` seulement, pour le scoring.
Piege evite : un `SELECT` sans borne sur l'hypertable complete juste pour scorer le prochain
pas horaire serait la meme erreur que celle corrigee sur `GET /readings` (fenetre non
plafonnee sur une table pouvant porter des annees d'historique).
"""
frame = pd.read_sql(_RECENT_READING_QUERY, connection, params={"since": since})
return _typer(frame[OUTPUT_COLUMNS])
def load_from_csv(csv_path: Path) -> pd.DataFrame: def load_from_csv(csv_path: Path) -> pd.DataFrame:
@@ -65,4 +107,25 @@ def load_from_csv(csv_path: Path) -> pd.DataFrame:
frame["capacity_kw"] = float("nan") frame["capacity_kw"] = float("nan")
frame["is_working_hours"] = frame["is_working_hours"].astype(bool) frame["is_working_hours"] = frame["is_working_hours"].astype(bool)
return frame[OUTPUT_COLUMNS] return _typer(frame[OUTPUT_COLUMNS])
def _typer(frame: pd.DataFrame) -> pd.DataFrame:
"""Force le typage numerique attendu par LightGBM.
Piege reel, pas theorique : `site.capacity_kw` n'est peuple par aucun pipeline d'ingestion
aujourd'hui (`historical_import.py` ne pose que `site_type`/`site_name`). Une colonne
entierement `NULL` revient de `pd.read_sql` en dtype `object` plutot que `float64`, ce que
LightGBM refuse ("pandas dtypes must be int, float or bool"). `pd.to_numeric` corrige aussi
n'importe quelle autre colonne mesuree entierement absente sur une fenetre de scoring, pas
seulement `capacity_kw`.
Piege additionnel : `NUMERIC_COLUMNS` inclut `consumption_kwh`, la cible du modele, pas
seulement des variables explicatives. Une valeur non numerique y devient donc silencieusement
`NaN` aussi bien a l'entrainement (ou `train.py` l'exclura ensuite via son `dropna`) qu'au
scoring -- ce n'est pas un effet de bord limite aux colonnes mesurees.
"""
typee = frame.copy()
for colonne in NUMERIC_COLUMNS:
typee[colonne] = pd.to_numeric(typee[colonne], errors="coerce")
return typee
+318
View File
@@ -0,0 +1,318 @@
"""Scoring du modele LightGBM : calcule et enregistre la consommation prevue du prochain pas
horaire, par site.
CLI autonome, sur le meme gabarit que `enervision_ml.train` et
`apps/backend/app/etl/historical_import.py`. Cf. `docs/ML-START.md`, section 2.
uv run python -m enervision_ml.score --csv ../ml/data/all_sites_combined.csv
uv run python -m enervision_ml.score # lit ML_DATABASE_URL, ecrit dans `prediction`
Reutilise `enervision_ml.features.build_features` tel quel (jamais reecrit) : c'est la garantie
contre le train/serve skew documentee dans ce module.
"""
import argparse
import hashlib
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Any, cast
import lightgbm as lgb
import pandas as pd
from sqlalchemy import create_engine, text
from sqlalchemy.engine import Connection
from enervision_ml import config
from enervision_ml.data import load_from_csv, load_recent_from_database
from enervision_ml.features import TARGET_COLUMN, WEATHER_COLUMNS, build_features, feature_columns
# Marge au-dessus des 168h necessaires au lag hebdomadaire, pour absorber les trous de mesure.
LOOKBACK = timedelta(days=21)
# Au-dela de ce seuil, la derniere lecture d'un site est trop vieille pour que "l'heure
# suivante" ait un sens operationnel : ce n'est plus une prevision a un pas, c'est un site dont
# l'ingestion s'est probablement arretee. Sans cette borne, `build_scoring_frame` produirait
# quand meme un `target_at` (derniere lecture + 1h), et rien en aval (ni l'API, ni le dashboard)
# ne distingue une prevision fraiche d'une prevision vieille de plusieurs jours.
MAX_STALENESS = timedelta(hours=24)
TARGET_METRIC = "consumption_kwh"
PERIOD_MINUTES = 60
LAG_168H_COLUMN = f"{TARGET_COLUMN}_lag_168h"
INSUFFICIENT_DATA_REASON = (
"Historique insuffisant : moins de 168h de consumption_kwh disponibles pour ce site."
)
def _stale_reason(age: pd.Timedelta) -> str:
return (
f"Dernière lecture vieille de {age.total_seconds() / 3600:.0f}h "
f"(seuil {MAX_STALENESS.total_seconds() / 3600:.0f}h) : ingestion probablement "
"arrêtée pour ce site."
)
@dataclass(frozen=True, slots=True)
class ScoredSite:
site_id: str
target_at: datetime
status: str
predicted_value: float | None
failure_reason: str | None
def model_reference(model_path: Path) -> str:
"""Identifiant stable du modele utilise, insensible au fait que `train.py` reecrive
toujours le meme nom de fichier a chaque entrainement (pas de versioning par nom, cf.
`ml/README.md`)."""
empreinte = hashlib.sha256(model_path.read_bytes()).hexdigest()
return f"lightgbm-{empreinte[:12]}"
def build_scoring_frame(recent: pd.DataFrame, *, site_id: str | None = None) -> pd.DataFrame:
"""Ajoute une ligne future (l'heure suivant la derniere lecture connue) par site, et calcule
ses features par `build_features` -- exactement comme a l'entrainement, seule la cible de
cette ligne est inconnue.
Piege assume : `is_working_hours` de la ligne future est copie de la derniere lecture reelle,
pas recalcule. Il n'existe aucune regle horaire ouvrable dans ce depot (elle vit dans le
generateur du jeu de donnees d'origine, hors de ce code) ; l'approximation n'est fausse
qu'aux heures de bascule (ouverture/fermeture), sur une seule feature parmi une dizaine, pour
une prevision a un pas seulement.
"""
travail = recent if site_id is None else recent[recent["site_id"] == site_id]
if travail.empty:
return build_features(travail)
dernieres = (
travail.sort_values("timestamp").groupby("site_id", as_index=False, sort=False).tail(1)
).copy()
dernieres["timestamp"] = dernieres["timestamp"] + pd.Timedelta(hours=1)
dernieres[TARGET_COLUMN] = float("nan")
# Meteo future inconnue (cf. piege documente dans `enervision_ml.features.build_features`) :
# laisser `NaN` ici n'a aucun effet sur les features utilisees, qui ne prennent la meteo que
# decalee.
for colonne in WEATHER_COLUMNS:
dernieres[colonne] = float("nan")
etendu = pd.concat([travail, dernieres], ignore_index=True)
features = build_features(etendu)
return features.groupby("site_id", as_index=False, sort=False).tail(1).reset_index(drop=True)
def score(
booster: lgb.Booster, scoring_frame: pd.DataFrame, *, instant: datetime
) -> list[ScoredSite]:
resultats: list[ScoredSite] = []
# `timestamp` de la ligne de scoring vaut derniere lecture + 1h (cf. `build_scoring_frame`) :
# on en deduit l'age de cette derniere lecture par rapport a `instant`.
travail = scoring_frame.copy()
travail["_age"] = instant - (travail["timestamp"] - pd.Timedelta(hours=1))
perimes = travail[travail["_age"] > MAX_STALENESS]
for enregistrement in _records(perimes):
resultats.append(
ScoredSite(
site_id=enregistrement["site_id"],
target_at=enregistrement["timestamp"].to_pydatetime(),
status="insufficient_data",
predicted_value=None,
failure_reason=_stale_reason(enregistrement["_age"]),
)
)
a_jour = travail[travail["_age"] <= MAX_STALENESS]
insuffisants = a_jour[a_jour[LAG_168H_COLUMN].isna()]
for enregistrement in _records(insuffisants):
resultats.append(
ScoredSite(
site_id=enregistrement["site_id"],
target_at=enregistrement["timestamp"].to_pydatetime(),
status="insufficient_data",
predicted_value=None,
failure_reason=INSUFFICIENT_DATA_REASON,
)
)
suffisants = a_jour[a_jour[LAG_168H_COLUMN].notna()]
if not suffisants.empty:
typee = suffisants.copy()
typee["site_type"] = typee["site_type"].astype("category")
predictions = booster.predict(typee[feature_columns()])
for enregistrement, valeur in zip(_records(suffisants), predictions, strict=True):
resultats.append(
ScoredSite(
site_id=enregistrement["site_id"],
target_at=enregistrement["timestamp"].to_pydatetime(),
status="available",
predicted_value=float(valeur),
failure_reason=None,
)
)
return resultats
def _records(frame: pd.DataFrame) -> list[dict[str, Any]]:
return cast(list[dict[str, Any]], frame.to_dict(orient="records"))
_INSERT_PREDICTION = text(
"""
INSERT INTO prediction (
site_id, target_at, target_metric, period_minutes,
predicted_value, model_reference, status, failure_reason
) VALUES (
:site_id, :target_at, :target_metric, :period_minutes,
:predicted_value, :model_reference, :status, :failure_reason
)
"""
)
def write_predictions(
connection: Connection, resultats: list[ScoredSite], *, reference: str
) -> None:
"""Ecrit une ligne par site score. Insertion seule, jamais de mise a jour : `prediction`
n'a pas de contrainte d'unicite sur `(site_id, target_at)`, chaque run garde sa propre trace
plutot que d'ecraser la precedente -- utile plus tard pour comparer prevision et realise
(surveillance de derive, #44/#45)."""
if not resultats:
return
lignes = [
{
"site_id": r.site_id,
"target_at": r.target_at,
"target_metric": TARGET_METRIC,
"period_minutes": PERIOD_MINUTES,
"predicted_value": r.predicted_value,
"model_reference": reference,
"status": r.status,
"failure_reason": r.failure_reason,
}
for r in resultats
]
connection.execute(_INSERT_PREDICTION, lignes)
def _load_recent_from_csv(csv_path: Path, *, now: datetime | None) -> tuple[pd.DataFrame, datetime]:
brute = load_from_csv(csv_path)
instant = now or (
brute["timestamp"].max().to_pydatetime() if not brute.empty else datetime.now(UTC)
)
return brute[brute["timestamp"] >= instant - LOOKBACK], instant
def _score_frame(
recent: pd.DataFrame, *, model_path: Path, site_id: str | None, instant: datetime
) -> list[ScoredSite]:
scoring_frame = build_scoring_frame(recent, site_id=site_id)
if scoring_frame.empty:
return []
booster = lgb.Booster(model_file=str(model_path))
return score(booster, scoring_frame, instant=instant)
def run_scoring(
*,
model_path: Path,
csv_path: Path | None = None,
site_id: str | None = None,
now: datetime | None = None,
) -> list[ScoredSite]:
"""Score le prochain pas horaire par site et l'ecrit dans `prediction`.
En mode `--csv`, rien n'est ecrit : c'est un instantane historique fige (l'heure "future"
calculee n'existe dans aucune base reelle), utile pour valider le pipeline sans base
joignable, cf. `ml/README.md`. `site_id` n'est filtre qu'une fois, dans
`build_scoring_frame` : le filtrer aussi ici serait redondant.
"""
if csv_path is not None:
recent, instant = _load_recent_from_csv(csv_path, now=now)
return _score_frame(recent, model_path=model_path, site_id=site_id, instant=instant)
# Un seul engine pour la lecture et l'ecriture de ce run, plutot qu'un par etape.
engine = create_engine(config.database_url())
try:
instant = now or datetime.now(UTC)
recent = load_recent_from_database(engine, since=instant - LOOKBACK)
resultats = _score_frame(recent, model_path=model_path, site_id=site_id, instant=instant)
reference = model_reference(model_path)
with engine.begin() as connection:
write_predictions(connection, resultats, reference=reference)
return resultats
finally:
engine.dispose()
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Scoring du modele LightGBM EnerVision")
parser.add_argument(
"--model",
type=Path,
default=Path("models/lightgbm-consumption.txt"),
help="Chemin du modele entraine. Defaut : models/lightgbm-consumption.txt.",
)
parser.add_argument(
"--csv",
type=Path,
default=None,
help=(
"Instantane historique de demarrage/demo, rien n'est ecrit en base. Omis, lit "
"ML_DATABASE_URL, se connecte a PostgreSQL et ecrit dans `prediction`."
),
)
parser.add_argument(
"--site-id",
default=None,
help="Ne score que ce site. Omis, tous les sites presents dans la fenetre recente.",
)
parser.add_argument(
"--now",
type=_parse_instant,
default=None,
help=(
"Instant de reference (ISO 8601), pour tester ou demontrer le scoring cote base sur "
"des donnees anciennes (ex. le jeu de donnees historique, qui s'arrete fin 2024). "
"Omis, horloge systeme reelle."
),
)
return parser.parse_args()
def _parse_instant(valeur: str) -> datetime:
instant = datetime.fromisoformat(valeur)
return instant if instant.tzinfo is not None else instant.replace(tzinfo=UTC)
def main() -> None:
args = parse_args()
resultats = run_scoring(
model_path=args.model, csv_path=args.csv, site_id=args.site_id, now=args.now
)
if not resultats:
print("Aucun site a scorer (aucune lecture recente dans la fenetre).")
return
for r in resultats:
if r.status == "available":
print(f"{r.site_id} @ {r.target_at} : {r.predicted_value:.2f} kWh")
else:
print(f"{r.site_id} @ {r.target_at} : {r.status} ({r.failure_reason})")
if args.csv is not None:
print("\nMode --csv : instantane historique, rien ecrit en base.")
if __name__ == "__main__":
main()
+3 -3
View File
@@ -47,9 +47,9 @@ select = [
"S", "S",
"PT", "PT",
] ]
# N806 : `X`/`y` (donnees/cible) est la convention scikit-learn/LightGBM, pas une variable mal # N806/N803 : `X`/`y` (donnees/cible) est la convention scikit-learn/LightGBM, pas une variable
# nommee. # ou un argument mal nomme.
ignore = ["B008", "N806"] ignore = ["B008", "N806", "N803"]
[tool.ruff.lint.per-file-ignores] [tool.ruff.lint.per-file-ignores]
"tests/**/*.py" = ["S101"] "tests/**/*.py" = ["S101"]
+55
View File
@@ -0,0 +1,55 @@
from pathlib import Path
import pandas as pd
from enervision_ml.data import NUMERIC_COLUMNS, load_from_csv
_CSV_HEADER = (
"site_id,timestamp,consumption_kwh,temperature_celsius,humidity_percent,"
"solar_irradiance_wm2,is_working_hours,site_type"
)
def write_csv(tmp_path: Path, *lignes: str) -> Path:
csv_path = tmp_path / "recent.csv"
csv_path.write_text("\n".join([_CSV_HEADER, *lignes]) + "\n")
return csv_path
def test_load_from_csv_types_every_numeric_column_as_float(tmp_path: Path) -> None:
csv_path = write_csv(tmp_path, "SITE001,2026-01-01T00:00:00,10.5,15.0,50.0,0.0,True,office")
frame = load_from_csv(csv_path)
for colonne in NUMERIC_COLUMNS:
assert frame[colonne].dtype == "float64"
def test_load_from_csv_coerces_a_corrupted_measurement_to_nan(tmp_path: Path) -> None:
# Reproduit une valeur de capteur corrompue plutot que vraiment manquante : `pandas` type
# alors la colonne entiere en `object`, pas en `float64` rempli de `NaN` -- le meme genre de
# divergence de typage que celle que `pd.read_sql` produit sur une colonne SQL entierement
# `NULL` (cf. `site.capacity_kw`, jamais peuplee par aucun pipeline d'ingestion aujourd'hui).
csv_path = write_csv(
tmp_path,
"SITE001,2026-01-01T00:00:00,10.5,15.0,50.0,0.0,True,office",
"SITE001,2026-01-01T01:00:00,capteur_hs,15.2,50.5,0.0,True,office",
)
frame = load_from_csv(csv_path)
assert frame["consumption_kwh"].dtype == "float64"
assert frame["consumption_kwh"].iloc[0] == 10.5
assert pd.isna(frame["consumption_kwh"].iloc[1])
def test_load_from_csv_always_types_capacity_kw_as_float(tmp_path: Path) -> None:
# `capacity_kw` n'existe pas dans ce CSV : `load_from_csv` la pose elle-meme a `NaN`. Cette
# affectation directe est deja un `float`, contrairement au cas `pd.read_sql` -- ce test
# garde le contrat visible malgre tout, au cas ou l'implementation changerait.
csv_path = write_csv(tmp_path, "SITE001,2026-01-01T00:00:00,10.5,15.0,50.0,0.0,True,office")
frame = load_from_csv(csv_path)
assert frame["capacity_kw"].dtype == "float64"
assert pd.isna(frame["capacity_kw"].iloc[0])
+277
View File
@@ -0,0 +1,277 @@
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Any
import pandas as pd
import pytest
from enervision_ml.features import TARGET_COLUMN
from enervision_ml.score import (
LAG_168H_COLUMN,
MAX_STALENESS,
ScoredSite,
build_scoring_frame,
model_reference,
run_scoring,
score,
write_predictions,
)
def make_recent(
site_id: str, *, heures: int, depart: datetime, valeur: float = 10.0
) -> pd.DataFrame:
instants = [depart + timedelta(hours=h) for h in range(heures)]
return pd.DataFrame(
{
"site_id": site_id,
"timestamp": instants,
TARGET_COLUMN: [valeur + h for h in range(heures)],
"temperature_celsius": 15.0,
"humidity_percent": 50.0,
"solar_irradiance_wm2": 0.0,
"is_working_hours": True,
"site_type": "office",
"capacity_kw": 100.0,
}
)
class FakeBooster:
def __init__(self, valeur: float = 42.0) -> None:
self.valeur = valeur
self.appels: list[int] = []
def predict(self, X: Any) -> list[float]:
self.appels.append(len(X))
return [self.valeur] * len(X)
class FakeConnection:
def __init__(self) -> None:
self.appels: list[tuple[Any, Any]] = []
def execute(self, statement: Any, parameters: Any = None) -> None:
self.appels.append((statement, parameters))
def test_build_scoring_frame_adds_one_row_per_site_one_hour_after_the_last_reading() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
recent = pd.concat(
[
make_recent("site-a", heures=200, depart=depart),
make_recent("site-b", heures=200, depart=depart),
],
ignore_index=True,
)
scoring_frame = build_scoring_frame(recent)
assert set(scoring_frame["site_id"]) == {"site-a", "site-b"}
derniere_lecture = depart + timedelta(hours=199)
assert (scoring_frame["timestamp"] == derniere_lecture + timedelta(hours=1)).all()
def test_build_scoring_frame_computes_lags_from_real_history() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
recent = make_recent("site-a", heures=200, depart=depart)
scoring_frame = build_scoring_frame(recent)
ligne = scoring_frame.iloc[0]
# La cible future n'existe pas : le lag d'1h doit valoir la toute derniere valeur reelle.
assert ligne[f"{TARGET_COLUMN}_lag_1h"] == recent[TARGET_COLUMN].iloc[-1]
def test_build_scoring_frame_flags_insufficient_history_under_168_hours() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
recent = make_recent("site-a", heures=100, depart=depart)
scoring_frame = build_scoring_frame(recent)
assert pd.isna(scoring_frame.iloc[0][LAG_168H_COLUMN])
def test_build_scoring_frame_accepts_a_full_week_of_history() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
recent = make_recent("site-a", heures=169, depart=depart)
scoring_frame = build_scoring_frame(recent)
assert not pd.isna(scoring_frame.iloc[0][LAG_168H_COLUMN])
def test_build_scoring_frame_filters_to_a_single_site() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
recent = pd.concat(
[
make_recent("site-a", heures=200, depart=depart),
make_recent("site-b", heures=200, depart=depart),
],
ignore_index=True,
)
scoring_frame = build_scoring_frame(recent, site_id="site-a")
assert scoring_frame["site_id"].tolist() == ["site-a"]
def test_build_scoring_frame_returns_empty_when_there_is_no_recent_reading() -> None:
recent = make_recent("site-a", heures=0, depart=datetime(2026, 1, 1, tzinfo=UTC))
scoring_frame = build_scoring_frame(recent)
assert scoring_frame.empty
def target_at_for(depart: datetime, heures: int) -> datetime:
"""`target_at` que produira `build_scoring_frame` pour ce jeu synthetique (derniere lecture
+ 1h) : l'utiliser comme `instant` donne un age d'1h, largement sous le seuil de peremption,
pour les tests qui ne visent pas ce filtre."""
return depart + timedelta(hours=heures)
def test_score_marks_insufficient_history_without_calling_the_model() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
scoring_frame = build_scoring_frame(make_recent("site-a", heures=100, depart=depart))
booster = FakeBooster()
resultats = score(
booster, # type: ignore[arg-type]
scoring_frame,
instant=target_at_for(depart, 100),
)
assert resultats == [
ScoredSite(
site_id="site-a",
target_at=resultats[0].target_at,
status="insufficient_data",
predicted_value=None,
failure_reason=resultats[0].failure_reason,
)
]
assert booster.appels == []
def test_score_predicts_when_history_is_sufficient() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
scoring_frame = build_scoring_frame(make_recent("site-a", heures=200, depart=depart))
booster = FakeBooster(valeur=99.5)
resultats = score(
booster, # type: ignore[arg-type]
scoring_frame,
instant=target_at_for(depart, 200),
)
assert len(resultats) == 1
assert resultats[0].status == "available"
assert resultats[0].predicted_value == 99.5
assert resultats[0].failure_reason is None
assert booster.appels == [1]
def test_score_marks_a_stale_site_as_insufficient_data_without_calling_the_model() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
# Historique largement suffisant (168h+), mais l'instant de reference est loin apres la
# derniere lecture : la fraicheur doit primer sur la disponibilite de l'historique.
scoring_frame = build_scoring_frame(make_recent("site-a", heures=200, depart=depart))
instant = target_at_for(depart, 200) + MAX_STALENESS + timedelta(hours=1)
booster = FakeBooster()
resultats = score(booster, scoring_frame, instant=instant) # type: ignore[arg-type]
assert len(resultats) == 1
assert resultats[0].status == "insufficient_data"
assert resultats[0].predicted_value is None
assert "vieille" in (resultats[0].failure_reason or "")
assert booster.appels == []
def test_score_accepts_a_reading_exactly_at_the_staleness_threshold() -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
scoring_frame = build_scoring_frame(make_recent("site-a", heures=200, depart=depart))
# `target_at_for(...)` donne deja un age d'1h (cf. sa docstring) : retrancher cette heure
# pour retomber exactement sur le seuil, ni en dessous ni au dessus.
instant = target_at_for(depart, 200) + MAX_STALENESS - timedelta(hours=1)
booster = FakeBooster(valeur=12.0)
resultats = score(booster, scoring_frame, instant=instant) # type: ignore[arg-type]
assert resultats[0].status == "available"
assert booster.appels == [1]
def test_write_predictions_does_nothing_when_there_is_nothing_to_write() -> None:
connection = FakeConnection()
write_predictions(connection, [], reference="lightgbm-test") # type: ignore[arg-type]
assert connection.appels == []
def test_write_predictions_sends_one_row_per_result() -> None:
connection = FakeConnection()
resultats = [
ScoredSite("site-a", datetime(2026, 1, 1, tzinfo=UTC), "available", 42.0, None),
ScoredSite(
"site-b",
datetime(2026, 1, 1, tzinfo=UTC),
"insufficient_data",
None,
"pas assez d'historique",
),
]
write_predictions(connection, resultats, reference="lightgbm-test") # type: ignore[arg-type]
assert len(connection.appels) == 1
_, lignes = connection.appels[0]
assert len(lignes) == 2
assert lignes[0]["model_reference"] == "lightgbm-test"
assert lignes[0]["target_metric"] == "consumption_kwh"
assert lignes[0]["period_minutes"] == 60
def test_model_reference_is_stable_for_the_same_file_content(tmp_path: Path) -> None:
model_path = tmp_path / "model.txt"
model_path.write_bytes(b"contenu-du-modele")
assert model_reference(model_path) == model_reference(model_path)
def test_model_reference_changes_with_the_file_content(tmp_path: Path) -> None:
premier = tmp_path / "model-a.txt"
premier.write_bytes(b"version-1")
second = tmp_path / "model-b.txt"
second.write_bytes(b"version-2")
assert model_reference(premier) != model_reference(second)
def test_run_scoring_in_csv_mode_scores_without_touching_a_database(tmp_path: Path) -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
frame = pd.concat(
[
make_recent("site-a", heures=400, depart=depart),
make_recent("site-b", heures=400, depart=depart),
],
ignore_index=True,
)
csv_path = tmp_path / "recent.csv"
frame.to_csv(csv_path, index=False)
model_path = tmp_path / "model.txt"
model_path.write_bytes(b"peu importe le contenu pour ce test")
with pytest.MonkeyPatch.context() as monkeypatch:
monkeypatch.setattr(
"enervision_ml.score.lgb.Booster", lambda model_file: FakeBooster(valeur=7.0)
)
resultats = run_scoring(model_path=model_path, csv_path=csv_path)
assert {r.site_id for r in resultats} == {"site-a", "site-b"}
assert all(r.status == "available" for r in resultats)
assert all(r.predicted_value == 7.0 for r in resultats)