Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f9c2a4610c | ||
|
|
b5fa7b0010 | ||
|
|
7f710c9084 |
@@ -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)
|
||||||
|
|
||||||
|
|||||||
@@ -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 |
|
||||||
|
|
||||||
|
|||||||
@@ -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)]
|
||||||
|
|||||||
@@ -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,
|
|
||||||
)
|
|
||||||
|
|||||||
@@ -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())
|
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
|
||||||
|
|||||||
@@ -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,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
@@ -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",
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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
|
|
||||||
|
|||||||
@@ -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",
|
|
||||||
}
|
|
||||||
@@ -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()
|
|
||||||
@@ -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.
|
|
||||||
@@ -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/` |
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
Reference in New Issue
Block a user