Compare commits

..
Author SHA1 Message Date
ValentinDeFariaandGitHub f9c2a4610c Update dashboard.ts
Frontend / Audit des dépendances (push) Successful in 6s
SonarQube / build-back (push) Successful in 1m6s
Frontend / build (push) Successful in 9m52s
SonarQube / build-front (push) Successful in 9m44s
SonarQube / test-back (push) Failing after 52s
Frontend / test (push) Failing after 5m0s
SonarQube / test-front (push) Failing after 5m2s
SonarQube / SonarQube (push) Skipped
2026-09-18 16:56:59 +02:00
ValentinDeFariaandGitHub b5fa7b0010 Merge branch 'dev' into feat/supervision-des-capteurs 2026-09-18 16:54:39 +02:00
Valentin 7f710c9084 feat(frontend): supervision des capteurs par site (admin) 2026-09-18 12:02:48 +02:00
44 changed files with 451 additions and 2085 deletions
+1 -4
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-score recommendations 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}'
@@ -77,9 +77,6 @@ ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DAT
ml-score: ## Score le prochain pas horaire et l'ecrit dans `prediction`. CSV=chemin optionnel 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),) cd $(ML) && uv run python -m enervision_ml.score $(if $(CSV),--csv $(CSV),)
recommendations: ## Genere les recommandations depuis les alertes en base. SITE=identifiant optionnel
cd $(BACKEND) && uv run python -m app.cli generate-recommendations $(if $(SITE),--site-id $(SITE),)
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)
-1
View File
@@ -113,7 +113,6 @@ Le sens de dependance est unique : `endpoints` vers `services` vers `repositorie
| `/api/v1/sites/{site_id}` | Décrit un site | `lecteur` | | `/api/v1/sites/{site_id}` | Décrit un site | `lecteur` |
| `/api/v1/recommendations` | Liste les recommandations | `lecteur` | | `/api/v1/recommendations` | Liste les recommandations | `lecteur` |
| `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation | `lecteur` | | `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation | `lecteur` |
| `/api/v1/recommendations/generate` | Génère les recommandations depuis les alertes (POST) | `admin` |
| `/metrics` | Métriques au format Prometheus | jeton si `APP_METRICS_TOKEN` | | `/metrics` | Métriques au format Prometheus | jeton si `APP_METRICS_TOKEN` |
| `/docs`, `/openapi.json` | Documentation, fermée en `staging` et `prod` | public sinon | | `/docs`, `/openapi.json` | Documentation, fermée en `staging` et `prod` | public sinon |
+2 -11
View File
@@ -180,23 +180,14 @@ SiteServiceDep = Annotated[SiteService, Depends(get_site_service)]
def get_alert_service(session: SessionDep) -> AlertService: def get_alert_service(session: SessionDep) -> AlertService:
return AlertService( return AlertService(alerts=AlertRepository(session))
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)]
def get_recommendation_service(session: SessionDep) -> RecommendationService: def get_recommendation_service(session: SessionDep) -> RecommendationService:
return RecommendationService( return RecommendationService(recommendations=RecommendationRepository(session))
recommendations=RecommendationRepository(session),
alerts=AlertRepository(session),
transaction=session,
)
RecommendationServiceDep = Annotated[RecommendationService, Depends(get_recommendation_service)] RecommendationServiceDep = Annotated[RecommendationService, Depends(get_recommendation_service)]
+1 -1
View File
@@ -63,7 +63,7 @@ TAGS: Final[list[dict[str, Any]]] = [
"name": "recommendations", "name": "recommendations",
"description": ( "description": (
"Consultation des recommandations issues des alertes. Accessible à partir du rôle " "Consultation des recommandations issues des alertes. Accessible à partir du rôle "
"`lecteur`. Leur génération par le moteur de règles est réservée au rôle `admin`." "`lecteur`."
), ),
}, },
{ {
@@ -1,18 +1,13 @@
from fastapi import APIRouter, HTTPException, status from fastapi import APIRouter, HTTPException, status
from app.api.deps import AdminDep, LecteurDep, RecommendationServiceDep from app.api.deps import LecteurDep, RecommendationServiceDep
from app.api.openapi import REPONSE_VALIDATION, REPONSES_ADMIN, Reponses from app.api.openapi import REPONSE_VALIDATION, Reponses
from app.schemas.errors import ErrorResponse from app.schemas.errors import ErrorResponse
from app.schemas.recommendation import ( from app.schemas.recommendation import RecommendationResponse
RecommendationGenerationResponse,
RecommendationResponse,
)
from app.services.recommendation import RecommendationNotFoundError from app.services.recommendation import RecommendationNotFoundError
router = APIRouter() router = APIRouter()
REPONSES_GENERATION: Reponses = {**REPONSES_ADMIN, **REPONSE_VALIDATION}
REPONSES_INTROUVABLE: Reponses = { REPONSES_INTROUVABLE: Reponses = {
**REPONSE_VALIDATION, **REPONSE_VALIDATION,
404: {"model": ErrorResponse, "description": "Aucune recommandation ne porte cet identifiant."}, 404: {"model": ErrorResponse, "description": "Aucune recommandation ne porte cet identifiant."},
@@ -43,22 +38,3 @@ async def get_recommendation(
status_code=status.HTTP_404_NOT_FOUND, detail="Recommandation introuvable" status_code=status.HTTP_404_NOT_FOUND, detail="Recommandation introuvable"
) from erreur ) from erreur
return RecommendationResponse.model_validate(recommendation) return RecommendationResponse.model_validate(recommendation)
@router.post(
"/generate",
response_model=RecommendationGenerationResponse,
summary="Génère les recommandations à partir des alertes",
responses=REPONSES_GENERATION,
)
async def generate_recommendations(
_: AdminDep,
service: RecommendationServiceDep,
site_id: str | None = None,
) -> RecommendationGenerationResponse:
rapport = await service.generate(site_id=site_id)
return RecommendationGenerationResponse(
alerts_examined=rapport.alertes_examinees,
recommendations_created=rapport.recommandations_creees,
already_present=rapport.deja_presentes,
)
-31
View File
@@ -22,11 +22,8 @@ from app.core.hashing import build_hasher
from app.core.roles import Role from app.core.roles import Role
from app.db.session import get_session_factory from app.db.session import get_session_factory
from app.main import create_app from app.main import create_app
from app.repositories.alert import AlertRepository
from app.repositories.recommendation import RecommendationRepository
from app.repositories.user import UserRepository from app.repositories.user import UserRepository
from app.schemas.auth import PASSWORD_MIN_LENGTH, SPECIAL_CHARACTERS, valide_complexite from app.schemas.auth import PASSWORD_MIN_LENGTH, SPECIAL_CHARACTERS, valide_complexite
from app.services.recommendation import RecommendationService
LONGUEUR_MOT_DE_PASSE_GENERE = 24 LONGUEUR_MOT_DE_PASSE_GENERE = 24
CHEMIN_CONTRAT = Path(__file__).resolve().parent.parent / "openapi.json" CHEMIN_CONTRAT = Path(__file__).resolve().parent.parent / "openapi.json"
@@ -66,22 +63,6 @@ async def create_admin(
) )
async def generate_recommendations(*, site_id: str | None) -> str:
async with get_session_factory()() as session:
service = RecommendationService(
recommendations=RecommendationRepository(session),
alerts=AlertRepository(session),
transaction=session,
)
rapport = await service.generate(site_id=site_id)
return (
f"{rapport.alertes_examinees} alerte(s) examinée(s), "
f"{rapport.recommandations_creees} recommandation(s) créée(s), "
f"{rapport.deja_presentes} déjà présente(s)"
)
# Piège : le schéma ne doit dépendre ni du `.env` du poste ni des variables `APP_*`, sinon le # Piège : le schéma ne doit dépendre ni du `.env` du poste ni des variables `APP_*`, sinon le
# fichier versionné changerait de machine en machine et le test de dérive deviendrait un oracle # fichier versionné changerait de machine en machine et le test de dérive deviendrait un oracle
# de configuration locale. Tout ce qui atteint le schéma est donc posé ici, `_env_file` compris. # de configuration locale. Tout ce qui atteint le schéma est donc posé ici, `_env_file` compris.
@@ -128,14 +109,6 @@ def build_parser() -> argparse.ArgumentParser:
"export-openapi", help="Écrit le contrat OpenAPI sur disque" "export-openapi", help="Écrit le contrat OpenAPI sur disque"
) )
contrat.add_argument("--output", default=str(CHEMIN_CONTRAT)) contrat.add_argument("--output", default=str(CHEMIN_CONTRAT))
recommandations = sous_commandes.add_parser(
"generate-recommendations",
help="Applique le moteur de règles aux alertes en base",
)
recommandations.add_argument(
"--site-id", default=None, help="Limite le traitement aux alertes d'un site"
)
return parser return parser
@@ -179,10 +152,6 @@ def main(argv: list[str] | None = None) -> int:
print(export_openapi(Path(arguments.output))) print(export_openapi(Path(arguments.output)))
return 0 return 0
if arguments.commande == "generate-recommendations":
print(asyncio.run(generate_recommendations(site_id=arguments.site_id)))
return 0
mot_de_passe = read_password(generate=arguments.generate) mot_de_passe = read_password(generate=arguments.generate)
succes, message = asyncio.run( succes, message = asyncio.run(
@@ -1,68 +0,0 @@
# 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,7 +1,6 @@
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
@@ -20,36 +19,3 @@ 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()
@@ -1,5 +1,4 @@
from collections.abc import Sequence from collections.abc import Sequence
from datetime import datetime
from sqlalchemy import select from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
@@ -11,26 +10,6 @@ class PredictionRepository:
def __init__(self, session: AsyncSession) -> None: def __init__(self, session: AsyncSession) -> None:
self._session = session 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]: async def latest_by_site(self) -> Sequence[Prediction]:
# `.distinct(site_id)` compile en `DISTINCT ON (site_id)` sous PostgreSQL : une seule # `.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 # ligne par site, la plus récente grâce à l'ordre composite qui suit. Même mécanisme que
-15
View File
@@ -34,21 +34,6 @@ 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,
*, *,
@@ -1,24 +1,11 @@
from collections.abc import Sequence from collections.abc import Sequence
from dataclasses import asdict, dataclass
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 Recommendation from app.models.energy import Recommendation
@dataclass(frozen=True, slots=True)
class NouvelleRecommandation:
alert_id: int
action: str
explanation: str
rule_reference: str
TAILLE_DE_LOT = 1000
class RecommendationRepository: class RecommendationRepository:
def __init__(self, session: AsyncSession) -> None: def __init__(self, session: AsyncSession) -> None:
self._session = session self._session = session
@@ -33,19 +20,3 @@ class RecommendationRepository:
) )
recommendation: Recommendation | None = await self._session.scalar(requete) recommendation: Recommendation | None = await self._session.scalar(requete)
return recommendation return recommendation
# Pourquoi : l'idempotence est déléguée à `uq_recommendation_alert_rule` plutôt qu'à une
# lecture préalable, qui laisserait une fenêtre entre le contrôle et l'insertion.
async def create_missing(self, nouvelles: Sequence[NouvelleRecommandation]) -> int:
creees = 0
# Piège : asyncpg plafonne une requête à 32 767 paramètres, soit 8 191 lignes de quatre
# colonnes. Au-delà de ce seuil un `INSERT` d'un seul tenant échouerait.
for debut in range(0, len(nouvelles), TAILLE_DE_LOT):
requete = (
insert(Recommendation)
.values([asdict(nouvelle) for nouvelle in nouvelles[debut : debut + TAILLE_DE_LOT]])
.on_conflict_do_nothing(constraint="uq_recommendation_alert_rule")
.returning(Recommendation.recommendation_id)
)
creees += len((await self._session.scalars(requete)).all())
return creees
@@ -12,9 +12,3 @@ class RecommendationResponse(BaseModel):
explanation: str explanation: str
rule_reference: str rule_reference: str
created_at: datetime created_at: datetime
class RecommendationGenerationResponse(BaseModel):
alerts_examined: int
recommendations_created: int
already_present: int
+2 -311
View File
@@ -1,323 +1,14 @@
from collections.abc import Sequence from collections.abc import Sequence
from datetime import UTC, datetime, timedelta
from app.models.energy import Alert, Prediction, Reading, Site from app.models.energy import Alert
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__( def __init__(self, *, alerts: AlertRepository) -> None:
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
+1 -37
View File
@@ -1,15 +1,7 @@
from collections.abc import Sequence from collections.abc import Sequence
from dataclasses import dataclass
from typing import Protocol
from app.models.energy import Recommendation from app.models.energy import Recommendation
from app.repositories.alert import AlertRepository
from app.repositories.recommendation import RecommendationRepository from app.repositories.recommendation import RecommendationRepository
from app.services.recommendation_rules import applique_les_regles
class Transaction(Protocol):
async def commit(self) -> None: ...
class RecommendationError(Exception): class RecommendationError(Exception):
@@ -20,24 +12,9 @@ class RecommendationNotFoundError(RecommendationError):
pass pass
@dataclass(frozen=True, slots=True)
class RapportGeneration:
alertes_examinees: int
recommandations_creees: int
deja_presentes: int
class RecommendationService: class RecommendationService:
def __init__( def __init__(self, *, recommendations: RecommendationRepository) -> None:
self,
*,
recommendations: RecommendationRepository,
alerts: AlertRepository,
transaction: Transaction,
) -> None:
self._recommendations = recommendations self._recommendations = recommendations
self._alerts = alerts
self._transaction = transaction
async def list_all(self) -> Sequence[Recommendation]: async def list_all(self) -> Sequence[Recommendation]:
return await self._recommendations.list_all() return await self._recommendations.list_all()
@@ -47,16 +24,3 @@ class RecommendationService:
if recommendation is None: if recommendation is None:
raise RecommendationNotFoundError(recommendation_id) raise RecommendationNotFoundError(recommendation_id)
return recommendation return recommendation
async def generate(self, *, site_id: str | None = None) -> RapportGeneration:
alertes = await self._alerts.list_all(site_id=site_id)
nouvelles = [nouvelle for alerte in alertes for nouvelle in applique_les_regles(alerte)]
creees = await self._recommendations.create_missing(nouvelles)
await self._transaction.commit()
return RapportGeneration(
alertes_examinees=len(alertes),
recommandations_creees=creees,
deja_presentes=len(nouvelles) - creees,
)
@@ -1,117 +0,0 @@
# Piège : `rule_reference` est la clé d'idempotence en base, portée par la contrainte
# `uq_recommendation_alert_rule`. Renommer une référence déjà livrée ne remplace pas les
# recommandations existantes, il en crée de nouvelles à côté. Une règle qui change de sens
# prend donc une référence suffixée `-v2` - REGLES.
from collections.abc import Callable
from dataclasses import dataclass
from typing import Final
from app.models.energy import Alert
from app.repositories.recommendation import NouvelleRecommandation
from app.schemas.alert import AlertSeverity, AlertType
FACTEUR_DEPASSEMENT_MAJEUR: Final = 1.2
POURCENTAGE_DEPASSEMENT_MAJEUR: Final = round((FACTEUR_DEPASSEMENT_MAJEUR - 1) * 100)
@dataclass(frozen=True, slots=True)
class Regle:
reference: str
action: str
declencheur: Callable[[Alert], bool]
motif: Callable[[Alert], str]
def _du_type(attendu: AlertType) -> Callable[[Alert], bool]:
return lambda alerte: alerte.type == attendu
def _de_severite(attendue: AlertSeverity) -> Callable[[Alert], bool]:
return lambda alerte: alerte.severity == attendue
# Un seuil nul ou négatif rendrait le rapport `value / threshold` arbitraire : l'alerte ne
# renseigne alors aucun dépassement exploitable, et la règle ne se déclenche pas.
def _depasse_largement_le_seuil(alerte: Alert) -> bool:
if alerte.value is None or alerte.threshold is None or alerte.threshold <= 0:
return False
return alerte.value >= alerte.threshold * FACTEUR_DEPASSEMENT_MAJEUR
REGLES: Final[tuple[Regle, ...]] = (
Regle(
reference="spike-delestage-v1",
action="Délester les équipements non prioritaires sur le créneau du pic",
declencheur=_du_type(AlertType.SPIKE),
motif=lambda alerte: f"Pic de consommation signalé sur le site {alerte.site_id}",
),
Regle(
reference="threshold-reduction-v1",
action="Ramener la puissance appelée sous le seuil contractuel",
declencheur=_du_type(AlertType.THRESHOLD),
motif=lambda alerte: f"Seuil de consommation dépassé sur le site {alerte.site_id}",
),
Regle(
reference="outage-secours-v1",
action="Basculer sur l'alimentation de secours et prévenir l'exploitant",
declencheur=_du_type(AlertType.OUTAGE),
motif=lambda alerte: (
f"Risque de surcharge ou de coupure imminente sur le site {alerte.site_id}"
),
),
Regle(
reference="sensor-maintenance-v1",
action="Planifier une intervention de maintenance sur le capteur",
declencheur=_du_type(AlertType.SENSOR),
motif=lambda alerte: (
f"Capteur défaillant sur le site {alerte.site_id}, les mesures ne sont plus fiables"
),
),
Regle(
reference="anomaly-verification-v1",
action="Confronter la mesure à la prévision et vérifier le paramétrage du site",
declencheur=_du_type(AlertType.ANOMALY),
motif=lambda alerte: (
f"Écart anormal entre la mesure et le comportement attendu du site {alerte.site_id}"
),
),
Regle(
reference="escalade-astreinte-v1",
action="Escalader à l'astreinte sous une heure",
declencheur=_de_severite(AlertSeverity.CRITICAL),
motif=lambda alerte: f"Alerte de sévérité critique sur le site {alerte.site_id}",
),
Regle(
reference="contrat-puissance-v1",
action="Réévaluer la puissance souscrite au contrat",
declencheur=_depasse_largement_le_seuil,
motif=lambda alerte: (
f"Dépassement d'au moins {POURCENTAGE_DEPASSEMENT_MAJEUR} % du seuil "
f"sur le site {alerte.site_id}"
),
),
)
def applique_les_regles(alerte: Alert) -> list[NouvelleRecommandation]:
contexte = _contexte_de_mesure(alerte)
return [
NouvelleRecommandation(
alert_id=alerte.alert_id,
action=regle.action,
explanation=f"{regle.motif(alerte)}{contexte}.",
rule_reference=regle.reference,
)
for regle in REGLES
if regle.declencheur(alerte)
]
def _contexte_de_mesure(alerte: Alert) -> str:
if alerte.value is None:
return ""
grandeur = alerte.metric or "valeur"
if alerte.threshold is None:
return f" ({grandeur} mesurée à {alerte.value})"
return f" ({grandeur} mesurée à {alerte.value}, seuil {alerte.threshold})"
+1 -108
View File
@@ -1453,90 +1453,6 @@
} }
} }
}, },
"/api/v1/recommendations/generate": {
"post": {
"tags": [
"recommendations"
],
"summary": "Génère les recommandations à partir des alertes",
"operationId": "generate_recommendations_api_v1_recommendations_generate_post",
"security": [
{
"Jeton d'accès": []
}
],
"parameters": [
{
"name": "site_id",
"in": "query",
"required": false,
"schema": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"title": "Site Id"
}
}
],
"responses": {
"200": {
"description": "Successful Response",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/RecommendationGenerationResponse"
}
}
}
},
"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": "Droits insuffisants, ou mot de passe provisoire à changer quand `detail` vaut `password_change_required`.",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/ErrorResponse"
}
}
}
},
"422": {
"description": "Corps invalide. Le détail nomme le champ fautif et le type d'erreur, jamais la valeur envoyée.",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/ValidationErrorResponse"
}
}
}
}
}
}
},
"/api/v1/stats/summary": { "/api/v1/stats/summary": {
"get": { "get": {
"tags": [ "tags": [
@@ -2428,29 +2344,6 @@
], ],
"title": "ReadingSource" "title": "ReadingSource"
}, },
"RecommendationGenerationResponse": {
"properties": {
"alerts_examined": {
"type": "integer",
"title": "Alerts Examined"
},
"recommendations_created": {
"type": "integer",
"title": "Recommendations Created"
},
"already_present": {
"type": "integer",
"title": "Already Present"
}
},
"type": "object",
"required": [
"alerts_examined",
"recommendations_created",
"already_present"
],
"title": "RecommendationGenerationResponse"
},
"RecommendationResponse": { "RecommendationResponse": {
"properties": { "properties": {
"recommendation_id": { "recommendation_id": {
@@ -3260,7 +3153,7 @@
}, },
{ {
"name": "recommendations", "name": "recommendations",
"description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`. Leur génération par le moteur de règles est réservée au rôle `admin`." "description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`."
}, },
{ {
"name": "stats", "name": "stats",
-1
View File
@@ -51,7 +51,6 @@ ROLE_MINIMUM: Final[dict[Route, Role]] = {
("GET", "/api/v1/alerts"): Role.LECTEUR, ("GET", "/api/v1/alerts"): Role.LECTEUR,
("GET", "/api/v1/recommendations"): Role.LECTEUR, ("GET", "/api/v1/recommendations"): Role.LECTEUR,
("GET", "/api/v1/recommendations/{recommendation_id}"): Role.LECTEUR, ("GET", "/api/v1/recommendations/{recommendation_id}"): Role.LECTEUR,
("POST", "/api/v1/recommendations/generate"): Role.ADMIN,
("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/predictions"): Role.LECTEUR,
+1 -59
View File
@@ -10,7 +10,7 @@ from app.api.deps import get_current_principal, get_recommendation_service
from app.core.principal import Principal from app.core.principal import Principal
from app.core.roles import AccountKind, Role from app.core.roles import AccountKind, Role
from app.models.energy import Recommendation from app.models.energy import Recommendation
from app.services.recommendation import RapportGeneration, RecommendationNotFoundError from app.services.recommendation import RecommendationNotFoundError
MOMENT = datetime(2024, 1, 1, tzinfo=UTC) MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
@@ -40,7 +40,6 @@ class FauxService:
def __init__(self, erreur: Exception | None = None) -> None: def __init__(self, erreur: Exception | None = None) -> None:
self._erreur = erreur self._erreur = erreur
self.recommendation = recommendation() self.recommendation = recommendation()
self.site_demande: str | None = None
async def list_all(self) -> list[Recommendation]: async def list_all(self) -> list[Recommendation]:
return [self.recommendation] return [self.recommendation]
@@ -50,10 +49,6 @@ class FauxService:
raise self._erreur raise self._erreur
return self.recommendation return self.recommendation
async def generate(self, *, site_id: str | None = None) -> RapportGeneration:
self.site_demande = site_id
return RapportGeneration(alertes_examinees=2, recommandations_creees=3, deja_presentes=1)
@pytest.fixture @pytest.fixture
def lecteur_connecte(app: FastAPI) -> Iterator[None]: def lecteur_connecte(app: FastAPI) -> Iterator[None]:
@@ -147,56 +142,3 @@ async def test_get_recommendation_returns_404_when_the_session_finds_nothing(
response = await client.get("/api/v1/recommendations/404") response = await client.get("/api/v1/recommendations/404")
assert response.status_code == 404 assert response.status_code == 404
@pytest.fixture
def admin_connecte(app: FastAPI) -> Iterator[None]:
app.dependency_overrides[get_current_principal] = lambda: principal(Role.ADMIN)
yield
app.dependency_overrides.pop(get_current_principal, None)
@pytest.fixture
def servi_en_admin(app: FastAPI, admin_connecte: None) -> Iterator[Callable[[], FauxService]]:
def installe() -> FauxService:
service = FauxService()
app.dependency_overrides[get_recommendation_service] = lambda: service
return service
yield installe
app.dependency_overrides.pop(get_recommendation_service, None)
async def test_generate_recommendations_returns_the_generation_report(
servi_en_admin: Callable[[], FauxService], client: AsyncClient
) -> None:
servi_en_admin()
response = await client.post("/api/v1/recommendations/generate")
assert response.status_code == 200
assert response.json() == {
"alerts_examined": 2,
"recommendations_created": 3,
"already_present": 1,
}
async def test_generate_recommendations_forwards_the_requested_site(
servi_en_admin: Callable[[], FauxService], client: AsyncClient
) -> None:
service = servi_en_admin()
await client.post("/api/v1/recommendations/generate", params={"site_id": "SITE002"})
assert service.site_demande == "SITE002"
async def test_generate_recommendations_refuses_a_reader(
servi: Callable[..., FauxService], client: AsyncClient
) -> None:
servi()
response = await client.post("/api/v1/recommendations/generate")
assert response.status_code == 403
@@ -89,60 +89,3 @@ 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 == []
@@ -29,85 +29,6 @@ async def creer_prediction(
return prediction return prediction
async def test_list_since_excludes_predictions_before_the_cutoff(session: AsyncSession) -> None:
site = await creer_site(session)
depot = PredictionRepository(session)
dedans = await creer_prediction(
session, site_id=site.site_id, target_at=datetime(2026, 9, 16, tzinfo=UTC)
)
await creer_prediction(
session, site_id=site.site_id, target_at=datetime(2026, 9, 1, tzinfo=UTC)
)
resultats = await depot.list_since(
since=datetime(2026, 9, 10, tzinfo=UTC), site_id=site.site_id
)
identifiants = [p.prediction_id for p in resultats]
await session.rollback()
assert identifiants == [dedans.prediction_id]
async def test_list_since_excludes_predictions_that_are_not_available(
session: AsyncSession,
) -> None:
site = await creer_site(session)
depot = PredictionRepository(session)
await creer_prediction(
session,
site_id=site.site_id,
target_at=datetime(2026, 9, 16, tzinfo=UTC),
status="insufficient_data",
predicted_value=None,
failure_reason="pas assez d'historique",
)
resultats = await depot.list_since(since=datetime(2026, 9, 1, tzinfo=UTC), site_id=site.site_id)
await session.rollback()
assert list(resultats) == []
async def test_list_since_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: async def test_latest_by_site_keeps_only_the_most_recent_target(session: AsyncSession) -> None:
site = await creer_site(session) site = await creer_site(session)
depot = PredictionRepository(session) depot = PredictionRepository(session)
@@ -155,79 +155,6 @@ 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:
@@ -5,8 +5,7 @@ import pytest
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.models.energy import Alert, Recommendation, Site from app.models.energy import Alert, Recommendation, Site
from app.repositories import recommendation as module_recommendation from app.repositories.recommendation import RecommendationRepository
from app.repositories.recommendation import NouvelleRecommandation, RecommendationRepository
pytestmark = pytest.mark.integration pytestmark = pytest.mark.integration
@@ -84,59 +83,3 @@ async def test_list_all_returns_the_recommendations_sorted_by_identifier(
await session.rollback() await session.rollback()
assert identifiants == sorted(identifiants) assert identifiants == sorted(identifiants)
def nouvelle(alert_id: int, reference: str = "spike-delestage-v1") -> NouvelleRecommandation:
return NouvelleRecommandation(
alert_id=alert_id,
action="Délester les équipements non prioritaires",
explanation="Pic de consommation signalé.",
rule_reference=reference,
)
async def test_create_missing_inserts_the_proposals(session: AsyncSession) -> None:
depot = RecommendationRepository(session)
alert_id = await creer_alerte(session)
creees = await depot.create_missing(
[nouvelle(alert_id), nouvelle(alert_id, "escalade-astreinte-v1")]
)
await session.rollback()
assert creees == 2
async def test_create_missing_ignores_a_rule_already_held_for_the_alert(
session: AsyncSession,
) -> None:
depot = RecommendationRepository(session)
alert_id = await creer_alerte(session)
await depot.create_missing([nouvelle(alert_id)])
creees = await depot.create_missing([nouvelle(alert_id)])
await session.rollback()
assert creees == 0
async def test_create_missing_returns_zero_without_any_proposal(session: AsyncSession) -> None:
creees = await RecommendationRepository(session).create_missing([])
assert creees == 0
async def test_create_missing_inserts_every_proposal_across_several_batches(
session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setattr(module_recommendation, "TAILLE_DE_LOT", 2)
depot = RecommendationRepository(session)
alert_id = await creer_alerte(session)
propositions = [nouvelle(alert_id, f"regle-{index}-v1") for index in range(5)]
creees = await depot.create_missing(propositions)
enregistrees = [r for r in await depot.list_all() if r.alert_id == alert_id]
await session.rollback()
assert creees == 5
assert len(enregistrees) == 5
+6 -408
View File
@@ -1,10 +1,7 @@
from dataclasses import dataclass, field from datetime import UTC, datetime
from datetime import UTC, datetime, timedelta
from app.models.energy import Alert from app.models.energy import Alert
from app.services.alert import OUTAGE_THRESHOLD, AlertService, _severity_from_ratio from app.services.alert import AlertService
NOW = datetime(2026, 9, 16, 12, 0, tzinfo=UTC)
def alert( def alert(
@@ -29,36 +26,10 @@ 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
@@ -66,392 +37,19 @@ 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:
svc, _ = service(sites=[], alerts=FakeRepository([alert(1), alert(2)])) service = AlertService(alerts=FakeRepository([alert(1), alert(2)]))
alertes = await svc.list_all() alertes = await service.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([])
svc, _ = service(sites=[], alerts=depot) service = AlertService(alerts=depot)
await svc.list_all(site_id="site-1", severity="critical") await service.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 == []
@@ -1,14 +1,10 @@
from collections.abc import Sequence
from datetime import UTC, datetime from datetime import UTC, datetime
import pytest import pytest
from app.models.energy import Alert, Recommendation from app.models.energy import Recommendation
from app.repositories.recommendation import NouvelleRecommandation
from app.services.recommendation import RecommendationNotFoundError, RecommendationService from app.services.recommendation import RecommendationNotFoundError, RecommendationService
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
def recommendation(recommendation_id: int = 1) -> Recommendation: def recommendation(recommendation_id: int = 1) -> Recommendation:
return Recommendation( return Recommendation(
@@ -17,33 +13,13 @@ def recommendation(recommendation_id: int = 1) -> Recommendation:
action="Vérifier la consommation", action="Vérifier la consommation",
explanation="Pic détecté", explanation="Pic détecté",
rule_reference="spike-v1", rule_reference="spike-v1",
created_at=MOMENT, created_at=datetime(2024, 1, 1, tzinfo=UTC),
)
def alerte(alert_id: int = 1, site_id: str = "SITE001", severity: str = "high") -> Alert:
return Alert(
alert_id=alert_id,
source_alert_id=f"ALR-{alert_id}",
site_id=site_id,
source="api_mock",
timestamp=MOMENT,
type="spike",
severity=severity,
message="Pic de consommation",
value=None,
threshold=None,
metric=None,
prediction_id=None,
raw_data={},
) )
class FakeRepository: class FakeRepository:
def __init__(self, recommendations: list[Recommendation], creees: int | None = None) -> None: def __init__(self, recommendations: list[Recommendation]) -> None:
self._recommendations = recommendations self._recommendations = recommendations
self._creees = creees
self.recues: list[NouvelleRecommandation] = []
async def list_all(self) -> list[Recommendation]: async def list_all(self) -> list[Recommendation]:
return self._recommendations return self._recommendations
@@ -53,111 +29,27 @@ class FakeRepository:
(r for r in self._recommendations if r.recommendation_id == recommendation_id), None (r for r in self._recommendations if r.recommendation_id == recommendation_id), None
) )
async def create_missing(self, nouvelles: Sequence[NouvelleRecommandation]) -> int:
self.recues = list(nouvelles)
return len(self.recues) if self._creees is None else self._creees
class FakeAlertRepository:
def __init__(self, alertes: list[Alert]) -> None:
self._alertes = alertes
self.site_demande: str | None = None
async def list_all(
self, *, site_id: str | None = None, severity: str | None = None
) -> list[Alert]:
self.site_demande = site_id
if site_id is None:
return self._alertes
return [a for a in self._alertes if a.site_id == site_id]
class FakeTransaction:
def __init__(self) -> None:
self.commits = 0
async def commit(self) -> None:
self.commits += 1
def service(
recommendations: FakeRepository | None = None,
alerts: FakeAlertRepository | None = None,
transaction: FakeTransaction | None = None,
) -> RecommendationService:
return RecommendationService(
recommendations=recommendations or FakeRepository([]),
alerts=alerts or FakeAlertRepository([]),
transaction=transaction or FakeTransaction(),
)
async def test_list_all_returns_the_repository_recommendations() -> None: async def test_list_all_returns_the_repository_recommendations() -> None:
depot = FakeRepository([recommendation(1), recommendation(2)]) service = RecommendationService(
recommendations=FakeRepository([recommendation(1), recommendation(2)])
)
recommendations = await service(recommendations=depot).list_all() recommendations = await service.list_all()
assert [r.recommendation_id for r in recommendations] == [1, 2] assert [r.recommendation_id for r in recommendations] == [1, 2]
async def test_get_by_id_returns_the_matching_recommendation() -> None: async def test_get_by_id_returns_the_matching_recommendation() -> None:
trouve = await service(recommendations=FakeRepository([recommendation(1)])).get_by_id(1) service = RecommendationService(recommendations=FakeRepository([recommendation(1)]))
trouve = await service.get_by_id(1)
assert trouve.recommendation_id == 1 assert trouve.recommendation_id == 1
async def test_get_by_id_raises_when_the_recommendation_is_unknown() -> None: async def test_get_by_id_raises_when_the_recommendation_is_unknown() -> None:
service = RecommendationService(recommendations=FakeRepository([]))
with pytest.raises(RecommendationNotFoundError): with pytest.raises(RecommendationNotFoundError):
await service().get_by_id(404) await service.get_by_id(404)
async def test_generate_persists_one_proposal_per_triggered_rule() -> None:
depot = FakeRepository([])
rapport = await service(
recommendations=depot, alerts=FakeAlertRepository([alerte(severity="critical")])
).generate()
assert {n.rule_reference for n in depot.recues} == {
"spike-delestage-v1",
"escalade-astreinte-v1",
}
assert rapport.recommandations_creees == 2
async def test_generate_commits_once() -> None:
transaction = FakeTransaction()
await service(alerts=FakeAlertRepository([alerte()]), transaction=transaction).generate()
assert transaction.commits == 1
async def test_generate_restricts_the_alerts_to_the_requested_site() -> None:
alertes = FakeAlertRepository([alerte(1, site_id="SITE001"), alerte(2, site_id="SITE002")])
depot = FakeRepository([])
rapport = await service(recommendations=depot, alerts=alertes).generate(site_id="SITE002")
assert alertes.site_demande == "SITE002"
assert rapport.alertes_examinees == 1
assert {n.alert_id for n in depot.recues} == {2}
async def test_generate_reports_nothing_when_no_alert_matches() -> None:
rapport = await service().generate()
assert rapport.alertes_examinees == 0
assert rapport.recommandations_creees == 0
assert rapport.deja_presentes == 0
async def test_generate_counts_the_proposals_the_database_already_held() -> None:
depot = FakeRepository([], creees=0)
rapport = await service(
recommendations=depot, alerts=FakeAlertRepository([alerte()])
).generate()
assert rapport.recommandations_creees == 0
assert rapport.deja_presentes == 1
@@ -1,142 +0,0 @@
from datetime import UTC, datetime
import pytest
from app.models.energy import Alert
from app.services.recommendation_rules import FACTEUR_DEPASSEMENT_MAJEUR, applique_les_regles
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
def alerte(
*,
alert_id: int = 1,
type_alerte: str = "spike",
severity: str = "high",
value: float | None = None,
threshold: float | None = None,
metric: str | None = None,
site_id: str = "SITE001",
) -> Alert:
return Alert(
alert_id=alert_id,
source_alert_id=f"ALR-{alert_id}",
site_id=site_id,
source="api_mock",
timestamp=MOMENT,
type=type_alerte,
severity=severity,
message="Alerte de test",
value=value,
threshold=threshold,
metric=metric,
prediction_id=None,
raw_data={},
)
@pytest.mark.parametrize(
("type_alerte", "attendue"),
[
("spike", "spike-delestage-v1"),
("threshold", "threshold-reduction-v1"),
("outage", "outage-secours-v1"),
("sensor", "sensor-maintenance-v1"),
("anomaly", "anomaly-verification-v1"),
],
ids=["pic", "seuil", "coupure", "capteur", "anomalie"],
)
def test_each_alert_type_yields_its_own_rule(type_alerte: str, attendue: str) -> None:
proposees = applique_les_regles(alerte(type_alerte=type_alerte))
assert [p.rule_reference for p in proposees] == [attendue]
def test_a_critical_alert_adds_the_escalation_rule() -> None:
proposees = applique_les_regles(alerte(severity="critical"))
assert "escalade-astreinte-v1" in {p.rule_reference for p in proposees}
@pytest.mark.parametrize("severity", ["low", "medium", "high"], ids=["faible", "moyenne", "haute"])
def test_a_non_critical_alert_does_not_escalate(severity: str) -> None:
proposees = applique_les_regles(alerte(severity=severity))
assert "escalade-astreinte-v1" not in {p.rule_reference for p in proposees}
def test_a_large_overshoot_adds_the_contract_rule() -> None:
proposees = applique_les_regles(
alerte(value=720.0 * FACTEUR_DEPASSEMENT_MAJEUR, threshold=720.0)
)
assert "contrat-puissance-v1" in {p.rule_reference for p in proposees}
def test_an_overshoot_below_the_factor_does_not_add_the_contract_rule() -> None:
proposees = applique_les_regles(alerte(value=800.0, threshold=720.0))
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
@pytest.mark.parametrize(
("value", "threshold"),
[(None, 720.0), (900.0, None), (900.0, 0.0), (900.0, -10.0)],
ids=["sans mesure", "sans seuil", "seuil nul", "seuil negatif"],
)
def test_the_contract_rule_stays_silent_without_an_exploitable_threshold(
value: float | None, threshold: float | None
) -> None:
proposees = applique_les_regles(alerte(value=value, threshold=threshold))
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
def test_the_explanation_quotes_the_measure_and_the_threshold() -> None:
proposees = applique_les_regles(alerte(value=812.5, threshold=720.0, metric="consumption_kw"))
assert "(consumption_kw mesurée à 812.5, seuil 720.0)" in proposees[0].explanation
def test_the_explanation_quotes_the_measure_alone_when_no_threshold_is_known() -> None:
proposees = applique_les_regles(alerte(value=812.5, metric="consumption_kw"))
assert "(consumption_kw mesurée à 812.5)" in proposees[0].explanation
def test_the_explanation_omits_the_measure_when_the_alert_carries_none() -> None:
proposees = applique_les_regles(alerte())
assert "(" not in proposees[0].explanation
def test_the_explanation_names_the_site() -> None:
proposees = applique_les_regles(alerte(site_id="SITE042"))
assert "SITE042" in proposees[0].explanation
def test_every_proposal_carries_the_alert_identifier() -> None:
proposees = applique_les_regles(alerte(alert_id=77, severity="critical"))
assert {p.alert_id for p in proposees} == {77}
def test_an_alert_never_yields_the_same_rule_twice() -> None:
proposees = applique_les_regles(
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
)
assert len(proposees) == len({p.rule_reference for p in proposees})
def test_a_critical_alert_over_the_threshold_yields_the_three_rules() -> None:
proposees = applique_les_regles(
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
)
assert {p.rule_reference for p in proposees} == {
"spike-delestage-v1",
"escalade-astreinte-v1",
"contrat-puissance-v1",
}
-27
View File
@@ -118,30 +118,3 @@ def test_main_exports_the_contract_without_asking_for_a_password(
assert code == 0 assert code == 0
assert destination.exists() assert destination.exists()
assert str(destination) in capsys.readouterr().out assert str(destination) in capsys.readouterr().out
def test_build_parser_reads_the_generate_recommendations_arguments() -> None:
arguments = cli.build_parser().parse_args(["generate-recommendations", "--site-id", "SITE002"])
assert arguments.commande == "generate-recommendations"
assert arguments.site_id == "SITE002"
def test_build_parser_defaults_the_generation_to_every_site() -> None:
arguments = cli.build_parser().parse_args(["generate-recommendations"])
assert arguments.site_id is None
def test_main_generates_the_recommendations_without_asking_for_a_password(
monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str]
) -> None:
async def fausse_generation(*, site_id: str | None) -> str:
return f"génération lancée pour {site_id}"
monkeypatch.setattr(cli, "generate_recommendations", fausse_generation)
code = cli.main(["generate-recommendations", "--site-id", "SITE002"])
assert code == 0
assert "SITE002" in capsys.readouterr().out
@@ -1,87 +0,0 @@
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()
+7
View File
@@ -23,4 +23,11 @@ export const routes: Routes = [
loadComponent: () => loadComponent: () =>
import('./features/sites/site-detail/site-detail').then((m) => m.SiteDetail), import('./features/sites/site-detail/site-detail').then((m) => m.SiteDetail),
}, },
{
path: 'monitoring/sensors',
canActivate: [authGuard],
data: { role: 'admin' },
loadComponent: () =>
import('./features/monitoring/sensor-status/sensor-status').then((m) => m.SensorStatusView),
},
]; ];
@@ -0,0 +1,48 @@
import { TestBed } from '@angular/core/testing';
import { provideHttpClient } from '@angular/common/http';
import { provideHttpClientTesting, HttpTestingController } from '@angular/common/http/testing';
import { SensorsService } from './sensors.service';
import { environment } from '../../../environments/environment';
describe('SensorsService', () => {
let service: SensorsService;
let httpMock: HttpTestingController;
beforeEach(() => {
TestBed.configureTestingModule({
providers: [provideHttpClient(), provideHttpClientTesting()],
});
service = TestBed.inject(SensorsService);
httpMock = TestBed.inject(HttpTestingController);
});
afterEach(() => httpMock.verify());
it("appelle l'endpoint /sensors/status et retourne la réponse", () => {
let result: unknown;
service.getStatus().subscribe((r) => (result = r));
const req = httpMock.expectOne(`${environment.apiUrl}/sensors/status`);
expect(req.request.method).toBe('GET');
req.flush({
timestamp: '2026-09-18T08:00:00',
sites: [
{
site_id: 'SITE001',
site_name: 'Test',
overall: 'ok',
sensors: {
consumption: { status: 'ok', since: null },
electrical: { status: 'ok', since: null },
temperature: { status: 'ok', since: null },
humidity: { status: 'ok', since: null },
network: { status: 'ok', since: 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 {SensorStatusResponse} from '../../shared/models/sensor-status.model';
@Service()
export class SensorsService {
private http = inject(HttpClient);
getStatus() {
return this.http.get<SensorStatusResponse>(`${environment.apiUrl}/sensors/status`);
}
}
@@ -10,6 +10,9 @@
</div> </div>
</div> </div>
<div class="dashboard__actions"> <div class="dashboard__actions">
@if (auth.principal()?.role === 'admin') {
<a routerLink="/monitoring/sensors" class="ev-link">Supervision des capteurs</a>
}
<a routerLink="/sites" class="ev-link">Voir les sites</a> <a routerLink="/sites" class="ev-link">Voir les sites</a>
<ev-button <ev-button
class="logout-button" class="logout-button"
@@ -68,6 +68,15 @@ h2 {
text-align: center; text-align: center;
} }
.card--link {
cursor: pointer;
transition: border-color 0.15s ease;
&:hover {
border-color: var(--color-primary);
}
}
.card__label { .card__label {
font-size: 0.8rem; font-size: 0.8rem;
color: var(--color-text-muted); color: var(--color-text-muted);
@@ -170,8 +170,11 @@ describe('Dashboard', () => {
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', () => {
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([])) }; const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([])) };
const authMock = { logout: vi.fn().mockReturnValue(of(undefined)), clearSession: vi.fn() }; const authMock = {
logout: vi.fn().mockReturnValue(of(undefined)),
clearSession: vi.fn(),
principal: vi.fn().mockReturnValue({ role: 'admin' }),
};
TestBed.configureTestingModule({ TestBed.configureTestingModule({
imports: [Dashboard], imports: [Dashboard],
providers: [ providers: [
@@ -198,9 +201,10 @@ describe('Dashboard', () => {
it('déconnecte localement et redirige vers /login même si logout échoue côté réseau', () => { it('déconnecte localement et redirige vers /login même si logout échoue côté réseau', () => {
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([])) }; const alertsMock = { getAlerts: vi.fn().mockReturnValue(of([])) };
const authMock = { const authMock = {
logout: vi.fn().mockReturnValue(throwError(() => new Error('réseau indisponible'))), logout: vi.fn().mockReturnValue(throwError(() => new Error('réseau indisponible'))),
clearSession: vi.fn(), clearSession: vi.fn(),
principal: vi.fn().mockReturnValue({ role: 'admin' }),
}; };
TestBed.configureTestingModule({ TestBed.configureTestingModule({
imports: [Dashboard], imports: [Dashboard],
@@ -59,8 +59,8 @@ const TON_PAR_STATUT_PREDICTION: Record<PredictionStatus, 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);
public auth = inject(AuthService);
private predictionsService = inject(PredictionsService); private predictionsService = inject(PredictionsService);
private auth = inject(AuthService);
private router = inject(Router); private router = inject(Router);
private destroyRef = inject(DestroyRef); private destroyRef = inject(DestroyRef);
@@ -0,0 +1,49 @@
<div class="sensor-status">
<nav class="ev-breadcrumb">
<a routerLink="/dashboard">Tableau de bord</a>
</nav>
<header class="sensor-status__header">
<a routerLink="/dashboard" class="ev-brand-link">
<ev-brand class="sensor-status__logo" />
</a>
<div>
<h1>Supervision des capteurs</h1>
<p class="sensor-status__subtitle">État de santé par capteur et par site</p>
</div>
</header>
@if (error(); as message) {
<ev-alert severity="danger" class="banner-error">{{ message }}</ev-alert>
}
@if (data(); as d) {
<div class="sites-grid">
@for (site of d.sites; track site.site_id) {
<ev-card class="site-card">
<div class="site-card__header">
<span class="site-card__name">{{ site.site_name }}</span>
<ev-badge [tone]="badgeToneForOverall(site.overall)">{{ site.overall }}</ev-badge>
</div>
<ul class="sensor-list">
@for (entry of sensorEntries; track entry[0]) {
<li class="sensor-item">
<span
class="sensor-dot"
[class]="'sensor-dot--' + sensorOf(site.sensors, entry[0]).status"
></span>
<span class="sensor-item__label">{{ entry[1] }}</span>
@if (sensorOf(site.sensors, entry[0]).status === 'failing') {
<span class="sensor-item__since">
depuis {{ sensorOf(site.sensors, entry[0]).since | date: 'short' }}
</span>
}
</li>
}
</ul>
</ev-card>
}
</div>
}
</div>
@@ -0,0 +1,90 @@
:host {
display: block;
color: var(--color-text);
padding: 2.5rem 2rem;
max-width: 1100px;
margin: 0 auto;
}
.sensor-status__header {
display: flex;
align-items: center;
gap: 0.85rem;
margin-bottom: 2rem;
h1 {
margin: 0;
font-size: 1.75rem;
font-weight: 700;
}
}
.sensor-status__logo {
font-size: 1.3rem;
}
.sensor-status__subtitle {
margin: 0.25rem 0 0;
color: var(--color-text-muted);
}
.banner-error {
display: block;
margin: 0 0 1.5rem;
}
.sites-grid {
display: grid;
grid-template-columns: repeat(auto-fit, minmax(260px, 1fr));
gap: 1rem;
}
.site-card__header {
display: flex;
align-items: center;
justify-content: space-between;
margin-bottom: 0.75rem;
}
.site-card__name {
font-weight: 600;
}
.sensor-list {
list-style: none;
margin: 0;
padding: 0;
display: flex;
flex-direction: column;
gap: 0.5rem;
}
.sensor-item {
display: flex;
align-items: center;
gap: 0.5rem;
font-size: 0.85rem;
}
.sensor-dot {
width: 8px;
height: 8px;
border-radius: 50%;
flex-shrink: 0;
&--ok {
background: var(--color-success);
}
&--failing {
background: var(--color-danger);
}
}
.sensor-item__label {
flex: 1;
}
.sensor-item__since {
color: var(--color-text-muted);
font-size: 0.75rem;
}
@@ -0,0 +1,99 @@
import { TestBed } from '@angular/core/testing';
import { of, throwError } from 'rxjs';
import { vi } from 'vitest';
import { SensorStatusView } from './sensor-status';
import { SensorsService } from '../../../core/services/sensors.service';
import { SiteSensors } from '../../../shared/models/sensor-status.model';
import {provideRouter} from '@angular/router';
const OK_SENSORS: SiteSensors = {
consumption: { status: 'ok', since: null },
electrical: { status: 'ok', since: null },
temperature: { status: 'ok', since: null },
humidity: { status: 'ok', since: null },
network: { status: 'ok', since: null },
};
describe('SensorStatusView', () => {
let sensorsMock: { getStatus: ReturnType<typeof vi.fn> };
beforeEach(() => {
sensorsMock = { getStatus: vi.fn() };
TestBed.configureTestingModule({
imports: [SensorStatusView],
providers: [
{ provide: SensorsService, useValue: sensorsMock },
provideRouter([]),
],
});
});
it('charge et affiche les données au démarrage', () => {
sensorsMock.getStatus.mockReturnValue(
of({
timestamp: '2026-09-18T08:00:00',
sites: [
{ site_id: 'SITE001', site_name: 'Bureau Test', overall: 'ok', sensors: OK_SENSORS },
],
})
);
const fixture = TestBed.createComponent(SensorStatusView);
fixture.detectChanges();
expect(fixture.componentInstance.data()?.sites.length).toBe(1);
expect(fixture.componentInstance.error()).toBeNull();
expect(fixture.nativeElement.textContent).toContain('Bureau Test');
});
it("affiche un message d'erreur si l'appel échoue", () => {
sensorsMock.getStatus.mockReturnValue(throwError(() => new Error('boom')));
const fixture = TestBed.createComponent(SensorStatusView);
fixture.detectChanges();
expect(fixture.componentInstance.error()).toBe(
'État des capteurs indisponible, réessayez plus tard.'
);
expect(fixture.componentInstance.data()).toBeNull();
expect(fixture.nativeElement.textContent).toContain('État des capteurs indisponible');
});
it('associe le bon ton de badge à chaque statut global', () => {
sensorsMock.getStatus.mockReturnValue(of({ timestamp: '2026-09-18T08:00:00', sites: [] }));
const fixture = TestBed.createComponent(SensorStatusView);
const component = fixture.componentInstance;
expect(component.badgeToneForOverall('ok')).toBe('success');
expect(component.badgeToneForOverall('degraded')).toBe('warning');
expect(component.badgeToneForOverall('critical')).toBe('critical');
expect(component.badgeToneForOverall('inconnu')).toBe('neutral');
});
it('retourne le bon diagnostic via sensorOf', () => {
sensorsMock.getStatus.mockReturnValue(of({ timestamp: '2026-09-18T08:00:00', sites: [] }));
const fixture = TestBed.createComponent(SensorStatusView);
const component = fixture.componentInstance;
expect(component.sensorOf(OK_SENSORS, 'temperature')).toEqual({ status: 'ok', since: null });
});
it('affiche la date depuis quand un capteur est en panne', () => {
const sensors: SiteSensors = {
...OK_SENSORS,
temperature: { status: 'failing', since: '2026-09-18T08:00:00' },
};
sensorsMock.getStatus.mockReturnValue(
of({
timestamp: '2026-09-18T08:00:00',
sites: [{ site_id: 'SITE001', site_name: 'Bureau Test', overall: 'degraded', sensors }],
})
);
const fixture = TestBed.createComponent(SensorStatusView);
fixture.detectChanges();
expect(fixture.nativeElement.textContent).toContain('depuis');
});
});
@@ -0,0 +1,65 @@
import { Component, OnInit, inject, signal } from '@angular/core';
import { RouterLink } from '@angular/router';
import { catchError, EMPTY, Observable } from 'rxjs';
import {Badge, BadgeTone} from '../../../shared/components/ui/badge/badge';
import {Card} from '../../../shared/components/ui/card/card';
import {Alert} from '../../../shared/components/ui/alert/alert';
import {Brand} from '../../../shared/components/ui/brand/brand';
import {SensorsService} from '../../../core/services/sensors.service';
import {SensorDiagnostic, SensorStatusResponse} from '../../../shared/models/sensor-status.model';
import { DatePipe } from '@angular/common';
const UNAVAILABLE_MESSAGE = 'État des capteurs indisponible, réessayez plus tard.';
const SENSOR_LABELS: Record<string, string> = {
consumption: 'Consommation',
electrical: 'Électrique',
temperature: 'Température',
humidity: 'Humidité',
network: 'Réseau',
};
const TON_PAR_OVERALL: Record<string, BadgeTone> = {
ok: 'success',
degraded: 'warning',
critical: 'critical',
};
@Component({
selector: 'app-sensor-status',
standalone: true,
imports: [RouterLink, Card, Alert, Badge, Brand, DatePipe],
templateUrl: './sensor-status.html',
styleUrl: './sensor-status.scss',
})
export class SensorStatusView implements OnInit {
private sensorsService = inject(SensorsService);
data = signal<SensorStatusResponse | null>(null);
error = signal<string | null>(null);
readonly sensorEntries = Object.entries(SENSOR_LABELS);
ngOnInit(): void {
this.sensorsService
.getStatus()
.pipe(catchError(() => this.reportUnavailable()))
.subscribe((response) => {
this.error.set(null);
this.data.set(response);
});
}
sensorOf(sensors: Record<string, SensorDiagnostic>, key: string): SensorDiagnostic {
return sensors[key];
}
badgeToneForOverall(overall: string): BadgeTone {
return TON_PAR_OVERALL[overall] ?? 'neutral';
}
private reportUnavailable(): Observable<never> {
this.error.set(UNAVAILABLE_MESSAGE);
return EMPTY;
}
}
@@ -0,0 +1,28 @@
export type SensorStatus = 'ok' | 'failing';
export type OverallStatus = 'ok' | 'degraded' | 'critical';
export interface SensorDiagnostic {
status: SensorStatus;
since: string | null;
}
export interface SiteSensors {
consumption: SensorDiagnostic;
electrical: SensorDiagnostic;
temperature: SensorDiagnostic;
humidity: SensorDiagnostic;
network: SensorDiagnostic;
[key: string]: SensorDiagnostic;
}
export interface SiteSensorStatus {
site_id: string;
site_name: string;
sensors: SiteSensors;
overall: OverallStatus;
}
export interface SensorStatusResponse {
timestamp: string;
sites: SiteSensorStatus[];
}
@@ -1,81 +0,0 @@
# 0006 - Le moteur de règles de recommandation vit dans le backend
- Statut : accepté
- Date : 2026-09-18
## Contexte
L'issue #38 demande un « moteur de règles pour recommandations », portée par le label `ml`. Le
schéma tranche déjà la forme du résultat : `recommendation(alert_id, action, explanation,
rule_reference)`, avec `alert_id` en clé étrangère `NOT NULL` et une contrainte d'unicité
`uq_recommendation_alert_rule` sur `(alert_id, rule_reference)`. Une recommandation est donc
**dérivée d'une alerte**, jamais d'une mesure brute ni d'une prévision.
Deux emplacements se disputaient le code :
1. `ml/enervision_ml/`, sur le patron de `enervision_ml.score` livré par #37 : un script autonome
qui se connecte par `ML_DATABASE_URL`, écrit une table, et que l'API se contente de lire.
L'[ADR 0005](0005-modele-prediction-lightgbm.md) annonce d'ailleurs #38 de ce côté, en écrivant
que le scoring, le moteur de recommandations et les tests de dérive « consommeront le même
module `enervision_ml.features` ».
2. `apps/backend/app/services/`, où `apps/backend/README.md` place les « regles metier ».
## Décision
**Le moteur vit dans `apps/backend/app/services/`**, sous la forme d'un module pur
`recommendation_rules.py` (le catalogue `REGLES`) et d'une méthode `RecommendationService.generate()`
qui l'applique, persiste et valide la transaction.
Trois raisons :
- **Il n'utilise rien du ML.** Le catalogue lit `alert.type`, `alert.severity`, `alert.value` et
`alert.threshold`. Aucun modèle, aucune feature, aucun `enervision_ml.features` : la phrase de
l'ADR 0005 vaut pour le scoring (#37) et les tests de dérive (#44/#45), qui manipulent bien des
features, pas pour des règles sur alertes. Le label `ml` de #38 désigne le lot fonctionnel
« prédiction et recommandation », pas l'emplacement du code.
- **Il lit et écrit deux tables déjà couvertes par des repositories.** `AlertRepository` sait déjà
filtrer par site. Le placer dans `ml/` obligerait à réécrire ces accès en SQL brut, et à
maintenir deux représentations du même domaine.
- **Le déclencheur HTTP n'a de sens que dans l'API.** `POST /recommendations/generate` doit passer
par `require_role(Role.ADMIN)` et par la session injectée : cela suppose d'être dans
l'application FastAPI.
Le moteur reste néanmoins **déclenchable hors HTTP**, par `python -m app.cli
generate-recommendations` (cible `make recommendations`), sur le patron de `make ml-score` : rien
n'oblige à exposer un port pour régénérer des recommandations.
## Conséquences
- L'API gagne sa première route d'écriture métier. La checklist de `20-backend.md` s'applique :
entrée dans `ROLE_MINIMUM` de `tests/api/acces.py`, et `openapi.json` régénéré dans le même
commit.
- `RecommendationService` n'est plus en lecture seule : il reçoit le `Transaction` Protocol déjà
utilisé par `AuthService` et `UserService`, et commite lui-même. Les repositories continuent de
ne pas commiter.
- **L'idempotence est déléguée à la base.** `create_missing()` insère en `ON CONFLICT DO NOTHING`
sur `uq_recommendation_alert_rule` plutôt que de relire avant d'écrire, ce qui supprime la
fenêtre entre le contrôle et l'insertion. Corollaire : `rule_reference` est une clé fonctionnelle.
Une règle dont le sens change prend une référence `-v2` ; renommer une référence livrée
ferait réapparaître ses recommandations à côté des anciennes.
- **Le moteur est branché sur la détection interne, et sur elle seule.** `alert` est alimentée
par `app/detection/internal_alerts.py` (#104), lancée à la main comme `enervision_ml.score` ;
l'ingestion de l'API Mock `/alerts` reste à faire. Le rapport de génération est donc à zéro tant
que la détection n'a pas tourné, sans que le moteur soit à retoucher.
- **L'insertion est découpée en lots.** `create_missing()` écrit par paquets de `TAILLE_DE_LOT`
lignes : asyncpg plafonne une requête à 32 767 paramètres, soit 8 191 lignes de quatre colonnes,
et la détection interne peut alimenter `alert` au fil de l'eau.
- Si le projet devait un jour pondérer les recommandations par un score appris, la décision serait
à rouvrir : le moteur redeviendrait consommateur du pipeline ML.
## Alternatives écartées
- **Module et CLI dans `ml/enervision_ml/`** : cohérent avec le label `ml` et avec la lettre de
l'ADR 0005, mais impose du SQL brut là où deux repositories existent, et laisse la génération
hors de portée de l'API. Redeviendrait le bon choix si les règles se mettaient à consommer des
features ou un modèle.
- **Génération à la volée, sans persistance**, calculée à chaque `GET /recommendations` : supprime
le besoin d'écriture, mais rend la table `recommendation` et sa contrainte d'unicité inutiles,
et interdit toute trace de ce qui a été proposé et quand.
- **Table de configuration des règles en base**, plutôt qu'un catalogue en Python : plus souple,
mais déplace la logique métier hors de la revue de code et hors des tests, pour un besoin que
rien n'exprime à ce stade.
-5
View File
@@ -165,8 +165,3 @@ Elles vivent dans `../adr/`, pas ici.
| ADR | Objet | | ADR | Objet |
|---|---| |---|---|
| [0001](../adr/0001-postgresql-timescaledb.md) | PostgreSQL 17 avec l'extension TimescaleDB, et la frontière `db/` vs `alembic/` | | [0001](../adr/0001-postgresql-timescaledb.md) | PostgreSQL 17 avec l'extension TimescaleDB, et la frontière `db/` vs `alembic/` |
| [0002](../adr/0002-authentification-jwt-et-refresh-opaque.md) | Authentification par JWT d'accès et jeton de rafraîchissement opaque |
| [0003](../adr/0003-autorisation-rbac-a-trois-roles.md) | Autorisation RBAC à trois rôles, avec relecture du compte à chaque requête |
| [0004](../adr/0004-journal-d-audit-en-ajout-seul.md) | Journal d'audit en ajout seul, garanti par PostgreSQL |
| [0005](../adr/0005-modele-prediction-lightgbm.md) | Modèle de prédiction de consommation : LightGBM |
| [0006](../adr/0006-moteur-de-regles-dans-le-backend.md) | Le moteur de règles de recommandation vit dans le backend, pas dans `ml/` |
-56
View File
@@ -146,7 +146,6 @@ Deux fichiers d'environnement, deux usages : `.env` à la racine alimente `docke
| GET | `/api/v1/alerts` | Liste les alertes, filtrable par `site_id` et `severity`. `lecteur` | 401, 403, 422, 500 | | GET | `/api/v1/alerts` | Liste les alertes, filtrable par `site_id` et `severity`. `lecteur` | 401, 403, 422, 500 |
| GET | `/api/v1/recommendations` | Liste les recommandations. `lecteur` | 401, 403, 500 | | GET | `/api/v1/recommendations` | Liste les recommandations. `lecteur` | 401, 403, 500 |
| GET | `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation. `lecteur` | 401, 403, 404, 422, 500 | | GET | `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation. `lecteur` | 401, 403, 404, 422, 500 |
| POST | `/api/v1/recommendations/generate` | Applique le moteur de règles aux alertes, filtrable par `site_id`. `admin` | 401, 403, 422, 500 |
| 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 |
@@ -195,21 +194,6 @@ plutôt qu'un statut inventé : le domaine `available`/`insufficient_data`/`erro
LightGBM elle-même ; elle lit ce que le pipeline de scoring a déjà écrit, cf. 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. [ML-START.md](../../ML-START.md) section 3.
`POST /recommendations/generate` est la seule route d'écriture métier du contrat. Elle applique
le moteur de règles d'`app/services/recommendation_rules.py` aux lignes d'`alert`, sans modèle ni
feature ML : le catalogue `REGLES` associe à chaque type et à chaque gravité d'alerte une action et
son explication, et une même alerte peut en déclencher plusieurs, comme le prévoit
[40-data.md](40-data.md). L'idempotence est portée par la base, pas par le service :
`RecommendationRepository.create_missing()` insère en `ON CONFLICT DO NOTHING` sur
`uq_recommendation_alert_rule`, donc rejouer la génération sur les mêmes alertes ne crée rien et
le rapport rendu distingue `recommendations_created` de `already_present`. Le même traitement est
disponible hors HTTP par `python -m app.cli generate-recommendations` (cible `make
recommendations`), sur le patron de `make ml-score`. Le choix de loger le moteur dans le backend
plutôt que dans `ml/` est justifié par l'[ADR 0006](../adr/0006-moteur-de-regles-dans-le-backend.md).
Les alertes traitées sont celles qu'écrit la détection interne (#104, section ci-dessous) : la
génération ne rend donc de recommandations qu'une fois la détection passée. L'insertion est
découpée en lots de `TAILLE_DE_LOT` lignes, asyncpg plafonnant une requête à 32 767 paramètres.
`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
@@ -223,46 +207,6 @@ 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
-6
View File
@@ -224,12 +224,6 @@ Les anomalies historiques décrites dans les JSON sont conservées
dans `dataset.metadata`. Elles servent à l’analyse des données dans `dataset.metadata`. Elles servent à l’analyse des données
et ne sont pas considérées comme des alertes actuelles. et ne sont pas considérées comme des alertes actuelles.
Les lignes de `recommendation` sont écrites par le moteur de règles du backend
(`app/services/recommendation_rules.py`), déclenché par `POST /api/v1/recommendations/generate`
ou par `make recommendations`, à partir des alertes déjà en base. Le couple
`(alert_id, rule_reference)` est unique : rejouer le moteur sur les mêmes alertes n'ajoute aucune
ligne.
### Relations entre les tables ### Relations entre les tables
- Un site possède plusieurs mesures, prévisions et alertes. - Un site possède plusieurs mesures, prévisions et alertes.