Compare commits

...
Author SHA1 Message Date
Johan LEROY 0a2ed5ad8f docs: distingue l'ingestion des mesures de celle des alertes
La #114 arrive sur dev avec un ADR 0006 qui note que l'ingestion de
l'API Mock /alerts reste à faire. « Les deux sources sont implémentées »
se lisait comme couvrant aussi les alertes.
2026-09-21 10:25:48 +02:00
Johan LEROY 459ddf1792 Merge remote-tracking branch 'origin/dev' into feat/mock-api-import 2026-09-21 10:24:48 +02:00
Johan LEROY de697b080d docs: rétablit la hiérarchie des titres et les motifs du document Data
40-data.md était passé à quatre titres de niveau 1 et etl/README.md à cinq,
alors que les huit autres documents d'architecture n'en ont qu'un. Les
sections ajoutées redescendent d'un niveau.

La réécriture de la section « Tables d'authentification » avait aussi vidé
quatre choix de modélisation de leur raison, dont le renvoi à l'ADR 0004 sur
audit_log.actor_id. Ces motifs sont rétablis, et les deux tables de
réinitialisation reçoivent le leur.

Documente enfin la frontière de confiance avec l'API Mock : les quatre
garde-fous, les plages de PHYSICAL_BOUNDS, et ce qu'il reste à faire.
2026-09-21 10:21:18 +02:00
Johan LEROY f238940867 fix(etl): borne la réponse de l'API Mock avant écriture en base
L'API Mock est le seul item OWASP API10 du projet, et ce script en est le
premier consommateur. Des quatre garde-fous exigés par la traçabilité OWASP,
seul le timeout était en place.

- plafonne la taille des réponses : MAX_SITES sites, au plus --limit mesures ;
- borne chaque grandeur physique par PHYSICAL_BOUNDS, une valeur hors plage,
  d'un type inattendu, NaN ou infinie devenant NULL avec sa raison dans
  null_reasons et data_quality à degraded ;
- ne recopie vers la base que les champs attendus, via build_site_row() et
  build_reading_row(), au lieu de passer les dictionnaires de l'API en
  paramètres SQL ;
- écarte une data_quality que ck_reading_quality refuserait, plutôt que de
  faire échouer le lot entier ;
- nomme la cible du ON CONFLICT, qui avalait jusqu'ici toute violation
  d'unicité, y compris celle de la clé primaire.

