feat(backend): moteur de règles de recommandations et route de génération
`recommendation` n'avait aucun écrivain : les quatre couches de lecture étaient livrées, mais rien ne produisait de ligne. Le moteur comble ce trou. Le catalogue `REGLES` vit dans `app/services/`, pas dans `ml/` : il lit `alert.type`, `alert.severity`, `alert.value` et `alert.threshold`, sans modèle ni feature, et s'appuie sur deux repositories existants. L'arbitrage avec l'ADR 0005, qui annonçait #38 du côté ML, est tranché par l'ADR 0006. Sept règles, cinq par type d'alerte et deux transverses (sévérité critique, dépassement d'au moins 20 % du seuil), donc une à trois recommandations par alerte. L'idempotence est portée par la base : `create_missing()` insère en `ON CONFLICT DO NOTHING` sur `uq_recommendation_alert_rule`, ce qui supprime la fenêtre entre un contrôle préalable et l'insertion. `rule_reference` devient de ce fait une clé fonctionnelle, d'où le suffixe de version sur chaque référence. Deux déclencheurs : `POST /api/v1/recommendations/generate` réservé `admin`, et `python -m app.cli generate-recommendations` (cible `make recommendations`). Limite connue : aucune source n'alimente `alert` aujourd'hui, ni détection interne (#104) ni ingestion de l'API Mock. La route répond, le rapport reste à zéro, et la chaîne s'allume sans retoucher le moteur le jour où les alertes existent. Tests : 80 unitaires et API verts, plus 6 d'intégration dont l'idempotence jouée contre PostgreSQL. Closes #38
This commit is contained in:
@@ -187,7 +187,11 @@ AlertServiceDep = Annotated[AlertService, Depends(get_alert_service)]
|
||||
|
||||
|
||||
def get_recommendation_service(session: SessionDep) -> RecommendationService:
|
||||
return RecommendationService(recommendations=RecommendationRepository(session))
|
||||
return RecommendationService(
|
||||
recommendations=RecommendationRepository(session),
|
||||
alerts=AlertRepository(session),
|
||||
transaction=session,
|
||||
)
|
||||
|
||||
|
||||
RecommendationServiceDep = Annotated[RecommendationService, Depends(get_recommendation_service)]
|
||||
|
||||
@@ -63,7 +63,7 @@ TAGS: Final[list[dict[str, Any]]] = [
|
||||
"name": "recommendations",
|
||||
"description": (
|
||||
"Consultation des recommandations issues des alertes. Accessible à partir du rôle "
|
||||
"`lecteur`."
|
||||
"`lecteur`. Leur génération par le moteur de règles est réservée au rôle `admin`."
|
||||
),
|
||||
},
|
||||
{
|
||||
|
||||
@@ -1,13 +1,18 @@
|
||||
from fastapi import APIRouter, HTTPException, status
|
||||
|
||||
from app.api.deps import LecteurDep, RecommendationServiceDep
|
||||
from app.api.openapi import REPONSE_VALIDATION, Reponses
|
||||
from app.api.deps import AdminDep, LecteurDep, RecommendationServiceDep
|
||||
from app.api.openapi import REPONSE_VALIDATION, REPONSES_ADMIN, Reponses
|
||||
from app.schemas.errors import ErrorResponse
|
||||
from app.schemas.recommendation import RecommendationResponse
|
||||
from app.schemas.recommendation import (
|
||||
RecommendationGenerationResponse,
|
||||
RecommendationResponse,
|
||||
)
|
||||
from app.services.recommendation import RecommendationNotFoundError
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
REPONSES_GENERATION: Reponses = {**REPONSES_ADMIN, **REPONSE_VALIDATION}
|
||||
|
||||
REPONSES_INTROUVABLE: Reponses = {
|
||||
**REPONSE_VALIDATION,
|
||||
404: {"model": ErrorResponse, "description": "Aucune recommandation ne porte cet identifiant."},
|
||||
@@ -38,3 +43,22 @@ async def get_recommendation(
|
||||
status_code=status.HTTP_404_NOT_FOUND, detail="Recommandation introuvable"
|
||||
) from erreur
|
||||
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,8 +22,11 @@ from app.core.hashing import build_hasher
|
||||
from app.core.roles import Role
|
||||
from app.db.session import get_session_factory
|
||||
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.schemas.auth import PASSWORD_MIN_LENGTH, SPECIAL_CHARACTERS, valide_complexite
|
||||
from app.services.recommendation import RecommendationService
|
||||
|
||||
LONGUEUR_MOT_DE_PASSE_GENERE = 24
|
||||
CHEMIN_CONTRAT = Path(__file__).resolve().parent.parent / "openapi.json"
|
||||
@@ -63,6 +66,22 @@ 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
|
||||
# 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.
|
||||
@@ -109,6 +128,14 @@ def build_parser() -> argparse.ArgumentParser:
|
||||
"export-openapi", help="Écrit le contrat OpenAPI sur disque"
|
||||
)
|
||||
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
|
||||
|
||||
|
||||
@@ -152,6 +179,10 @@ def main(argv: list[str] | None = None) -> int:
|
||||
print(export_openapi(Path(arguments.output)))
|
||||
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)
|
||||
|
||||
succes, message = asyncio.run(
|
||||
|
||||
@@ -1,11 +1,21 @@
|
||||
from collections.abc import Sequence
|
||||
from dataclasses import asdict, dataclass
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.dialects.postgresql import insert
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.energy import Recommendation
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class NouvelleRecommandation:
|
||||
alert_id: int
|
||||
action: str
|
||||
explanation: str
|
||||
rule_reference: str
|
||||
|
||||
|
||||
class RecommendationRepository:
|
||||
def __init__(self, session: AsyncSession) -> None:
|
||||
self._session = session
|
||||
@@ -20,3 +30,18 @@ class RecommendationRepository:
|
||||
)
|
||||
recommendation: Recommendation | None = await self._session.scalar(requete)
|
||||
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:
|
||||
if not nouvelles:
|
||||
return 0
|
||||
|
||||
requete = (
|
||||
insert(Recommendation)
|
||||
.values([asdict(nouvelle) for nouvelle in nouvelles])
|
||||
.on_conflict_do_nothing(constraint="uq_recommendation_alert_rule")
|
||||
.returning(Recommendation.recommendation_id)
|
||||
)
|
||||
creees = (await self._session.scalars(requete)).all()
|
||||
return len(creees)
|
||||
|
||||
@@ -12,3 +12,9 @@ class RecommendationResponse(BaseModel):
|
||||
explanation: str
|
||||
rule_reference: str
|
||||
created_at: datetime
|
||||
|
||||
|
||||
class RecommendationGenerationResponse(BaseModel):
|
||||
alerts_examined: int
|
||||
recommendations_created: int
|
||||
already_present: int
|
||||
|
||||
@@ -1,7 +1,15 @@
|
||||
from collections.abc import Sequence
|
||||
from dataclasses import dataclass
|
||||
from typing import Protocol
|
||||
|
||||
from app.models.energy import Recommendation
|
||||
from app.repositories.alert import AlertRepository
|
||||
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):
|
||||
@@ -12,9 +20,24 @@ class RecommendationNotFoundError(RecommendationError):
|
||||
pass
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class RapportGeneration:
|
||||
alertes_examinees: int
|
||||
recommandations_creees: int
|
||||
deja_presentes: int
|
||||
|
||||
|
||||
class RecommendationService:
|
||||
def __init__(self, *, recommendations: RecommendationRepository) -> None:
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
recommendations: RecommendationRepository,
|
||||
alerts: AlertRepository,
|
||||
transaction: Transaction,
|
||||
) -> None:
|
||||
self._recommendations = recommendations
|
||||
self._alerts = alerts
|
||||
self._transaction = transaction
|
||||
|
||||
async def list_all(self) -> Sequence[Recommendation]:
|
||||
return await self._recommendations.list_all()
|
||||
@@ -24,3 +47,16 @@ class RecommendationService:
|
||||
if recommendation is None:
|
||||
raise RecommendationNotFoundError(recommendation_id)
|
||||
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,
|
||||
)
|
||||
|
||||
@@ -0,0 +1,117 @@
|
||||
# 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})"
|
||||
Reference in New Issue
Block a user