raw_data conserve la réponse d'origine intacte : rien n'est perdu, seule son
exploitation est bornée.
2026-09-21 10:21:10 +02:00
Johan LEROYandGitHub 9d2384a639 Merge pull request #114 from ineszang/feat/moteur-regles-recommandations
feat(backend): moteur de règles de recommandations et route de génération
2026-09-21 09:51:34 +02:00
Johan LEROY 19c38fe571 fix(backend): decoupe l'insertion des recommandations en lots et remet les docs a jour
Backend / Tests exigeant une base (push) Failing after 34s
Backend / Lint, typage et tests (push) Successful in 1m24s
Backend / Audit des dépendances (push) Successful in 57s
SonarQube / build-back (push) Successful in 1m5s
SonarQube / build-front (push) Successful in 9m39s
SonarQube / test-back (push) Failing after 51s
SonarQube / test-front (push) Failing after 5m6s
SonarQube / SonarQube (push) Skipped
`create_missing()` construisait un seul `INSERT ... VALUES` pour la totalite des
propositions. Avec quatre colonnes par ligne et le plafond asyncpg de 32 767
parametres, la route echouait au-dela de 8 191 recommandations par appel, cas
devenu realiste maintenant que la detection interne (#104) alimente `alert` en
continu. L'insertion passe par des lots de `TAILLE_DE_LOT` lignes, sur le patron
de `app/etl/historical_import.py`.

L'ADR 0006, `20-backend.md` et la description de la PR annoncaient qu'aucune
source n'alimentait `alert` et que #104 n'etait pas commencee. #104 est livree
sur `dev` depuis la #113 : les phrases sont corrigees plutot que laissees a
vieillir dans un ADR.
2026-09-21 09:45:25 +02:00
Johan LEROY 9a1af94d88 Merge remote-tracking branch 'origin/dev' into feat/moteur-regles-recommandations 2026-09-21 09:40:46 +02:00
Johan LEROY aeb07e14db 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
2026-09-18 15:49:54 +02:00
25 changed files with 2260 additions and 1014 deletions
+4 -1
View File
@@ -6,7 +6,7 @@ ML := ml
.PHONY: help install install-backend install-frontend install-ml dev dev-backend dev-frontend \
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 \
ml-lint ml-typecheck ml-test ml-check ml-train ml-score
ml-lint ml-typecheck ml-test ml-check ml-train ml-score recommendations
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}'
@@ -77,6 +77,9 @@ 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
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 -t enervision-backend:local $(BACKEND)
+1
View File
@@ -113,6 +113,7 @@ Le sens de dependance est unique : `endpoints` vers `services` vers `repositorie
| `/api/v1/sites/{site_id}` | Décrit un site | `lecteur` |
| `/api/v1/recommendations` | Liste les recommandations | `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` |
| `/docs`, `/openapi.json` | Documentation, fermée en `staging` et `prod` | public sinon |
+5 -1
View File
@@ -192,7 +192,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)]
+1 -1
View File
@@ -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,
)
+31
View File
@@ -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(
+121 -20
View File
@@ -1,3 +1,11 @@
# Contrainte : la réponse de l'API Mock est une entrée hostile, pas une source de confiance.
# Voir OWASP API10 dans docs/architecture/owasp-traceabilite.md. Rien de ce qu'elle renvoie
# n'atteint la base sans passer par build_site_row() ou build_reading_row() : seuls les champs
# attendus sont recopiés, les grandeurs physiques sont bornées par PHYSICAL_BOUNDS et la taille
# des tableaux est plafonnée par MAX_SITES et par --limit. Une valeur hors bornes devient NULL
# et laisse sa trace dans null_reasons plutôt que de lever : le mock émet des anomalies par
# construction, et raw_data conserve de toute façon la réponse d'origine intacte.
from __future__ import annotations
import argparse
@@ -14,6 +22,25 @@ from app.core.config import get_settings
SOURCE_HISTORY = "api_history"
MAX_SITES = 100
MAX_LIMIT = 1000
# Les quatre seules valeurs que la contrainte ck_reading_quality accepte.
ACCEPTED_QUALITIES = frozenset({"good", "partial", "degraded", "critical"})
PHYSICAL_BOUNDS: dict[str, tuple[float, float]] = {
"consumption_kw": (0.0, 100_000.0),
"consumption_kwh": (0.0, 100_000.0),
"voltage_v": (0.0, 1_000.0),
"current_a": (0.0, 10_000.0),
"power_factor": (0.0, 1.0),
"temperature_celsius": (-90.0, 60.0),
"humidity_percent": (0.0, 100.0),
}
CAPACITY_BOUNDS = (0.0, 100_000.0)
def create_mock_api_client() -> httpx.AsyncClient:
settings = get_settings()
@@ -31,6 +58,53 @@ def create_mock_api_client() -> httpx.AsyncClient:
)
def read_text(payload: dict[str, Any], key: str) -> str:
value = payload.get(key)
if not isinstance(value, str) or not value:
raise ValueError(f"Champ {key} absent ou invalide dans la réponse de l'API Mock.")
return value
def optional_text(value: Any) -> str | None:
return value if isinstance(value, str) else None
def coerce_measure(
value: Any,
bounds: tuple[float, float],
) -> float | None:
if isinstance(value, bool) or not isinstance(value, int | float):
return None
lower, upper = bounds
# Écarte aussi NaN et les infinis, qu'aucune comparaison de bornes ne retient.
return float(value) if lower <= value <= upper else None
def resolve_quality(
value: Any,
rejected: list[str],
) -> str | None:
quality = value if isinstance(value, str) and value in ACCEPTED_QUALITIES else None
if rejected:
return "critical" if quality == "critical" else "degraded"
return quality
def resolve_null_reasons(
value: Any,
rejected: list[str],
) -> list[str]:
reported = [str(reason) for reason in value] if isinstance(value, list) else []
return reported + rejected
async def fetch_sites(
client: httpx.AsyncClient,
) -> list[dict[str, Any]]:
@@ -43,14 +117,32 @@ async def fetch_sites(
if not isinstance(payload, list):
raise ValueError("La réponse /api/v1/sites doit être une liste.")
if len(payload) > MAX_SITES:
raise ValueError(f"La réponse /api/v1/sites dépasse le plafond de {MAX_SITES} sites.")
return payload
def build_site_row(
site: dict[str, Any],
) -> dict[str, Any]:
return {
"site_id": read_text(site, "site_id"),
"site_type": read_text(site, "site_type"),
"site_name": read_text(site, "site_name"),
"location": optional_text(site.get("location")),
"capacity_kw": coerce_measure(site.get("capacity_kw"), CAPACITY_BOUNDS),
"status": optional_text(site.get("status")),
}
async def upsert_sites(
connection: AsyncConnection,
sites: list[dict[str, Any]],
) -> None:
if not sites:
rows = [build_site_row(site) for site in sites]
if not rows:
return
await connection.execute(
@@ -81,7 +173,7 @@ async def upsert_sites(
status = EXCLUDED.status
"""
),
sites,
rows,
)
@@ -90,7 +182,7 @@ async def fetch_readings(
site_id: str,
start_time: datetime,
end_time: datetime,
limit: int = 1000,
limit: int = MAX_LIMIT,
) -> list[dict[str, Any]]:
response = await client.get(
"/api/v1/readings",
@@ -109,30 +201,36 @@ async def fetch_readings(
if not isinstance(payload, list):
raise ValueError("La réponse /api/v1/readings doit être une liste.")
if len(payload) > limit:
raise ValueError(f"La réponse /api/v1/readings dépasse la limite demandée de {limit}.")
return payload
def build_reading_row(
reading: dict[str, Any],
) -> dict[str, Any]:
timestamp = datetime.fromisoformat(reading["timestamp"].replace("Z", "+00:00"))
measures: dict[str, float | None] = {}
rejected: list[str] = []
for name, bounds in PHYSICAL_BOUNDS.items():
received = reading.get(name)
measures[name] = coerce_measure(received, bounds)
if received is not None and measures[name] is None:
rejected.append(f"out_of_physical_bounds:{name}")
return {
"site_id": reading["site_id"],
"timestamp": timestamp,
"site_id": read_text(reading, "site_id"),
"timestamp": parse_datetime(read_text(reading, "timestamp")),
"source": SOURCE_HISTORY,
"dataset_id": None,
"consumption_kw": reading.get("consumption_kw"),
"consumption_kwh": reading.get("consumption_kwh"),
**measures,
"consumption_euros": None,
"voltage_v": reading.get("voltage_v"),
"current_a": reading.get("current_a"),
"power_factor": reading.get("power_factor"),
"temperature_celsius": reading.get("temperature_celsius"),
"humidity_percent": reading.get("humidity_percent"),
"solar_irradiance_wm2": None,
"is_working_hours": None,
"data_quality": reading.get("data_quality"),
"null_reasons": reading.get("null_reasons"),
"data_quality": resolve_quality(reading.get("data_quality"), rejected),
"null_reasons": resolve_null_reasons(reading.get("null_reasons"), rejected),
"imputed_values": None,
"imputation_method": None,
"raw_data": json.dumps(
@@ -142,6 +240,8 @@ def build_reading_row(
}
# Le conflit vise l'index unique uq_reading_source plutôt que la table entière : sans cible
# nommée, DO NOTHING avalerait aussi une violation de clé primaire.
READING_INSERT = text(
"""
INSERT INTO reading (
@@ -186,7 +286,8 @@ READING_INSERT = text(
:imputation_method,
CAST(:raw_data AS jsonb)
)
ON CONFLICT DO NOTHING
ON CONFLICT (site_id, timestamp, source, (coalesce(dataset_id, 0)))
DO NOTHING
"""
)
@@ -213,7 +314,7 @@ async def import_mock_api_history(
all_readings: list[dict[str, Any]] = []
for site in sites:
site_id = site["site_id"]
site_id = read_text(site, "site_id")
readings = await fetch_readings(
client=client,
@@ -281,7 +382,7 @@ def parse_args() -> argparse.Namespace:
parser.add_argument(
"--limit",
type=int,
default=1000,
default=MAX_LIMIT,
)
parser.add_argument(
@@ -295,8 +396,8 @@ def parse_args() -> argparse.Namespace:
def main() -> None:
args = parse_args()
if args.limit < 1 or args.limit > 1000:
raise ValueError("--limit doit être compris entre 1 et 1000.")
if args.limit < 1 or args.limit > MAX_LIMIT:
raise ValueError(f"--limit doit être compris entre 1 et {MAX_LIMIT}.")
if args.start_time >= args.end_time:
raise ValueError("--start-time doit être antérieur à --end-time.")
@@ -1,11 +1,24 @@
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
TAILLE_DE_LOT = 1000
class RecommendationRepository:
def __init__(self, session: AsyncSession) -> None:
self._session = session
@@ -20,3 +33,19 @@ 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:
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,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
+37 -1
View File
@@ -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})"
+108 -1
View File
@@ -1453,6 +1453,90 @@
}
}
},
"/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": {
"get": {
"tags": [
@@ -2344,6 +2428,29 @@
],
"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": {
"properties": {
"recommendation_id": {
@@ -3153,7 +3260,7 @@
},
{
"name": "recommendations",
"description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`."
"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`."
},
{
"name": "stats",
+1
View File
@@ -51,6 +51,7 @@ ROLE_MINIMUM: Final[dict[Route, Role]] = {
("GET", "/api/v1/alerts"): Role.LECTEUR,
("GET", "/api/v1/recommendations"): 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/readings"): Role.LECTEUR,
("GET", "/api/v1/predictions"): Role.LECTEUR,
+59 -1
View File
@@ -10,7 +10,7 @@ from app.api.deps import get_current_principal, get_recommendation_service
from app.core.principal import Principal
from app.core.roles import AccountKind, Role
from app.models.energy import Recommendation
from app.services.recommendation import RecommendationNotFoundError
from app.services.recommendation import RapportGeneration, RecommendationNotFoundError
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
@@ -40,6 +40,7 @@ class FauxService:
def __init__(self, erreur: Exception | None = None) -> None:
self._erreur = erreur
self.recommendation = recommendation()
self.site_demande: str | None = None
async def list_all(self) -> list[Recommendation]:
return [self.recommendation]
@@ -49,6 +50,10 @@ class FauxService:
raise self._erreur
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
def lecteur_connecte(app: FastAPI) -> Iterator[None]:
@@ -142,3 +147,56 @@ async def test_get_recommendation_returns_404_when_the_session_finds_nothing(
response = await client.get("/api/v1/recommendations/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
@@ -13,10 +13,12 @@ from sqlalchemy.ext.asyncio import AsyncSession
import app.etl.mock_api_import as mock_api_import
from app.etl.mock_api_import import (
MAX_SITES,
READING_INSERT,
SOURCE_HISTORY,
build_reading_batch,
build_reading_row,
build_site_row,
fetch_readings,
fetch_sites,
upsert_sites,
@@ -73,6 +75,26 @@ async def test_fetch_sites_returns_sites() -> None:
assert sites[0]["site_type"] == "office"
async def test_fetch_sites_rejects_non_list_response() -> None:
def handler(request: Request) -> Response:
return Response(
status_code=200,
json={"unexpected": "payload"},
)
transport = MockTransport(handler)
async with AsyncClient(
transport=transport,
base_url="https://mock.test",
) as client:
with pytest.raises(
ValueError,
match="La réponse /api/v1/sites doit être une liste",
):
await fetch_sites(client)
async def test_fetch_readings_sends_expected_query_parameters() -> None:
captured_params: dict[str, str] = {}
@@ -576,6 +598,148 @@ def test_main_runs_import(
)
async def test_fetch_sites_rejects_a_response_above_the_cap() -> None:
def handler(request: Request) -> Response:
return Response(
status_code=200,
json=[make_site() for _ in range(MAX_SITES + 1)],
)
transport = MockTransport(handler)
async with AsyncClient(
transport=transport,
base_url="https://mock.test",
) as client:
with pytest.raises(
ValueError,
match=f"dépasse le plafond de {MAX_SITES} sites",
):
await fetch_sites(client)
async def test_fetch_readings_rejects_a_response_above_the_requested_limit() -> None:
def handler(request: Request) -> Response:
return Response(
status_code=200,
json=[make_reading(), make_reading(), make_reading()],
)
transport = MockTransport(handler)
async with AsyncClient(
transport=transport,
base_url="https://mock.test",
) as client:
with pytest.raises(
ValueError,
match="dépasse la limite demandée de 2",
):
await fetch_readings(
client=client,
site_id="SITE001",
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
limit=2,
)
def test_build_reading_row_neutralises_values_outside_physical_bounds() -> None:
reading = make_reading()
reading["power_factor"] = 42.0
reading["temperature_celsius"] = 1e30
reading["humidity_percent"] = -1.0
row = build_reading_row(reading)
assert row["power_factor"] is None
assert row["temperature_celsius"] is None
assert row["humidity_percent"] is None
assert row["null_reasons"] == [
"out_of_physical_bounds:power_factor",
"out_of_physical_bounds:temperature_celsius",
"out_of_physical_bounds:humidity_percent",
]
assert row["data_quality"] == "degraded"
assert json.loads(row["raw_data"])["power_factor"] == 42.0
def test_build_reading_row_rejects_a_measure_that_is_not_a_number() -> None:
reading = make_reading()
reading["consumption_kw"] = "87.34"
row = build_reading_row(reading)
assert row["consumption_kw"] is None
assert "out_of_physical_bounds:consumption_kw" in row["null_reasons"]
def test_build_reading_row_drops_a_quality_the_database_refuses() -> None:
reading = make_reading()
reading["data_quality"] = "unknown"
row = build_reading_row(reading)
assert row["data_quality"] is None
def test_build_reading_row_requires_an_identifier() -> None:
reading = make_reading()
del reading["site_id"]
with pytest.raises(
ValueError,
match="Champ site_id absent ou invalide",
):
build_reading_row(reading)
def test_build_site_row_keeps_only_the_expected_columns() -> None:
site = make_site()
site["unexpected"] = "valeur hostile"
site["capacity_kw"] = -5.0
site["status"] = 12
row = build_site_row(site)
assert set(row) == {
"site_id",
"site_type",
"site_name",
"location",
"capacity_kw",
"status",
}
assert row["capacity_kw"] is None
assert row["status"] is None
async def test_upsert_sites_sends_only_the_expected_columns() -> None:
connection = AsyncMock()
site = make_site()
site["unexpected"] = "valeur hostile"
await upsert_sites(
connection,
[site],
)
rows = connection.execute.await_args.args[1]
assert "unexpected" not in rows[0]
assert rows[0]["site_id"] == "SITE001"
@pytest.mark.integration
async def test_reading_insert_is_idempotent(
session: AsyncSession,
@@ -620,3 +784,51 @@ async def test_reading_insert_is_idempotent(
assert result.scalar_one() == 1
await session.rollback()
@pytest.mark.integration
async def test_out_of_bounds_reading_is_stored_neutralised(
session: AsyncSession,
) -> None:
reading = make_reading()
reading["power_factor"] = 42.0
row = build_reading_row(reading)
connection = await session.connection()
await upsert_sites(
connection,
[make_site()],
)
await session.execute(
READING_INSERT,
[row],
)
result = await session.execute(
text(
"""
SELECT power_factor, data_quality, null_reasons, raw_data ->> 'power_factor'
FROM reading
WHERE site_id = :site_id
AND timestamp = :timestamp
AND source = :source
"""
),
{
"site_id": row["site_id"],
"timestamp": row["timestamp"],
"source": row["source"],
},
)
stored = result.one()
await session.rollback()
assert stored[0] is None
assert stored[1] == "degraded"
assert stored[2] == ["out_of_physical_bounds:power_factor"]
assert stored[3] == "42.0"
@@ -5,7 +5,8 @@ import pytest
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.energy import Alert, Recommendation, Site
from app.repositories.recommendation import RecommendationRepository
from app.repositories import recommendation as module_recommendation
from app.repositories.recommendation import NouvelleRecommandation, RecommendationRepository
pytestmark = pytest.mark.integration
@@ -83,3 +84,59 @@ async def test_list_all_returns_the_recommendations_sorted_by_identifier(
await session.rollback()
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,14 @@
from collections.abc import Sequence
from datetime import UTC, datetime
import pytest
from app.models.energy import Recommendation
from app.models.energy import Alert, Recommendation
from app.repositories.recommendation import NouvelleRecommandation
from app.services.recommendation import RecommendationNotFoundError, RecommendationService
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
def recommendation(recommendation_id: int = 1) -> Recommendation:
return Recommendation(
@@ -13,13 +17,33 @@ def recommendation(recommendation_id: int = 1) -> Recommendation:
action="Vérifier la consommation",
explanation="Pic détecté",
rule_reference="spike-v1",
created_at=datetime(2024, 1, 1, tzinfo=UTC),
created_at=MOMENT,
)
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:
def __init__(self, recommendations: list[Recommendation]) -> None:
def __init__(self, recommendations: list[Recommendation], creees: int | None = None) -> None:
self._recommendations = recommendations
self._creees = creees
self.recues: list[NouvelleRecommandation] = []
async def list_all(self) -> list[Recommendation]:
return self._recommendations
@@ -29,27 +53,111 @@ class FakeRepository:
(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
async def test_list_all_returns_the_repository_recommendations() -> None:
service = RecommendationService(
recommendations=FakeRepository([recommendation(1), recommendation(2)])
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(),
)
recommendations = await service.list_all()
async def test_list_all_returns_the_repository_recommendations() -> None:
depot = FakeRepository([recommendation(1), recommendation(2)])
recommendations = await service(recommendations=depot).list_all()
assert [r.recommendation_id for r in recommendations] == [1, 2]
async def test_get_by_id_returns_the_matching_recommendation() -> None:
service = RecommendationService(recommendations=FakeRepository([recommendation(1)]))
trouve = await service.get_by_id(1)
trouve = await service(recommendations=FakeRepository([recommendation(1)])).get_by_id(1)
assert trouve.recommendation_id == 1
async def test_get_by_id_raises_when_the_recommendation_is_unknown() -> None:
service = RecommendationService(recommendations=FakeRepository([]))
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
@@ -0,0 +1,142 @@
from datetime import UTC, datetime
import pytest
from app.models.energy import Alert
from app.services.recommendation_rules import FACTEUR_DEPASSEMENT_MAJEUR, applique_les_regles
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
def alerte(
*,
alert_id: int = 1,
type_alerte: str = "spike",
severity: str = "high",
value: float | None = None,
threshold: float | None = None,
metric: str | None = None,
site_id: str = "SITE001",
) -> Alert:
return Alert(
alert_id=alert_id,
source_alert_id=f"ALR-{alert_id}",
site_id=site_id,
source="api_mock",
timestamp=MOMENT,
type=type_alerte,
severity=severity,
message="Alerte de test",
value=value,
threshold=threshold,
metric=metric,
prediction_id=None,
raw_data={},
)
@pytest.mark.parametrize(
("type_alerte", "attendue"),
[
("spike", "spike-delestage-v1"),
("threshold", "threshold-reduction-v1"),
("outage", "outage-secours-v1"),
("sensor", "sensor-maintenance-v1"),
("anomaly", "anomaly-verification-v1"),
],
ids=["pic", "seuil", "coupure", "capteur", "anomalie"],
)
def test_each_alert_type_yields_its_own_rule(type_alerte: str, attendue: str) -> None:
proposees = applique_les_regles(alerte(type_alerte=type_alerte))
assert [p.rule_reference for p in proposees] == [attendue]
def test_a_critical_alert_adds_the_escalation_rule() -> None:
proposees = applique_les_regles(alerte(severity="critical"))
assert "escalade-astreinte-v1" in {p.rule_reference for p in proposees}
@pytest.mark.parametrize("severity", ["low", "medium", "high"], ids=["faible", "moyenne", "haute"])
def test_a_non_critical_alert_does_not_escalate(severity: str) -> None:
proposees = applique_les_regles(alerte(severity=severity))
assert "escalade-astreinte-v1" not in {p.rule_reference for p in proposees}
def test_a_large_overshoot_adds_the_contract_rule() -> None:
proposees = applique_les_regles(
alerte(value=720.0 * FACTEUR_DEPASSEMENT_MAJEUR, threshold=720.0)
)
assert "contrat-puissance-v1" in {p.rule_reference for p in proposees}
def test_an_overshoot_below_the_factor_does_not_add_the_contract_rule() -> None:
proposees = applique_les_regles(alerte(value=800.0, threshold=720.0))
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
@pytest.mark.parametrize(
("value", "threshold"),
[(None, 720.0), (900.0, None), (900.0, 0.0), (900.0, -10.0)],
ids=["sans mesure", "sans seuil", "seuil nul", "seuil negatif"],
)
def test_the_contract_rule_stays_silent_without_an_exploitable_threshold(
value: float | None, threshold: float | None
) -> None:
proposees = applique_les_regles(alerte(value=value, threshold=threshold))
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
def test_the_explanation_quotes_the_measure_and_the_threshold() -> None:
proposees = applique_les_regles(alerte(value=812.5, threshold=720.0, metric="consumption_kw"))
assert "(consumption_kw mesurée à 812.5, seuil 720.0)" in proposees[0].explanation
def test_the_explanation_quotes_the_measure_alone_when_no_threshold_is_known() -> None:
proposees = applique_les_regles(alerte(value=812.5, metric="consumption_kw"))
assert "(consumption_kw mesurée à 812.5)" in proposees[0].explanation
def test_the_explanation_omits_the_measure_when_the_alert_carries_none() -> None:
proposees = applique_les_regles(alerte())
assert "(" not in proposees[0].explanation
def test_the_explanation_names_the_site() -> None:
proposees = applique_les_regles(alerte(site_id="SITE042"))
assert "SITE042" in proposees[0].explanation
def test_every_proposal_carries_the_alert_identifier() -> None:
proposees = applique_les_regles(alerte(alert_id=77, severity="critical"))
assert {p.alert_id for p in proposees} == {77}
def test_an_alert_never_yields_the_same_rule_twice() -> None:
proposees = applique_les_regles(
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
)
assert len(proposees) == len({p.rule_reference for p in proposees})
def test_a_critical_alert_over_the_threshold_yields_the_three_rules() -> None:
proposees = applique_les_regles(
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
)
assert {p.rule_reference for p in proposees} == {
"spike-delestage-v1",
"escalade-astreinte-v1",
"contrat-puissance-v1",
}
+27
View File
@@ -118,3 +118,30 @@ def test_main_exports_the_contract_without_asking_for_a_password(
assert code == 0
assert destination.exists()
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
@@ -0,0 +1,81 @@
# 0006 - Le moteur de règles de recommandation vit dans le backend
- Statut : accepté
- Date : 2026-09-18
## Contexte
L'issue #38 demande un « moteur de règles pour recommandations », portée par le label `ml`. Le
schéma tranche déjà la forme du résultat : `recommendation(alert_id, action, explanation,
rule_reference)`, avec `alert_id` en clé étrangère `NOT NULL` et une contrainte d'unicité
`uq_recommendation_alert_rule` sur `(alert_id, rule_reference)`. Une recommandation est donc
**dérivée d'une alerte**, jamais d'une mesure brute ni d'une prévision.
Deux emplacements se disputaient le code :
1. `ml/enervision_ml/`, sur le patron de `enervision_ml.score` livré par #37 : un script autonome
qui se connecte par `ML_DATABASE_URL`, écrit une table, et que l'API se contente de lire.
L'[ADR 0005](0005-modele-prediction-lightgbm.md) annonce d'ailleurs #38 de ce côté, en écrivant
que le scoring, le moteur de recommandations et les tests de dérive « consommeront le même
module `enervision_ml.features` ».
2. `apps/backend/app/services/`, où `apps/backend/README.md` place les « regles metier ».
## Décision
**Le moteur vit dans `apps/backend/app/services/`**, sous la forme d'un module pur
`recommendation_rules.py` (le catalogue `REGLES`) et d'une méthode `RecommendationService.generate()`
qui l'applique, persiste et valide la transaction.
Trois raisons :
- **Il n'utilise rien du ML.** Le catalogue lit `alert.type`, `alert.severity`, `alert.value` et
`alert.threshold`. Aucun modèle, aucune feature, aucun `enervision_ml.features` : la phrase de
l'ADR 0005 vaut pour le scoring (#37) et les tests de dérive (#44/#45), qui manipulent bien des
features, pas pour des règles sur alertes. Le label `ml` de #38 désigne le lot fonctionnel
« prédiction et recommandation », pas l'emplacement du code.
- **Il lit et écrit deux tables déjà couvertes par des repositories.** `AlertRepository` sait déjà
filtrer par site. Le placer dans `ml/` obligerait à réécrire ces accès en SQL brut, et à
maintenir deux représentations du même domaine.
- **Le déclencheur HTTP n'a de sens que dans l'API.** `POST /recommendations/generate` doit passer
par `require_role(Role.ADMIN)` et par la session injectée : cela suppose d'être dans
l'application FastAPI.
Le moteur reste néanmoins **déclenchable hors HTTP**, par `python -m app.cli
generate-recommendations` (cible `make recommendations`), sur le patron de `make ml-score` : rien
n'oblige à exposer un port pour régénérer des recommandations.
## Conséquences
- L'API gagne sa première route d'écriture métier. La checklist de `20-backend.md` s'applique :
entrée dans `ROLE_MINIMUM` de `tests/api/acces.py`, et `openapi.json` régénéré dans le même
commit.
- `RecommendationService` n'est plus en lecture seule : il reçoit le `Transaction` Protocol déjà
utilisé par `AuthService` et `UserService`, et commite lui-même. Les repositories continuent de
ne pas commiter.
- **L'idempotence est déléguée à la base.** `create_missing()` insère en `ON CONFLICT DO NOTHING`
sur `uq_recommendation_alert_rule` plutôt que de relire avant d'écrire, ce qui supprime la
fenêtre entre le contrôle et l'insertion. Corollaire : `rule_reference` est une clé fonctionnelle.
Une règle dont le sens change prend une référence `-v2` ; renommer une référence livrée
ferait réapparaître ses recommandations à côté des anciennes.
- **Le moteur est branché sur la détection interne, et sur elle seule.** `alert` est alimentée
par `app/detection/internal_alerts.py` (#104), lancée à la main comme `enervision_ml.score` ;
l'ingestion de l'API Mock `/alerts` reste à faire. Le rapport de génération est donc à zéro tant
que la détection n'a pas tourné, sans que le moteur soit à retoucher.
- **L'insertion est découpée en lots.** `create_missing()` écrit par paquets de `TAILLE_DE_LOT`
lignes : asyncpg plafonne une requête à 32 767 paramètres, soit 8 191 lignes de quatre colonnes,
et la détection interne peut alimenter `alert` au fil de l'eau.
- Si le projet devait un jour pondérer les recommandations par un score appris, la décision serait
à rouvrir : le moteur redeviendrait consommateur du pipeline ML.
## Alternatives écartées
- **Module et CLI dans `ml/enervision_ml/`** : cohérent avec le label `ml` et avec la lettre de
l'ADR 0005, mais impose du SQL brut là où deux repositories existent, et laisse la génération
hors de portée de l'API. Redeviendrait le bon choix si les règles se mettaient à consommer des
features ou un modèle.
- **Génération à la volée, sans persistance**, calculée à chaque `GET /recommendations` : supprime
le besoin d'écriture, mais rend la table `recommendation` et sa contrainte d'unicité inutiles,
et interdit toute trace de ce qui a été proposé et quand.
- **Table de configuration des règles en base**, plutôt qu'un catalogue en Python : plus souple,
mais déplace la logique métier hors de la revue de code et hors des tests, pour un besoin que
rien n'exprime à ce stade.
+5
View File
@@ -165,3 +165,8 @@ Elles vivent dans `../adr/`, pas ici.
| ADR | Objet |
|---|---|
| [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/` |
+16
View File
@@ -146,6 +146,7 @@ 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/recommendations` | Liste les recommandations. `lecteur` | 401, 403, 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/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 |
@@ -194,6 +195,21 @@ 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.
[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,
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
+67 -26
View File
@@ -12,8 +12,10 @@ d'énergie, dont l'hypertable `reading`.
Les sections marquées `Fait` relèvent du code déjà implémenté. Les sections marquées `Cible`
décrivent les éléments prévus mais pas encore réalisés.
L'ingestion des deux sources de données du MVP est maintenant implémentée. L'orchestration
Airflow, les agrégats continus, la compression et la rétention restent des cibles.
L'ingestion des **mesures** est implémentée pour les deux sources du MVP, le dataset CSV/JSON et
l'API Mock. Celle des **alertes** de l'API Mock, `/alerts`, reste à faire : voir
l'[ADR 0006](../adr/0006-moteur-de-regles-dans-le-backend.md). L'orchestration Airflow, les
agrégats continus, la compression et la rétention restent des cibles.
## Trois emplacements, trois rôles
@@ -96,8 +98,8 @@ Les flèches pleines représentent les traitements actuellement implémentés.
Les flèches pointillées représentent les éléments encore prévus comme cibles.
Les lectures futures de l'API et de Grafana visent l'agrégat continu plutôt que la table brute
lorsque cette partie TimescaleDB sera mise en place.
Les lectures de l'API et de Grafana viseront l'agrégat continu, pas la table brute : c'est tout
l'intérêt de TimescaleDB, et cela doit rester vrai quand les volumes augmenteront.
## Tables d'authentification
@@ -170,21 +172,28 @@ erDiagram
}
```
Plusieurs choix de modélisation portent une intention précise :
Six choix de modélisation portent une intention et se défendent seuls :
- **`app_user` et non `user`** : `user` est un mot réservé PostgreSQL, raccourci de
`CURRENT_USER`. Le nom rappelle aussi qu'il s'agit d'un compte applicatif.
- **`credentials_changed_at`, une seule colonne**, couvre notamment le changement de mot de passe,
le changement de rôle et la désactivation.
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation.
- **`audit_log.actor_id` n'a aucune clé étrangère** afin de conserver les informations d'audit
même si l'entité d'origine évolue.
- `password_reset_token` ne stocke que l'empreinte du jeton et jamais sa valeur directement.
- `password_reset_attempt` est séparée de `audit_log`, car son volume peut être piloté
par des demandes externes répétées.
`CURRENT_USER`. Le nom rappelle en prime qu'il s'agit d'un compte applicatif, par opposition
au rôle PostgreSQL qui portera le cantonnement de l'ETL.
- **`credentials_changed_at`, une seule colonne**, couvre le changement de mot de passe, le
changement de rôle et la désactivation. Un compteur de version ne dirait rien à un humain qui
lit un audit.
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation. S'il
glissait, la promesse de sept jours serait fictive et une session active ne finirait jamais.
- **`audit_log.actor_id` n'a aucune clé étrangère**, et `actor_email` comme `actor_role` sont
dénormalisés. Une contrainte `ON DELETE SET NULL` déclencherait un `UPDATE` que le déclencheur
d'ajout seul refuserait. Voir l'[ADR 0004](../adr/0004-journal-d-audit-en-ajout-seul.md).
- **`password_reset_token` ne stocke que l'empreinte du jeton**, jamais sa valeur. Une fuite de
la table ne donne donc rien à rejouer.
- **`password_reset_attempt` est séparée de `audit_log`** : son volume est piloté par le
demandeur, comme celui de `login_attempt`, donc elle doit pouvoir se purger.
`audit_log` porte des déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`.
Elle n'est donc **pas** une hypertable.
`audit_log` porte deux déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`. Elle n'est
donc **pas** une hypertable : une politique de rétention émettrait des `DELETE` qu'ils
refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est
piloté par l'attaquant.
## Gabarit de révision créant une hypertable
@@ -234,7 +243,8 @@ colonne de temps : les index déclarés dans la révision le couvrent déjà.
## Questions ouvertes
Elles portent maintenant principalement sur l'exploitation du schéma :
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma et l'ingestion sont
livrés : ce qui suit porte sur leur exploitation, plus sur leur forme.
- **Quelle granularité** conserver à long terme à l'ingestion : seconde, minute ou quart d'heure.
- **Quels agrégats continus** créer et sur quelles fenêtres.
@@ -279,6 +289,12 @@ Les anomalies historiques décrites dans les JSON sont conservées dans `dataset
Elles servent à l'analyse des données 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
- Un site possède plusieurs mesures, prévisions et alertes.
@@ -287,7 +303,7 @@ Elles servent à l'analyse des données et ne sont pas considérées comme des a
- Une alerte peut être associée à une prévision du même site.
- Une alerte peut donner lieu à plusieurs recommandations.
# Ingestion des données historiques
## Ingestion des données historiques
Statut : `Fait`.
@@ -301,7 +317,7 @@ Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisat
Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
## Architecture du flux historique
### Architecture du flux historique
```text
Dataset CSV + métadonnées JSON
@@ -351,7 +367,7 @@ source = "csv"
dataset_id = identifiant du dataset
```
## Résultats validés pour l'historique
### Résultats validés pour l'historique
Le chargement de référence a permis d'obtenir :
@@ -366,7 +382,7 @@ aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à
La procédure détaillée d'installation, d'exécution, de validation et de contrôle du pipeline
est disponible dans `etl/README.md`.
# Ingestion depuis l'API Mock
## Ingestion depuis l'API Mock
Statut : `Fait`.
@@ -378,7 +394,7 @@ Le traitement est implémenté dans :
apps/backend/app/etl/mock_api_import.py
```
## Endpoints utilisés
### Endpoints utilisés
Le pipeline récupère les informations des sites depuis :
@@ -410,7 +426,7 @@ Les paramètres de ligne de commande disponibles pour l'import sont :
--dry-run
```
## Flux d'ingestion API Mock
### Flux d'ingestion API Mock
```text
API Mock
@@ -454,7 +470,32 @@ La réponse source reçue depuis l'API est conservée dans :
raw_data
```
## Qualité des données de l'API Mock
### Frontière de confiance avec l'API Mock
L'API Mock de l'école n'a aucune authentification et expose un endpoint mutatif à quiconque. Sa
réponse est donc traitée comme une entrée hostile, conformément à API10 dans
[la traçabilité OWASP](owasp-traceabilite.md). Le risque premier n'est pas la fausse alerte,
c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction.
Quatre garde-fous, tous dans `mock_api_import.py` :
| Garde-fou | Mise en œuvre |
|---|---|
| Timeout | `APP_MOCK_API_TIMEOUT_SECONDS`, dix secondes par défaut |
| Taille de tableau plafonnée | `MAX_SITES` sites, et au plus `--limit` mesures par site |
| Bornes physiques | `PHYSICAL_BOUNDS`, une plage par grandeur |
| Frontière d'anti-corruption | `build_site_row()` et `build_reading_row()`, qui ne recopient que les champs attendus |
Une valeur hors bornes, d'un type inattendu, `NaN` ou infinie devient `NULL`. Elle laisse sa
trace dans `null_reasons` sous la forme `out_of_physical_bounds:<colonne>`, et `data_quality`
descend à `degraded`. Une `data_quality` que `ck_reading_quality` refuserait devient `NULL`
plutôt que de faire échouer le lot entier. Dans tous les cas `raw_data` conserve la réponse
d'origine intacte : rien n'est perdu, seule son exploitation est bornée.
Le plafond de taille s'applique après désérialisation de la réponse. Borner le corps HTTP
lui-même demanderait une lecture en flux, et reste à faire.
### Qualité des données de l'API Mock
Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion.
@@ -475,7 +516,7 @@ imputed_values = NULL
imputation_method = NULL
```
## Validation de l'import API Mock
### Validation de l'import API Mock
Un scénario de validation a été exécuté pour les 7 sites sur la période :
@@ -526,7 +567,7 @@ Les tests automatisés couvrent également :
- la conservation des données sources ;
- l'idempotence en base.
# Évolution prévue
## Évolution prévue
La prochaine étape consiste à orchestrer les deux mécanismes d'ingestion avec Apache Airflow.
+2 -1
View File
@@ -41,6 +41,7 @@ lecture seule ; plusieurs lignes resteront à compléter une fois les endpoints
| En-têtes `nosniff`, `DENY`, `no-referrer`, et `no-store` sur les routes d'authentification | `app/api/middleware.py` | A05 |
| Refus de rétrograder ou désactiver le dernier administrateur actif | `app/services/user.py` | A04 Insecure Design |
| Amorçage du premier administrateur hors dépôt, mot de passe jamais dans `argv` ni dans Git | `app/cli.py` | A02, A05 |
| Réponse de l'API Mock bornée avant écriture : timeout, plafond de sites et de mesures, bornes physiques par grandeur, recopie des seuls champs attendus | `app/etl/mock_api_import.py` | API10 Unsafe Consumption of APIs |
| CI bloquante : format, lint avec règles Bandit, typage strict, tests avec seuil de couverture | `.github/workflows/backend.yml` | A06 Vulnerable and Outdated Components |
Note sur A06 : le jeu de règles `S` de ruff, déjà actif dans `pyproject.toml`, est le portage des
@@ -53,7 +54,7 @@ règles Bandit. Ajouter Bandit à la CI serait redondant, contrairement à ce qu
| **API1 Broken Object Level Authorization** | **ouvert** | Les rôles sont globaux, il n'y a pas de portée par site : `GET /sites/{site_id}` et `GET /recommendations/{recommendation_id}` répondent à tout compte `lecteur` pour n'importe quel site ou recommandation, sans vérifier une affectation compte-site qui n'existe pas encore. Un opérateur du site A pourra agir sur le site B dès que les endpoints d'écriture métier existeront. Correctif prévu : table d'affectation compte-site, contrôle d'appartenance dans la même dépendance que le contrôle de rôle. |
| **API4, lectures de séries temporelles** | **partiel** | `GET /readings` plafonne la fenêtre temporelle (90 jours) et la pagination (`limit` ≤ 2000), voir plus haut. Reste ouvert : pagination en `limit`/`offset` simple plutôt qu'en curseur (un `offset` élevé sur une fenêtre dense reste coûteux), et aucun `statement_timeout` au niveau de la connexion pour borner une requête individuelle si les plafonds au-dessus s'avéraient insuffisants. |
| **API8 Security Misconfiguration, transport** | **ouvert** | Pas de TLS, donc ni HSTS, ni cookie `Secure` réellement posé en production. Ils appartiennent au terminateur TLS, qui n'existe pas. |
| **API10 Unsafe Consumption of APIs** | **ouvert, et spécifique à ce projet** | L'API Mock de l'école n'a aucune authentification, tourne en HTTP clair sur le réseau de l'école, et expose un endpoint mutatif à quiconque. Sa réponse doit être traitée comme une entrée hostile : bornes physiques, taille de tableau plafonnée, timeout, et frontière d'anti-corruption. La conséquence la plus sérieuse n'est pas la fausse alerte, c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction. |
| **API10 Unsafe Consumption of APIs** | **partiel, et spécifique à ce projet** | L'API Mock de l'école n'a aucune authentification, tourne en HTTP clair sur le réseau de l'école, et expose un endpoint mutatif à quiconque. Sa réponse est traitée comme une entrée hostile par `app/etl/mock_api_import.py`, son seul consommateur à ce jour : les quatre garde-fous attendus sont en place, voir la ligne correspondante plus haut. Reste ouvert : le plafond de taille s'applique après désérialisation de la réponse, borner le corps HTTP lui-même demanderait une lecture en flux ; et `APP_MOCK_API_BASE_URL` n'impose pas `https`, donc les identifiants Basic partiraient en clair sur une URL en `http`. La conséquence la plus sérieuse n'est pas la fausse alerte, c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction. |
| **A08 Software and Data Integrity Failures** | **partiel** | La CI vérifie le code mais n'analyse ni les dépendances ni les images. `.terraform.lock.hcl` reste ignoré par git, ce qui contredit une chaîne d'approvisionnement maîtrisée. |
| **A10 Server-Side Request Forgery** | **sans objet aujourd'hui** | Aucune URL sortante n'est pilotée par une donnée utilisateur. Le jour où l'adresse d'une source devient un champ de configuration, il faudra une liste blanche de schémas et d'hôtes, sans suivi de redirection. |
| **Cantonnement des accès ETL et ML** | **dette assumée** | Le compte applicatif porte l'identité, le rôle PostgreSQL porterait le cantonnement. Voir ADR 0003. |
+63 -25
View File
@@ -81,9 +81,9 @@ L'API Mock est utilisée pour compléter les données historiques avec des mesur
| mypy | Vérification du typage |
| Pytest | Tests automatisés |
# Import du dataset historique
## Import du dataset historique
## Fonctionnement du pipeline historique
### Fonctionnement du pipeline historique
Le script d'import se trouve dans :
@@ -115,14 +115,14 @@ CSV + métadonnées JSON
PostgreSQL / TimescaleDB
```
### 1. Extraction
#### 1. Extraction
Le pipeline charge :
- `all_sites_combined.csv` avec Pandas ;
- `dataset_metadata.json` avec le module JSON de Python.
### 2. Validation
#### 2. Validation
Avant toute écriture en base, le pipeline contrôle notamment :
@@ -136,7 +136,7 @@ Avant toute écriture en base, le pipeline contrôle notamment :
Une incohérence détectée pendant cette étape interrompt l'import avant le chargement.
### 3. Dry-run
#### 3. Dry-run
Un mode `--dry-run` permet d'exécuter les contrôles sans écrire de données dans PostgreSQL.
@@ -149,7 +149,7 @@ Il permet notamment de vérifier :
- les valeurs NULL ;
- l'empreinte SHA-256.
### 4. Traçabilité
#### 4. Traçabilité
Une empreinte SHA-256 est calculée à partir du fichier CSV afin d'identifier le dataset utilisé.
@@ -161,7 +161,7 @@ Empreinte SHA-256 du dataset validé :
Cette empreinte participe à la traçabilité du dataset chargé.
### 5. Transformation
#### 5. Transformation
Les timestamps sont normalisés avec la timezone :
@@ -180,7 +180,7 @@ imputed_values = NULL
imputation_method = NULL
```
### 6. Chargement
#### 6. Chargement
Le chargement est réalisé avec SQLAlchemy Async dans PostgreSQL/TimescaleDB.
@@ -207,7 +207,7 @@ dataset_id = identifiant du dataset
Cette représentation respecte les contraintes définies dans le schéma de la base.
## Dataset validé
### Dataset validé
Le dataset traité contient :
@@ -226,7 +226,7 @@ Valeurs manquantes identifiées :
| `humidity_percent` | 3 423 |
| `solar_irradiance_wm2` | 3 964 |
## Exécution historique en dry-run
### Exécution historique en dry-run
Depuis le dossier :
@@ -246,7 +246,7 @@ uv run python -m app.etl.historical_import `
Aucune donnée n'est écrite dans la base pendant cette exécution.
## Chargement historique réel
### Chargement historique réel
Depuis `apps/backend/` :
@@ -268,7 +268,7 @@ Chargement : 2000/122647
Chargement : 122647/122647
```
## Résultats obtenus pour le dataset historique
### Résultats obtenus pour le dataset historique
Après le chargement initial, les contrôles en base ont confirmé :
@@ -285,7 +285,7 @@ Le premier import a créé :
nouvelles lectures : 122647
```
## Idempotence du dataset historique
### Idempotence du dataset historique
Le pipeline a été exécuté une deuxième fois avec exactement le même dataset afin de vérifier son idempotence.
@@ -299,7 +299,7 @@ nouvelles lectures : 0
Une nouvelle exécution du même import ne crée donc pas de mesures supplémentaires pour le dataset testé.
## Vérifications SQL du dataset historique
### Vérifications SQL du dataset historique
Depuis la racine du projet, vérifier le nombre d'enregistrements avec :
@@ -327,9 +327,9 @@ Résultat attendu pour le dataset historique :
csv | 122647
```
# Import depuis l'API Mock
## Import depuis l'API Mock
## Fonctionnement
### Fonctionnement
Le script d'import de l'API Mock se trouve dans :
@@ -386,7 +386,7 @@ limit
Le paramètre `limit` doit être compris entre 1 et 1000.
## Configuration de l'API Mock
### Configuration de l'API Mock
La connexion à l'API Mock est configurée avec les variables d'environnement suivantes :
@@ -401,7 +401,7 @@ Les identifiants réels ne sont pas versionnés dans Git.
Les fichiers `.env.example` indiquent uniquement les variables nécessaires à l'exécution.
## Transformation des mesures API
### Transformation des mesures API
Les mesures provenant de l'API Mock sont enregistrées dans `reading` avec :
@@ -422,7 +422,7 @@ raw_data
afin de préserver la donnée reçue et faciliter la traçabilité.
## Qualité des données API
### Qualité des données API
Les valeurs `NULL` fournies par l'API sont conservées telles quelles.
@@ -444,6 +444,9 @@ degraded
critical
```
Ce sont les quatre seules valeurs que la contrainte `ck_reading_quality` accepte. Toute autre
valeur renvoyée par l'API est remplacée par `NULL` plutôt que de faire échouer le lot entier.
Aucune imputation n'est réalisée pendant l'ingestion :
```text
@@ -453,7 +456,42 @@ imputation_method = NULL
Cette stratégie permet de distinguer une véritable valeur nulle ou manquante d'une consommation égale à zéro et de conserver les informations liées aux défaillances de capteurs.
## Dry-run de l'API Mock
### Bornes physiques et frontière de confiance
La réponse de l'API Mock est traitée comme une entrée hostile : l'API n'a pas
d'authentification et expose un endpoint mutatif à quiconque. Voir API10 dans
`docs/architecture/owasp-traceabilite.md`.
Les plages acceptées sont déclarées dans `PHYSICAL_BOUNDS` :
| Grandeur | Plage acceptée |
|---|---|
| `consumption_kw` | 0 à 100 000 |
| `consumption_kwh` | 0 à 100 000 |
| `voltage_v` | 0 à 1 000 |
| `current_a` | 0 à 10 000 |
| `power_factor` | 0 à 1 |
| `temperature_celsius` | -90 à 60 |
| `humidity_percent` | 0 à 100 |
| `capacity_kw` | 0 à 100 000 |
Une valeur hors plage, d'un type inattendu, `NaN` ou infinie devient `NULL` :
```text
null_reasons += "out_of_physical_bounds:<colonne>"
data_quality = "degraded"
```
L'import ne s'interrompt pas pour autant : le mock émet des anomalies par construction, et
`raw_data` conserve la réponse d'origine.
La taille des réponses est plafonnée : au plus `MAX_SITES` sites, et au plus `--limit` mesures
par site. Au-delà, l'import échoue au lieu de charger.
Enfin, seuls les champs attendus sont recopiés vers la base. Une clé supplémentaire renvoyée par
l'API n'atteint jamais une colonne.
### Dry-run de l'API Mock
Le mode `--dry-run` permet de tester la connexion, la récupération des sites et la récupération des mesures sans écrire dans PostgreSQL.
@@ -467,7 +505,7 @@ uv run python -m app.etl.mock_api_import `
--dry-run
```
## Chargement réel depuis l'API Mock
### Chargement réel depuis l'API Mock
Depuis `apps/backend/` :
@@ -478,7 +516,7 @@ uv run python -m app.etl.mock_api_import `
--limit 60
```
## Résultat validé pour l'API Mock
### Résultat validé pour l'API Mock
Le scénario de validation utilisé couvre la période :
@@ -511,7 +549,7 @@ Les contrôles effectués directement dans PostgreSQL/TimescaleDB ont confirmé
- la conservation de `null_reasons` ;
- la conservation de la donnée source dans `raw_data`.
## Idempotence de l'import API Mock
### Idempotence de l'import API Mock
Le même import a été exécuté plusieurs fois afin de vérifier qu'une mesure déjà présente n'est pas créée une seconde fois.
@@ -519,7 +557,7 @@ L'idempotence repose sur la contrainte d'unicité de la table `reading` et sur l
Un test d'intégration automatisé vérifie également ce comportement.
# Tests et qualité
## Tests et qualité
Les tests automatisés des pipelines ETL sont situés dans :
@@ -599,7 +637,7 @@ Lors de la validation de l'import API Mock :
La suite backend complète a également été validée avec une couverture supérieure au seuil de 85 %.
# Suite du pipeline Data
## Suite du pipeline Data
Deux sources de données sont maintenant prises en charge :