Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
56c6b79a5a | ||
|
|
452cfdef85 | ||
|
|
96dd1f834c | ||
|
|
e66ef86729 | ||
|
|
0318ee6cc5 | ||
|
|
0ddfb1997d |
@@ -17,3 +17,9 @@ APP_LOG_LEVEL=INFO
|
||||
APP_SECRET_KEY=change_me
|
||||
APP_CORS_ORIGINS=http://localhost:4200
|
||||
BACKEND_PORT=8000
|
||||
|
||||
# API Mock EnerVision
|
||||
APP_MOCK_API_BASE_URL=https://api-mock.charlieandre.fr
|
||||
APP_MOCK_API_USERNAME=change_me
|
||||
APP_MOCK_API_PASSWORD=change_me
|
||||
APP_MOCK_API_TIMEOUT_SECONDS=10
|
||||
|
||||
@@ -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 recommendations
|
||||
ml-lint ml-typecheck ml-test ml-check ml-train ml-score
|
||||
|
||||
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,9 +77,6 @@ ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DAT
|
||||
ml-score: ## Score le prochain pas horaire et l'ecrit dans `prediction`. CSV=chemin optionnel
|
||||
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)
|
||||
|
||||
|
||||
@@ -18,3 +18,7 @@ APP_SMTP_HOST=localhost
|
||||
APP_SMTP_PORT=1025
|
||||
APP_SMTP_USE_TLS=false
|
||||
APP_SMTP_FROM_ADDRESS=no-reply@enervision.fr
|
||||
APP_MOCK_API_BASE_URL=https://api-mock.charlieandre.fr
|
||||
APP_MOCK_API_USERNAME=change_me
|
||||
APP_MOCK_API_PASSWORD=change_me
|
||||
APP_MOCK_API_TIMEOUT_SECONDS=10
|
||||
|
||||
@@ -113,7 +113,6 @@ Le sens de dependance est unique : `endpoints` vers `services` vers `repositorie
|
||||
| `/api/v1/sites/{site_id}` | Décrit un site | `lecteur` |
|
||||
| `/api/v1/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 |
|
||||
|
||||
|
||||
@@ -192,11 +192,7 @@ AlertServiceDep = Annotated[AlertService, Depends(get_alert_service)]
|
||||
|
||||
|
||||
def get_recommendation_service(session: SessionDep) -> RecommendationService:
|
||||
return RecommendationService(
|
||||
recommendations=RecommendationRepository(session),
|
||||
alerts=AlertRepository(session),
|
||||
transaction=session,
|
||||
)
|
||||
return RecommendationService(recommendations=RecommendationRepository(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`. Leur génération par le moteur de règles est réservée au rôle `admin`."
|
||||
"`lecteur`."
|
||||
),
|
||||
},
|
||||
{
|
||||
|
||||
@@ -1,18 +1,13 @@
|
||||
from fastapi import APIRouter, HTTPException, status
|
||||
|
||||
from app.api.deps import AdminDep, LecteurDep, RecommendationServiceDep
|
||||
from app.api.openapi import REPONSE_VALIDATION, REPONSES_ADMIN, Reponses
|
||||
from app.api.deps import LecteurDep, RecommendationServiceDep
|
||||
from app.api.openapi import REPONSE_VALIDATION, Reponses
|
||||
from app.schemas.errors import ErrorResponse
|
||||
from app.schemas.recommendation import (
|
||||
RecommendationGenerationResponse,
|
||||
RecommendationResponse,
|
||||
)
|
||||
from app.schemas.recommendation import 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."},
|
||||
@@ -43,22 +38,3 @@ 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,11 +22,8 @@ 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"
|
||||
@@ -66,22 +63,6 @@ async def create_admin(
|
||||
)
|
||||
|
||||
|
||||
async def generate_recommendations(*, site_id: str | None) -> str:
|
||||
async with get_session_factory()() as session:
|
||||
service = RecommendationService(
|
||||
recommendations=RecommendationRepository(session),
|
||||
alerts=AlertRepository(session),
|
||||
transaction=session,
|
||||
)
|
||||
rapport = await service.generate(site_id=site_id)
|
||||
|
||||
return (
|
||||
f"{rapport.alertes_examinees} alerte(s) examinée(s), "
|
||||
f"{rapport.recommandations_creees} recommandation(s) créée(s), "
|
||||
f"{rapport.deja_presentes} déjà présente(s)"
|
||||
)
|
||||
|
||||
|
||||
# Piège : le schéma ne doit dépendre ni du `.env` du poste ni des variables `APP_*`, sinon le
|
||||
# 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.
|
||||
@@ -128,14 +109,6 @@ 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
|
||||
|
||||
|
||||
@@ -179,10 +152,6 @@ 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(
|
||||
|
||||
@@ -34,6 +34,11 @@ class Settings(BaseSettings):
|
||||
database_pool_size: int = 5
|
||||
database_max_overflow: int = 10
|
||||
|
||||
mock_api_base_url: str = "https://api-mock.charlieandre.fr"
|
||||
mock_api_username: str | None = None
|
||||
mock_api_password: SecretStr | None = None
|
||||
mock_api_timeout_seconds: float = Field(default=10.0, gt=0)
|
||||
|
||||
jwt_issuer: str = "enervision-api"
|
||||
jwt_audience: str = "enervision-web"
|
||||
access_token_ttl_seconds: int = Field(default=900, ge=60, le=3600)
|
||||
|
||||
@@ -0,0 +1,315 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import asyncio
|
||||
import json
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
|
||||
|
||||
from app.core.config import get_settings
|
||||
|
||||
SOURCE_HISTORY = "api_history"
|
||||
|
||||
|
||||
def create_mock_api_client() -> httpx.AsyncClient:
|
||||
settings = get_settings()
|
||||
|
||||
if settings.mock_api_username is None or settings.mock_api_password is None:
|
||||
raise ValueError("Les identifiants de l'API Mock ne sont pas configurés.")
|
||||
|
||||
return httpx.AsyncClient(
|
||||
base_url=settings.mock_api_base_url.rstrip("/"),
|
||||
auth=(
|
||||
settings.mock_api_username,
|
||||
settings.mock_api_password.get_secret_value(),
|
||||
),
|
||||
timeout=settings.mock_api_timeout_seconds,
|
||||
)
|
||||
|
||||
|
||||
async def fetch_sites(
|
||||
client: httpx.AsyncClient,
|
||||
) -> list[dict[str, Any]]:
|
||||
response = await client.get("/api/v1/sites")
|
||||
|
||||
response.raise_for_status()
|
||||
|
||||
payload = response.json()
|
||||
|
||||
if not isinstance(payload, list):
|
||||
raise ValueError("La réponse /api/v1/sites doit être une liste.")
|
||||
|
||||
return payload
|
||||
|
||||
|
||||
async def upsert_sites(
|
||||
connection: AsyncConnection,
|
||||
sites: list[dict[str, Any]],
|
||||
) -> None:
|
||||
if not sites:
|
||||
return
|
||||
|
||||
await connection.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO site (
|
||||
site_id,
|
||||
site_type,
|
||||
site_name,
|
||||
location,
|
||||
capacity_kw,
|
||||
status
|
||||
)
|
||||
VALUES (
|
||||
:site_id,
|
||||
:site_type,
|
||||
:site_name,
|
||||
:location,
|
||||
:capacity_kw,
|
||||
:status
|
||||
)
|
||||
ON CONFLICT (site_id)
|
||||
DO UPDATE SET
|
||||
site_type = EXCLUDED.site_type,
|
||||
site_name = EXCLUDED.site_name,
|
||||
location = EXCLUDED.location,
|
||||
capacity_kw = EXCLUDED.capacity_kw,
|
||||
status = EXCLUDED.status
|
||||
"""
|
||||
),
|
||||
sites,
|
||||
)
|
||||
|
||||
|
||||
async def fetch_readings(
|
||||
client: httpx.AsyncClient,
|
||||
site_id: str,
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
limit: int = 1000,
|
||||
) -> list[dict[str, Any]]:
|
||||
response = await client.get(
|
||||
"/api/v1/readings",
|
||||
params={
|
||||
"site_id": site_id,
|
||||
"start_time": start_time.isoformat(),
|
||||
"end_time": end_time.isoformat(),
|
||||
"limit": limit,
|
||||
},
|
||||
)
|
||||
|
||||
response.raise_for_status()
|
||||
|
||||
payload = response.json()
|
||||
|
||||
if not isinstance(payload, list):
|
||||
raise ValueError("La réponse /api/v1/readings doit être une liste.")
|
||||
|
||||
return payload
|
||||
|
||||
|
||||
def build_reading_row(
|
||||
reading: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
timestamp = datetime.fromisoformat(reading["timestamp"].replace("Z", "+00:00"))
|
||||
return {
|
||||
"site_id": reading["site_id"],
|
||||
"timestamp": timestamp,
|
||||
"source": SOURCE_HISTORY,
|
||||
"dataset_id": None,
|
||||
"consumption_kw": reading.get("consumption_kw"),
|
||||
"consumption_kwh": reading.get("consumption_kwh"),
|
||||
"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"),
|
||||
"imputed_values": None,
|
||||
"imputation_method": None,
|
||||
"raw_data": json.dumps(
|
||||
reading,
|
||||
ensure_ascii=False,
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
READING_INSERT = text(
|
||||
"""
|
||||
INSERT INTO reading (
|
||||
site_id,
|
||||
timestamp,
|
||||
source,
|
||||
dataset_id,
|
||||
consumption_kw,
|
||||
consumption_kwh,
|
||||
consumption_euros,
|
||||
voltage_v,
|
||||
current_a,
|
||||
power_factor,
|
||||
temperature_celsius,
|
||||
humidity_percent,
|
||||
solar_irradiance_wm2,
|
||||
is_working_hours,
|
||||
data_quality,
|
||||
null_reasons,
|
||||
imputed_values,
|
||||
imputation_method,
|
||||
raw_data
|
||||
)
|
||||
VALUES (
|
||||
:site_id,
|
||||
:timestamp,
|
||||
:source,
|
||||
:dataset_id,
|
||||
:consumption_kw,
|
||||
:consumption_kwh,
|
||||
:consumption_euros,
|
||||
:voltage_v,
|
||||
:current_a,
|
||||
:power_factor,
|
||||
:temperature_celsius,
|
||||
:humidity_percent,
|
||||
:solar_irradiance_wm2,
|
||||
:is_working_hours,
|
||||
:data_quality,
|
||||
:null_reasons,
|
||||
CAST(:imputed_values AS jsonb),
|
||||
:imputation_method,
|
||||
CAST(:raw_data AS jsonb)
|
||||
)
|
||||
ON CONFLICT DO NOTHING
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
def build_reading_batch(
|
||||
readings: list[dict[str, Any]],
|
||||
) -> list[dict[str, Any]]:
|
||||
return [build_reading_row(reading) for reading in readings]
|
||||
|
||||
|
||||
async def import_mock_api_history(
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
limit: int,
|
||||
dry_run: bool,
|
||||
) -> None:
|
||||
settings = get_settings()
|
||||
|
||||
async with create_mock_api_client() as client:
|
||||
sites = await fetch_sites(client)
|
||||
|
||||
print(f"Sites récupérés : {len(sites)}")
|
||||
|
||||
all_readings: list[dict[str, Any]] = []
|
||||
|
||||
for site in sites:
|
||||
site_id = site["site_id"]
|
||||
|
||||
readings = await fetch_readings(
|
||||
client=client,
|
||||
site_id=site_id,
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=limit,
|
||||
)
|
||||
|
||||
print(f"{site_id}: {len(readings)} lectures")
|
||||
|
||||
all_readings.extend(readings)
|
||||
|
||||
print(f"Lectures récupérées : {len(all_readings)}")
|
||||
|
||||
if dry_run:
|
||||
print("Dry-run terminé : aucune donnée écrite.")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
str(settings.database_url),
|
||||
pool_pre_ping=True,
|
||||
)
|
||||
|
||||
try:
|
||||
async with engine.begin() as connection:
|
||||
await upsert_sites(
|
||||
connection,
|
||||
sites,
|
||||
)
|
||||
|
||||
rows = build_reading_batch(all_readings)
|
||||
|
||||
if rows:
|
||||
await connection.execute(
|
||||
READING_INSERT,
|
||||
rows,
|
||||
)
|
||||
|
||||
finally:
|
||||
await engine.dispose()
|
||||
|
||||
print("Import API Mock terminé.")
|
||||
|
||||
|
||||
def parse_datetime(value: str) -> datetime:
|
||||
return datetime.fromisoformat(value.replace("Z", "+00:00"))
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(description=("Import historique depuis l'API Mock EnerVision"))
|
||||
|
||||
parser.add_argument(
|
||||
"--start-time",
|
||||
required=True,
|
||||
type=parse_datetime,
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--end-time",
|
||||
required=True,
|
||||
type=parse_datetime,
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--limit",
|
||||
type=int,
|
||||
default=1000,
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--dry-run",
|
||||
action="store_true",
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
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.start_time >= args.end_time:
|
||||
raise ValueError("--start-time doit être antérieur à --end-time.")
|
||||
|
||||
asyncio.run(
|
||||
import_mock_api_history(
|
||||
start_time=args.start_time,
|
||||
end_time=args.end_time,
|
||||
limit=args.limit,
|
||||
dry_run=args.dry_run,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -1,24 +1,11 @@
|
||||
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
|
||||
@@ -33,19 +20,3 @@ 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,9 +12,3 @@ class RecommendationResponse(BaseModel):
|
||||
explanation: str
|
||||
rule_reference: str
|
||||
created_at: datetime
|
||||
|
||||
|
||||
class RecommendationGenerationResponse(BaseModel):
|
||||
alerts_examined: int
|
||||
recommendations_created: int
|
||||
already_present: int
|
||||
|
||||
@@ -1,15 +1,7 @@
|
||||
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):
|
||||
@@ -20,24 +12,9 @@ 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,
|
||||
alerts: AlertRepository,
|
||||
transaction: Transaction,
|
||||
) -> None:
|
||||
def __init__(self, *, recommendations: RecommendationRepository) -> None:
|
||||
self._recommendations = recommendations
|
||||
self._alerts = alerts
|
||||
self._transaction = transaction
|
||||
|
||||
async def list_all(self) -> Sequence[Recommendation]:
|
||||
return await self._recommendations.list_all()
|
||||
@@ -47,16 +24,3 @@ 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,
|
||||
)
|
||||
|
||||
@@ -1,117 +0,0 @@
|
||||
# Piège : `rule_reference` est la clé d'idempotence en base, portée par la contrainte
|
||||
# `uq_recommendation_alert_rule`. Renommer une référence déjà livrée ne remplace pas les
|
||||
# recommandations existantes, il en crée de nouvelles à côté. Une règle qui change de sens
|
||||
# prend donc une référence suffixée `-v2` - REGLES.
|
||||
|
||||
from collections.abc import Callable
|
||||
from dataclasses import dataclass
|
||||
from typing import Final
|
||||
|
||||
from app.models.energy import Alert
|
||||
from app.repositories.recommendation import NouvelleRecommandation
|
||||
from app.schemas.alert import AlertSeverity, AlertType
|
||||
|
||||
FACTEUR_DEPASSEMENT_MAJEUR: Final = 1.2
|
||||
POURCENTAGE_DEPASSEMENT_MAJEUR: Final = round((FACTEUR_DEPASSEMENT_MAJEUR - 1) * 100)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Regle:
|
||||
reference: str
|
||||
action: str
|
||||
declencheur: Callable[[Alert], bool]
|
||||
motif: Callable[[Alert], str]
|
||||
|
||||
|
||||
def _du_type(attendu: AlertType) -> Callable[[Alert], bool]:
|
||||
return lambda alerte: alerte.type == attendu
|
||||
|
||||
|
||||
def _de_severite(attendue: AlertSeverity) -> Callable[[Alert], bool]:
|
||||
return lambda alerte: alerte.severity == attendue
|
||||
|
||||
|
||||
# Un seuil nul ou négatif rendrait le rapport `value / threshold` arbitraire : l'alerte ne
|
||||
# renseigne alors aucun dépassement exploitable, et la règle ne se déclenche pas.
|
||||
def _depasse_largement_le_seuil(alerte: Alert) -> bool:
|
||||
if alerte.value is None or alerte.threshold is None or alerte.threshold <= 0:
|
||||
return False
|
||||
return alerte.value >= alerte.threshold * FACTEUR_DEPASSEMENT_MAJEUR
|
||||
|
||||
|
||||
REGLES: Final[tuple[Regle, ...]] = (
|
||||
Regle(
|
||||
reference="spike-delestage-v1",
|
||||
action="Délester les équipements non prioritaires sur le créneau du pic",
|
||||
declencheur=_du_type(AlertType.SPIKE),
|
||||
motif=lambda alerte: f"Pic de consommation signalé sur le site {alerte.site_id}",
|
||||
),
|
||||
Regle(
|
||||
reference="threshold-reduction-v1",
|
||||
action="Ramener la puissance appelée sous le seuil contractuel",
|
||||
declencheur=_du_type(AlertType.THRESHOLD),
|
||||
motif=lambda alerte: f"Seuil de consommation dépassé sur le site {alerte.site_id}",
|
||||
),
|
||||
Regle(
|
||||
reference="outage-secours-v1",
|
||||
action="Basculer sur l'alimentation de secours et prévenir l'exploitant",
|
||||
declencheur=_du_type(AlertType.OUTAGE),
|
||||
motif=lambda alerte: (
|
||||
f"Risque de surcharge ou de coupure imminente sur le site {alerte.site_id}"
|
||||
),
|
||||
),
|
||||
Regle(
|
||||
reference="sensor-maintenance-v1",
|
||||
action="Planifier une intervention de maintenance sur le capteur",
|
||||
declencheur=_du_type(AlertType.SENSOR),
|
||||
motif=lambda alerte: (
|
||||
f"Capteur défaillant sur le site {alerte.site_id}, les mesures ne sont plus fiables"
|
||||
),
|
||||
),
|
||||
Regle(
|
||||
reference="anomaly-verification-v1",
|
||||
action="Confronter la mesure à la prévision et vérifier le paramétrage du site",
|
||||
declencheur=_du_type(AlertType.ANOMALY),
|
||||
motif=lambda alerte: (
|
||||
f"Écart anormal entre la mesure et le comportement attendu du site {alerte.site_id}"
|
||||
),
|
||||
),
|
||||
Regle(
|
||||
reference="escalade-astreinte-v1",
|
||||
action="Escalader à l'astreinte sous une heure",
|
||||
declencheur=_de_severite(AlertSeverity.CRITICAL),
|
||||
motif=lambda alerte: f"Alerte de sévérité critique sur le site {alerte.site_id}",
|
||||
),
|
||||
Regle(
|
||||
reference="contrat-puissance-v1",
|
||||
action="Réévaluer la puissance souscrite au contrat",
|
||||
declencheur=_depasse_largement_le_seuil,
|
||||
motif=lambda alerte: (
|
||||
f"Dépassement d'au moins {POURCENTAGE_DEPASSEMENT_MAJEUR} % du seuil "
|
||||
f"sur le site {alerte.site_id}"
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def applique_les_regles(alerte: Alert) -> list[NouvelleRecommandation]:
|
||||
contexte = _contexte_de_mesure(alerte)
|
||||
return [
|
||||
NouvelleRecommandation(
|
||||
alert_id=alerte.alert_id,
|
||||
action=regle.action,
|
||||
explanation=f"{regle.motif(alerte)}{contexte}.",
|
||||
rule_reference=regle.reference,
|
||||
)
|
||||
for regle in REGLES
|
||||
if regle.declencheur(alerte)
|
||||
]
|
||||
|
||||
|
||||
def _contexte_de_mesure(alerte: Alert) -> str:
|
||||
if alerte.value is None:
|
||||
return ""
|
||||
grandeur = alerte.metric or "valeur"
|
||||
if alerte.threshold is None:
|
||||
return f" ({grandeur} mesurée à {alerte.value})"
|
||||
return f" ({grandeur} mesurée à {alerte.value}, seuil {alerte.threshold})"
|
||||
+1
-108
@@ -1453,90 +1453,6 @@
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/recommendations/generate": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"recommendations"
|
||||
],
|
||||
"summary": "Génère les recommandations à partir des alertes",
|
||||
"operationId": "generate_recommendations_api_v1_recommendations_generate_post",
|
||||
"security": [
|
||||
{
|
||||
"Jeton d'accès": []
|
||||
}
|
||||
],
|
||||
"parameters": [
|
||||
{
|
||||
"name": "site_id",
|
||||
"in": "query",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "string"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Site Id"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "Successful Response",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/RecommendationGenerationResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"500": {
|
||||
"description": "Erreur interne. `correlation` identifie la trace côté serveur, qui n'est pas renvoyée au client.",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/InternalErrorResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"401": {
|
||||
"description": "Jeton absent, illisible, périmé, ou rendu caduc par un changement de rôle ou une désactivation. L'en-tête `WWW-Authenticate` porte la cause dans `error=`.",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"403": {
|
||||
"description": "Droits insuffisants, ou mot de passe provisoire à changer quand `detail` vaut `password_change_required`.",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"422": {
|
||||
"description": "Corps invalide. Le détail nomme le champ fautif et le type d'erreur, jamais la valeur envoyée.",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/ValidationErrorResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/stats/summary": {
|
||||
"get": {
|
||||
"tags": [
|
||||
@@ -2428,29 +2344,6 @@
|
||||
],
|
||||
"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": {
|
||||
@@ -3260,7 +3153,7 @@
|
||||
},
|
||||
{
|
||||
"name": "recommendations",
|
||||
"description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`. Leur génération par le moteur de règles est réservée au rôle `admin`."
|
||||
"description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`."
|
||||
},
|
||||
{
|
||||
"name": "stats",
|
||||
|
||||
@@ -17,6 +17,7 @@ dependencies = [
|
||||
"argon2-cffi>=23.1",
|
||||
"anyio>=4.0",
|
||||
"aiosmtplib>=5.1.3",
|
||||
"httpx>=0.28.1",
|
||||
"pandas>=3.0.5",
|
||||
]
|
||||
|
||||
@@ -27,7 +28,6 @@ dev = [
|
||||
"pytest>=9.1.1",
|
||||
"pytest-asyncio>=1.4.0",
|
||||
"pytest-cov>=7.1.0",
|
||||
"httpx>=0.28.1",
|
||||
"pandas-stubs>=3.0.5.260914",
|
||||
]
|
||||
|
||||
|
||||
@@ -51,7 +51,6 @@ 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,
|
||||
|
||||
@@ -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 RapportGeneration, RecommendationNotFoundError
|
||||
from app.services.recommendation import RecommendationNotFoundError
|
||||
|
||||
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
||||
|
||||
@@ -40,7 +40,6 @@ 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]
|
||||
@@ -50,10 +49,6 @@ 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]:
|
||||
@@ -147,56 +142,3 @@ 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
|
||||
|
||||
@@ -0,0 +1,622 @@
|
||||
import json
|
||||
import sys
|
||||
from datetime import datetime
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
from httpx import AsyncClient, MockTransport, Request, Response
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
import app.etl.mock_api_import as mock_api_import
|
||||
from app.etl.mock_api_import import (
|
||||
READING_INSERT,
|
||||
SOURCE_HISTORY,
|
||||
build_reading_batch,
|
||||
build_reading_row,
|
||||
fetch_readings,
|
||||
fetch_sites,
|
||||
upsert_sites,
|
||||
)
|
||||
|
||||
|
||||
def make_site() -> dict[str, Any]:
|
||||
return {
|
||||
"site_id": "SITE001",
|
||||
"site_type": "office",
|
||||
"site_name": "Bureau Paris La Défense",
|
||||
"location": "Paris, France",
|
||||
"capacity_kw": 200,
|
||||
"status": "active",
|
||||
}
|
||||
|
||||
|
||||
def make_reading() -> dict[str, Any]:
|
||||
return {
|
||||
"timestamp": "2024-06-15T12:00:00Z",
|
||||
"site_id": "SITE001",
|
||||
"site_type": "office",
|
||||
"consumption_kw": 87.34,
|
||||
"consumption_kwh": 87.34,
|
||||
"voltage_v": 401.2,
|
||||
"current_a": 132.5,
|
||||
"power_factor": 0.923,
|
||||
"temperature_celsius": 22.1,
|
||||
"humidity_percent": 58.4,
|
||||
"null_reasons": [],
|
||||
"data_quality": "good",
|
||||
}
|
||||
|
||||
|
||||
async def test_fetch_sites_returns_sites() -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
assert request.url.path == "/api/v1/sites"
|
||||
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_site()],
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
sites = await fetch_sites(client)
|
||||
|
||||
assert len(sites) == 1
|
||||
assert sites[0]["site_id"] == "SITE001"
|
||||
assert sites[0]["site_type"] == "office"
|
||||
|
||||
|
||||
async def test_fetch_readings_sends_expected_query_parameters() -> None:
|
||||
captured_params: dict[str, str] = {}
|
||||
|
||||
def handler(request: Request) -> Response:
|
||||
nonlocal captured_params
|
||||
|
||||
captured_params = dict(request.url.params)
|
||||
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_reading()],
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
start_time = datetime.fromisoformat("2024-06-15T12:00:00")
|
||||
end_time = datetime.fromisoformat("2024-06-15T13:00:00")
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
readings = await fetch_readings(
|
||||
client=client,
|
||||
site_id="SITE001",
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=60,
|
||||
)
|
||||
|
||||
assert len(readings) == 1
|
||||
assert captured_params["site_id"] == "SITE001"
|
||||
assert captured_params["start_time"] == "2024-06-15T12:00:00"
|
||||
assert captured_params["end_time"] == "2024-06-15T13:00:00"
|
||||
assert captured_params["limit"] == "60"
|
||||
|
||||
|
||||
async def test_fetch_readings_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/readings doit être une liste",
|
||||
):
|
||||
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=60,
|
||||
)
|
||||
|
||||
|
||||
async def test_fetch_readings_raises_on_http_error() -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
return Response(
|
||||
status_code=404,
|
||||
json={"detail": "Site non trouvé"},
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
with pytest.raises(httpx.HTTPStatusError):
|
||||
await fetch_readings(
|
||||
client=client,
|
||||
site_id="SITE999",
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
)
|
||||
|
||||
|
||||
def test_build_reading_row_respects_database_contract() -> None:
|
||||
reading = make_reading()
|
||||
|
||||
row = build_reading_row(reading)
|
||||
|
||||
assert row["site_id"] == "SITE001"
|
||||
assert row["source"] == SOURCE_HISTORY
|
||||
assert row["source"] == "api_history"
|
||||
assert row["dataset_id"] is None
|
||||
|
||||
assert row["timestamp"] == datetime.fromisoformat("2024-06-15T12:00:00+00:00")
|
||||
|
||||
assert row["consumption_kw"] == 87.34
|
||||
assert row["consumption_kwh"] == 87.34
|
||||
assert row["data_quality"] == "good"
|
||||
assert row["null_reasons"] == []
|
||||
|
||||
assert row["imputed_values"] is None
|
||||
assert row["imputation_method"] is None
|
||||
|
||||
|
||||
def test_build_reading_row_keeps_null_values_and_quality() -> None:
|
||||
reading = make_reading()
|
||||
|
||||
reading["consumption_kw"] = None
|
||||
reading["consumption_kwh"] = None
|
||||
reading["voltage_v"] = None
|
||||
reading["current_a"] = None
|
||||
reading["power_factor"] = None
|
||||
reading["data_quality"] = "degraded"
|
||||
reading["null_reasons"] = [
|
||||
"consumption_sensor_failure",
|
||||
"electrical_sensor_failure",
|
||||
]
|
||||
|
||||
row = build_reading_row(reading)
|
||||
|
||||
assert row["consumption_kw"] is None
|
||||
assert row["consumption_kwh"] is None
|
||||
assert row["voltage_v"] is None
|
||||
assert row["current_a"] is None
|
||||
assert row["power_factor"] is None
|
||||
|
||||
assert row["data_quality"] == "degraded"
|
||||
assert row["null_reasons"] == [
|
||||
"consumption_sensor_failure",
|
||||
"electrical_sensor_failure",
|
||||
]
|
||||
|
||||
assert row["imputed_values"] is None
|
||||
assert row["imputation_method"] is None
|
||||
|
||||
|
||||
def test_build_reading_row_keeps_raw_source_data() -> None:
|
||||
reading = make_reading()
|
||||
|
||||
row = build_reading_row(reading)
|
||||
|
||||
raw_data = json.loads(row["raw_data"])
|
||||
|
||||
assert raw_data == reading
|
||||
|
||||
|
||||
def test_build_reading_batch_transforms_all_readings() -> None:
|
||||
first = make_reading()
|
||||
|
||||
second = make_reading()
|
||||
second["timestamp"] = "2024-06-15T12:01:00Z"
|
||||
second["consumption_kw"] = 90.5
|
||||
|
||||
rows = build_reading_batch([first, second])
|
||||
|
||||
assert len(rows) == 2
|
||||
|
||||
assert rows[0]["site_id"] == "SITE001"
|
||||
assert rows[0]["consumption_kw"] == 87.34
|
||||
|
||||
assert rows[1]["site_id"] == "SITE001"
|
||||
assert rows[1]["consumption_kw"] == 90.5
|
||||
|
||||
|
||||
def test_create_mock_api_client_requires_credentials(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
settings = SimpleNamespace(
|
||||
mock_api_username=None,
|
||||
mock_api_password=None,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: settings,
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="Les identifiants de l'API Mock ne sont pas configurés",
|
||||
):
|
||||
mock_api_import.create_mock_api_client()
|
||||
|
||||
|
||||
async def test_create_mock_api_client_uses_configuration(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
password = MagicMock()
|
||||
password.get_secret_value.return_value = "test-password"
|
||||
|
||||
settings = SimpleNamespace(
|
||||
mock_api_base_url="https://mock.test/",
|
||||
mock_api_username="test-user",
|
||||
mock_api_password=password,
|
||||
mock_api_timeout_seconds=10.0,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: settings,
|
||||
)
|
||||
|
||||
client = mock_api_import.create_mock_api_client()
|
||||
|
||||
try:
|
||||
assert str(client.base_url) == "https://mock.test"
|
||||
assert client.timeout.connect == 10.0
|
||||
finally:
|
||||
await client.aclose()
|
||||
|
||||
|
||||
async def test_upsert_sites_with_empty_list_does_nothing() -> None:
|
||||
connection = AsyncMock()
|
||||
|
||||
await upsert_sites(
|
||||
connection,
|
||||
[],
|
||||
)
|
||||
|
||||
connection.execute.assert_not_awaited()
|
||||
|
||||
|
||||
async def test_import_mock_api_history_dry_run_does_not_write(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
if request.url.path == "/api/v1/sites":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_site()],
|
||||
)
|
||||
|
||||
if request.url.path == "/api/v1/readings":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_reading()],
|
||||
)
|
||||
|
||||
return Response(status_code=404)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
client = AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_mock_api_client",
|
||||
lambda: client,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: SimpleNamespace(
|
||||
database_url="postgresql+asyncpg://unused",
|
||||
),
|
||||
)
|
||||
|
||||
create_engine_mock = MagicMock()
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_async_engine",
|
||||
create_engine_mock,
|
||||
)
|
||||
|
||||
await mock_api_import.import_mock_api_history(
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
dry_run=True,
|
||||
)
|
||||
|
||||
create_engine_mock.assert_not_called()
|
||||
|
||||
|
||||
async def test_import_mock_api_history_loads_data(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
if request.url.path == "/api/v1/sites":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_site()],
|
||||
)
|
||||
|
||||
if request.url.path == "/api/v1/readings":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_reading()],
|
||||
)
|
||||
|
||||
return Response(status_code=404)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
client = AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_mock_api_client",
|
||||
lambda: client,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: SimpleNamespace(
|
||||
database_url="postgresql+asyncpg://test:test@localhost/test",
|
||||
),
|
||||
)
|
||||
|
||||
connection = AsyncMock()
|
||||
|
||||
transaction_context = MagicMock()
|
||||
transaction_context.__aenter__ = AsyncMock(
|
||||
return_value=connection,
|
||||
)
|
||||
transaction_context.__aexit__ = AsyncMock(
|
||||
return_value=None,
|
||||
)
|
||||
|
||||
engine = MagicMock()
|
||||
engine.begin.return_value = transaction_context
|
||||
engine.dispose = AsyncMock()
|
||||
|
||||
create_engine_mock = MagicMock(
|
||||
return_value=engine,
|
||||
)
|
||||
|
||||
upsert_sites_mock = AsyncMock()
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_async_engine",
|
||||
create_engine_mock,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"upsert_sites",
|
||||
upsert_sites_mock,
|
||||
)
|
||||
|
||||
await mock_api_import.import_mock_api_history(
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
dry_run=False,
|
||||
)
|
||||
|
||||
create_engine_mock.assert_called_once_with(
|
||||
"postgresql+asyncpg://test:test@localhost/test",
|
||||
pool_pre_ping=True,
|
||||
)
|
||||
|
||||
upsert_sites_mock.assert_awaited_once_with(
|
||||
connection,
|
||||
[make_site()],
|
||||
)
|
||||
|
||||
connection.execute.assert_awaited_once()
|
||||
engine.dispose.assert_awaited_once()
|
||||
|
||||
|
||||
def test_parse_datetime_accepts_z_suffix() -> None:
|
||||
result = mock_api_import.parse_datetime(
|
||||
"2024-06-15T12:00:00Z",
|
||||
)
|
||||
|
||||
assert result == datetime.fromisoformat(
|
||||
"2024-06-15T12:00:00+00:00",
|
||||
)
|
||||
|
||||
|
||||
def test_parse_args_reads_cli_parameters(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(
|
||||
sys,
|
||||
"argv",
|
||||
[
|
||||
"mock_api_import",
|
||||
"--start-time",
|
||||
"2024-06-15T12:00:00Z",
|
||||
"--end-time",
|
||||
"2024-06-15T13:00:00Z",
|
||||
"--limit",
|
||||
"60",
|
||||
"--dry-run",
|
||||
],
|
||||
)
|
||||
|
||||
args = mock_api_import.parse_args()
|
||||
|
||||
assert args.start_time == datetime.fromisoformat(
|
||||
"2024-06-15T12:00:00+00:00",
|
||||
)
|
||||
assert args.end_time == datetime.fromisoformat(
|
||||
"2024-06-15T13:00:00+00:00",
|
||||
)
|
||||
assert args.limit == 60
|
||||
assert args.dry_run is True
|
||||
|
||||
|
||||
def test_main_rejects_limit_out_of_bounds(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(
|
||||
sys,
|
||||
"argv",
|
||||
[
|
||||
"mock_api_import",
|
||||
"--start-time",
|
||||
"2024-06-15T12:00:00Z",
|
||||
"--end-time",
|
||||
"2024-06-15T13:00:00Z",
|
||||
"--limit",
|
||||
"0",
|
||||
],
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="--limit doit être compris entre 1 et 1000",
|
||||
):
|
||||
mock_api_import.main()
|
||||
|
||||
|
||||
def test_main_rejects_invalid_period(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(
|
||||
sys,
|
||||
"argv",
|
||||
[
|
||||
"mock_api_import",
|
||||
"--start-time",
|
||||
"2024-06-15T14:00:00Z",
|
||||
"--end-time",
|
||||
"2024-06-15T13:00:00Z",
|
||||
"--limit",
|
||||
"60",
|
||||
],
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="--start-time doit être antérieur à --end-time",
|
||||
):
|
||||
mock_api_import.main()
|
||||
|
||||
|
||||
def test_main_runs_import(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
start_time = datetime.fromisoformat(
|
||||
"2024-06-15T12:00:00+00:00",
|
||||
)
|
||||
end_time = datetime.fromisoformat(
|
||||
"2024-06-15T13:00:00+00:00",
|
||||
)
|
||||
|
||||
import_mock = AsyncMock()
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"parse_args",
|
||||
lambda: SimpleNamespace(
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=60,
|
||||
dry_run=True,
|
||||
),
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"import_mock_api_history",
|
||||
import_mock,
|
||||
)
|
||||
|
||||
mock_api_import.main()
|
||||
|
||||
import_mock.assert_awaited_once_with(
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=60,
|
||||
dry_run=True,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.integration
|
||||
async def test_reading_insert_is_idempotent(
|
||||
session: AsyncSession,
|
||||
) -> None:
|
||||
reading = make_reading()
|
||||
row = build_reading_row(reading)
|
||||
|
||||
connection = await session.connection()
|
||||
|
||||
await upsert_sites(
|
||||
connection,
|
||||
[make_site()],
|
||||
)
|
||||
|
||||
await session.execute(
|
||||
READING_INSERT,
|
||||
[row],
|
||||
)
|
||||
|
||||
await session.execute(
|
||||
READING_INSERT,
|
||||
[row],
|
||||
)
|
||||
|
||||
result = await session.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT COUNT(*)
|
||||
FROM reading
|
||||
WHERE site_id = :site_id
|
||||
AND timestamp = :timestamp
|
||||
AND source = :source
|
||||
"""
|
||||
),
|
||||
{
|
||||
"site_id": row["site_id"],
|
||||
"timestamp": row["timestamp"],
|
||||
"source": row["source"],
|
||||
},
|
||||
)
|
||||
|
||||
assert result.scalar_one() == 1
|
||||
|
||||
await session.rollback()
|
||||
@@ -5,8 +5,7 @@ import pytest
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.energy import Alert, Recommendation, Site
|
||||
from app.repositories import recommendation as module_recommendation
|
||||
from app.repositories.recommendation import NouvelleRecommandation, RecommendationRepository
|
||||
from app.repositories.recommendation import RecommendationRepository
|
||||
|
||||
pytestmark = pytest.mark.integration
|
||||
|
||||
@@ -84,59 +83,3 @@ 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,14 +1,10 @@
|
||||
from collections.abc import Sequence
|
||||
from datetime import UTC, datetime
|
||||
|
||||
import pytest
|
||||
|
||||
from app.models.energy import Alert, Recommendation
|
||||
from app.repositories.recommendation import NouvelleRecommandation
|
||||
from app.models.energy import Recommendation
|
||||
from app.services.recommendation import RecommendationNotFoundError, RecommendationService
|
||||
|
||||
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
||||
|
||||
|
||||
def recommendation(recommendation_id: int = 1) -> Recommendation:
|
||||
return Recommendation(
|
||||
@@ -17,33 +13,13 @@ def recommendation(recommendation_id: int = 1) -> Recommendation:
|
||||
action="Vérifier la consommation",
|
||||
explanation="Pic détecté",
|
||||
rule_reference="spike-v1",
|
||||
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={},
|
||||
created_at=datetime(2024, 1, 1, tzinfo=UTC),
|
||||
)
|
||||
|
||||
|
||||
class FakeRepository:
|
||||
def __init__(self, recommendations: list[Recommendation], creees: int | None = None) -> None:
|
||||
def __init__(self, recommendations: list[Recommendation]) -> None:
|
||||
self._recommendations = recommendations
|
||||
self._creees = creees
|
||||
self.recues: list[NouvelleRecommandation] = []
|
||||
|
||||
async def list_all(self) -> list[Recommendation]:
|
||||
return self._recommendations
|
||||
@@ -53,111 +29,27 @@ 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
|
||||
|
||||
|
||||
class FakeAlertRepository:
|
||||
def __init__(self, alertes: list[Alert]) -> None:
|
||||
self._alertes = alertes
|
||||
self.site_demande: str | None = None
|
||||
|
||||
async def list_all(
|
||||
self, *, site_id: str | None = None, severity: str | None = None
|
||||
) -> list[Alert]:
|
||||
self.site_demande = site_id
|
||||
if site_id is None:
|
||||
return self._alertes
|
||||
return [a for a in self._alertes if a.site_id == site_id]
|
||||
|
||||
|
||||
class FakeTransaction:
|
||||
def __init__(self) -> None:
|
||||
self.commits = 0
|
||||
|
||||
async def commit(self) -> None:
|
||||
self.commits += 1
|
||||
|
||||
|
||||
def service(
|
||||
recommendations: FakeRepository | None = None,
|
||||
alerts: FakeAlertRepository | None = None,
|
||||
transaction: FakeTransaction | None = None,
|
||||
) -> RecommendationService:
|
||||
return RecommendationService(
|
||||
recommendations=recommendations or FakeRepository([]),
|
||||
alerts=alerts or FakeAlertRepository([]),
|
||||
transaction=transaction or FakeTransaction(),
|
||||
)
|
||||
|
||||
|
||||
async def test_list_all_returns_the_repository_recommendations() -> None:
|
||||
depot = FakeRepository([recommendation(1), recommendation(2)])
|
||||
service = RecommendationService(
|
||||
recommendations=FakeRepository([recommendation(1), recommendation(2)])
|
||||
)
|
||||
|
||||
recommendations = await service(recommendations=depot).list_all()
|
||||
recommendations = await service.list_all()
|
||||
|
||||
assert [r.recommendation_id for r in recommendations] == [1, 2]
|
||||
|
||||
|
||||
async def test_get_by_id_returns_the_matching_recommendation() -> None:
|
||||
trouve = await service(recommendations=FakeRepository([recommendation(1)])).get_by_id(1)
|
||||
service = RecommendationService(recommendations=FakeRepository([recommendation(1)]))
|
||||
|
||||
trouve = await service.get_by_id(1)
|
||||
|
||||
assert trouve.recommendation_id == 1
|
||||
|
||||
|
||||
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)
|
||||
|
||||
|
||||
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
|
||||
await service.get_by_id(404)
|
||||
|
||||
@@ -1,142 +0,0 @@
|
||||
from datetime import UTC, datetime
|
||||
|
||||
import pytest
|
||||
|
||||
from app.models.energy import Alert
|
||||
from app.services.recommendation_rules import FACTEUR_DEPASSEMENT_MAJEUR, applique_les_regles
|
||||
|
||||
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
||||
|
||||
|
||||
def alerte(
|
||||
*,
|
||||
alert_id: int = 1,
|
||||
type_alerte: str = "spike",
|
||||
severity: str = "high",
|
||||
value: float | None = None,
|
||||
threshold: float | None = None,
|
||||
metric: str | None = None,
|
||||
site_id: str = "SITE001",
|
||||
) -> Alert:
|
||||
return Alert(
|
||||
alert_id=alert_id,
|
||||
source_alert_id=f"ALR-{alert_id}",
|
||||
site_id=site_id,
|
||||
source="api_mock",
|
||||
timestamp=MOMENT,
|
||||
type=type_alerte,
|
||||
severity=severity,
|
||||
message="Alerte de test",
|
||||
value=value,
|
||||
threshold=threshold,
|
||||
metric=metric,
|
||||
prediction_id=None,
|
||||
raw_data={},
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("type_alerte", "attendue"),
|
||||
[
|
||||
("spike", "spike-delestage-v1"),
|
||||
("threshold", "threshold-reduction-v1"),
|
||||
("outage", "outage-secours-v1"),
|
||||
("sensor", "sensor-maintenance-v1"),
|
||||
("anomaly", "anomaly-verification-v1"),
|
||||
],
|
||||
ids=["pic", "seuil", "coupure", "capteur", "anomalie"],
|
||||
)
|
||||
def test_each_alert_type_yields_its_own_rule(type_alerte: str, attendue: str) -> None:
|
||||
proposees = applique_les_regles(alerte(type_alerte=type_alerte))
|
||||
|
||||
assert [p.rule_reference for p in proposees] == [attendue]
|
||||
|
||||
|
||||
def test_a_critical_alert_adds_the_escalation_rule() -> None:
|
||||
proposees = applique_les_regles(alerte(severity="critical"))
|
||||
|
||||
assert "escalade-astreinte-v1" in {p.rule_reference for p in proposees}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("severity", ["low", "medium", "high"], ids=["faible", "moyenne", "haute"])
|
||||
def test_a_non_critical_alert_does_not_escalate(severity: str) -> None:
|
||||
proposees = applique_les_regles(alerte(severity=severity))
|
||||
|
||||
assert "escalade-astreinte-v1" not in {p.rule_reference for p in proposees}
|
||||
|
||||
|
||||
def test_a_large_overshoot_adds_the_contract_rule() -> None:
|
||||
proposees = applique_les_regles(
|
||||
alerte(value=720.0 * FACTEUR_DEPASSEMENT_MAJEUR, threshold=720.0)
|
||||
)
|
||||
|
||||
assert "contrat-puissance-v1" in {p.rule_reference for p in proposees}
|
||||
|
||||
|
||||
def test_an_overshoot_below_the_factor_does_not_add_the_contract_rule() -> None:
|
||||
proposees = applique_les_regles(alerte(value=800.0, threshold=720.0))
|
||||
|
||||
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("value", "threshold"),
|
||||
[(None, 720.0), (900.0, None), (900.0, 0.0), (900.0, -10.0)],
|
||||
ids=["sans mesure", "sans seuil", "seuil nul", "seuil negatif"],
|
||||
)
|
||||
def test_the_contract_rule_stays_silent_without_an_exploitable_threshold(
|
||||
value: float | None, threshold: float | None
|
||||
) -> None:
|
||||
proposees = applique_les_regles(alerte(value=value, threshold=threshold))
|
||||
|
||||
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
|
||||
|
||||
|
||||
def test_the_explanation_quotes_the_measure_and_the_threshold() -> None:
|
||||
proposees = applique_les_regles(alerte(value=812.5, threshold=720.0, metric="consumption_kw"))
|
||||
|
||||
assert "(consumption_kw mesurée à 812.5, seuil 720.0)" in proposees[0].explanation
|
||||
|
||||
|
||||
def test_the_explanation_quotes_the_measure_alone_when_no_threshold_is_known() -> None:
|
||||
proposees = applique_les_regles(alerte(value=812.5, metric="consumption_kw"))
|
||||
|
||||
assert "(consumption_kw mesurée à 812.5)" in proposees[0].explanation
|
||||
|
||||
|
||||
def test_the_explanation_omits_the_measure_when_the_alert_carries_none() -> None:
|
||||
proposees = applique_les_regles(alerte())
|
||||
|
||||
assert "(" not in proposees[0].explanation
|
||||
|
||||
|
||||
def test_the_explanation_names_the_site() -> None:
|
||||
proposees = applique_les_regles(alerte(site_id="SITE042"))
|
||||
|
||||
assert "SITE042" in proposees[0].explanation
|
||||
|
||||
|
||||
def test_every_proposal_carries_the_alert_identifier() -> None:
|
||||
proposees = applique_les_regles(alerte(alert_id=77, severity="critical"))
|
||||
|
||||
assert {p.alert_id for p in proposees} == {77}
|
||||
|
||||
|
||||
def test_an_alert_never_yields_the_same_rule_twice() -> None:
|
||||
proposees = applique_les_regles(
|
||||
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
|
||||
)
|
||||
|
||||
assert len(proposees) == len({p.rule_reference for p in proposees})
|
||||
|
||||
|
||||
def test_a_critical_alert_over_the_threshold_yields_the_three_rules() -> None:
|
||||
proposees = applique_les_regles(
|
||||
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
|
||||
)
|
||||
|
||||
assert {p.rule_reference for p in proposees} == {
|
||||
"spike-delestage-v1",
|
||||
"escalade-astreinte-v1",
|
||||
"contrat-puissance-v1",
|
||||
}
|
||||
@@ -118,30 +118,3 @@ def test_main_exports_the_contract_without_asking_for_a_password(
|
||||
assert code == 0
|
||||
assert 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
|
||||
|
||||
Generated
+2
-2
@@ -326,6 +326,7 @@ dependencies = [
|
||||
{ name = "argon2-cffi" },
|
||||
{ name = "asyncpg" },
|
||||
{ name = "fastapi" },
|
||||
{ name = "httpx" },
|
||||
{ name = "pandas" },
|
||||
{ name = "prometheus-fastapi-instrumentator" },
|
||||
{ name = "pydantic", extra = ["email"] },
|
||||
@@ -338,7 +339,6 @@ dependencies = [
|
||||
|
||||
[package.dev-dependencies]
|
||||
dev = [
|
||||
{ name = "httpx" },
|
||||
{ name = "mypy" },
|
||||
{ name = "pandas-stubs" },
|
||||
{ name = "pytest" },
|
||||
@@ -355,6 +355,7 @@ requires-dist = [
|
||||
{ name = "argon2-cffi", specifier = ">=23.1" },
|
||||
{ name = "asyncpg", specifier = ">=0.31.0" },
|
||||
{ name = "fastapi", specifier = ">=0.141.1" },
|
||||
{ name = "httpx", specifier = ">=0.28.1" },
|
||||
{ name = "pandas", specifier = ">=3.0.5" },
|
||||
{ name = "prometheus-fastapi-instrumentator", specifier = ">=8.1.0" },
|
||||
{ name = "pydantic", extras = ["email"], specifier = ">=2.13.5" },
|
||||
@@ -367,7 +368,6 @@ requires-dist = [
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
dev = [
|
||||
{ name = "httpx", specifier = ">=0.28.1" },
|
||||
{ name = "mypy", specifier = ">=2.3.1" },
|
||||
{ name = "pandas-stubs", specifier = ">=3.0.5.260914" },
|
||||
{ name = "pytest", specifier = ">=9.1.1" },
|
||||
|
||||
@@ -50,6 +50,12 @@ services:
|
||||
APP_SECRET_KEY: ${APP_SECRET_KEY:?}
|
||||
APP_CORS_ORIGINS: ${APP_CORS_ORIGINS:-http://localhost:4200}
|
||||
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
||||
|
||||
APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr}
|
||||
APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-}
|
||||
APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-}
|
||||
APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10}
|
||||
|
||||
APP_FRONTEND_RESET_PASSWORD_URL: ${APP_FRONTEND_RESET_PASSWORD_URL:-http://localhost:4200/reset-password}
|
||||
APP_SMTP_HOST: mailpit
|
||||
APP_SMTP_PORT: "1025"
|
||||
|
||||
@@ -1,81 +0,0 @@
|
||||
# 0006 - Le moteur de règles de recommandation vit dans le backend
|
||||
|
||||
- Statut : accepté
|
||||
- Date : 2026-09-18
|
||||
|
||||
## Contexte
|
||||
|
||||
L'issue #38 demande un « moteur de règles pour recommandations », portée par le label `ml`. Le
|
||||
schéma tranche déjà la forme du résultat : `recommendation(alert_id, action, explanation,
|
||||
rule_reference)`, avec `alert_id` en clé étrangère `NOT NULL` et une contrainte d'unicité
|
||||
`uq_recommendation_alert_rule` sur `(alert_id, rule_reference)`. Une recommandation est donc
|
||||
**dérivée d'une alerte**, jamais d'une mesure brute ni d'une prévision.
|
||||
|
||||
Deux emplacements se disputaient le code :
|
||||
|
||||
1. `ml/enervision_ml/`, sur le patron de `enervision_ml.score` livré par #37 : un script autonome
|
||||
qui se connecte par `ML_DATABASE_URL`, écrit une table, et que l'API se contente de lire.
|
||||
L'[ADR 0005](0005-modele-prediction-lightgbm.md) annonce d'ailleurs #38 de ce côté, en écrivant
|
||||
que le scoring, le moteur de recommandations et les tests de dérive « consommeront le même
|
||||
module `enervision_ml.features` ».
|
||||
2. `apps/backend/app/services/`, où `apps/backend/README.md` place les « regles metier ».
|
||||
|
||||
## Décision
|
||||
|
||||
**Le moteur vit dans `apps/backend/app/services/`**, sous la forme d'un module pur
|
||||
`recommendation_rules.py` (le catalogue `REGLES`) et d'une méthode `RecommendationService.generate()`
|
||||
qui l'applique, persiste et valide la transaction.
|
||||
|
||||
Trois raisons :
|
||||
|
||||
- **Il n'utilise rien du ML.** Le catalogue lit `alert.type`, `alert.severity`, `alert.value` et
|
||||
`alert.threshold`. Aucun modèle, aucune feature, aucun `enervision_ml.features` : la phrase de
|
||||
l'ADR 0005 vaut pour le scoring (#37) et les tests de dérive (#44/#45), qui manipulent bien des
|
||||
features, pas pour des règles sur alertes. Le label `ml` de #38 désigne le lot fonctionnel
|
||||
« prédiction et recommandation », pas l'emplacement du code.
|
||||
- **Il lit et écrit deux tables déjà couvertes par des repositories.** `AlertRepository` sait déjà
|
||||
filtrer par site. Le placer dans `ml/` obligerait à réécrire ces accès en SQL brut, et à
|
||||
maintenir deux représentations du même domaine.
|
||||
- **Le déclencheur HTTP n'a de sens que dans l'API.** `POST /recommendations/generate` doit passer
|
||||
par `require_role(Role.ADMIN)` et par la session injectée : cela suppose d'être dans
|
||||
l'application FastAPI.
|
||||
|
||||
Le moteur reste néanmoins **déclenchable hors HTTP**, par `python -m app.cli
|
||||
generate-recommendations` (cible `make recommendations`), sur le patron de `make ml-score` : rien
|
||||
n'oblige à exposer un port pour régénérer des recommandations.
|
||||
|
||||
## Conséquences
|
||||
|
||||
- L'API gagne sa première route d'écriture métier. La checklist de `20-backend.md` s'applique :
|
||||
entrée dans `ROLE_MINIMUM` de `tests/api/acces.py`, et `openapi.json` régénéré dans le même
|
||||
commit.
|
||||
- `RecommendationService` n'est plus en lecture seule : il reçoit le `Transaction` Protocol déjà
|
||||
utilisé par `AuthService` et `UserService`, et commite lui-même. Les repositories continuent de
|
||||
ne pas commiter.
|
||||
- **L'idempotence est déléguée à la base.** `create_missing()` insère en `ON CONFLICT DO NOTHING`
|
||||
sur `uq_recommendation_alert_rule` plutôt que de relire avant d'écrire, ce qui supprime la
|
||||
fenêtre entre le contrôle et l'insertion. Corollaire : `rule_reference` est une clé fonctionnelle.
|
||||
Une règle dont le sens change prend une référence `-v2` ; renommer une référence livrée
|
||||
ferait réapparaître ses recommandations à côté des anciennes.
|
||||
- **Le moteur est branché sur la détection interne, et sur elle seule.** `alert` est alimentée
|
||||
par `app/detection/internal_alerts.py` (#104), lancée à la main comme `enervision_ml.score` ;
|
||||
l'ingestion de l'API Mock `/alerts` reste à faire. Le rapport de génération est donc à zéro tant
|
||||
que la détection n'a pas tourné, sans que le moteur soit à retoucher.
|
||||
- **L'insertion est découpée en lots.** `create_missing()` écrit par paquets de `TAILLE_DE_LOT`
|
||||
lignes : asyncpg plafonne une requête à 32 767 paramètres, soit 8 191 lignes de quatre colonnes,
|
||||
et la détection interne peut alimenter `alert` au fil de l'eau.
|
||||
- Si le projet devait un jour pondérer les recommandations par un score appris, la décision serait
|
||||
à rouvrir : le moteur redeviendrait consommateur du pipeline ML.
|
||||
|
||||
## Alternatives écartées
|
||||
|
||||
- **Module et CLI dans `ml/enervision_ml/`** : cohérent avec le label `ml` et avec la lettre de
|
||||
l'ADR 0005, mais impose du SQL brut là où deux repositories existent, et laisse la génération
|
||||
hors de portée de l'API. Redeviendrait le bon choix si les règles se mettaient à consommer des
|
||||
features ou un modèle.
|
||||
- **Génération à la volée, sans persistance**, calculée à chaque `GET /recommendations` : supprime
|
||||
le besoin d'écriture, mais rend la table `recommendation` et sa contrainte d'unicité inutiles,
|
||||
et interdit toute trace de ce qui a été proposé et quand.
|
||||
- **Table de configuration des règles en base**, plutôt qu'un catalogue en Python : plus souple,
|
||||
mais déplace la logique métier hors de la revue de code et hors des tests, pour un besoin que
|
||||
rien n'exprime à ce stade.
|
||||
@@ -165,8 +165,3 @@ Elles vivent dans `../adr/`, pas ici.
|
||||
| ADR | Objet |
|
||||
|---|---|
|
||||
| [0001](../adr/0001-postgresql-timescaledb.md) | PostgreSQL 17 avec l'extension TimescaleDB, et la frontière `db/` vs `alembic/` |
|
||||
| [0002](../adr/0002-authentification-jwt-et-refresh-opaque.md) | Authentification par JWT d'accès et jeton de rafraîchissement opaque |
|
||||
| [0003](../adr/0003-autorisation-rbac-a-trois-roles.md) | Autorisation RBAC à trois rôles, avec relecture du compte à chaque requête |
|
||||
| [0004](../adr/0004-journal-d-audit-en-ajout-seul.md) | Journal d'audit en ajout seul, garanti par PostgreSQL |
|
||||
| [0005](../adr/0005-modele-prediction-lightgbm.md) | Modèle de prédiction de consommation : LightGBM |
|
||||
| [0006](../adr/0006-moteur-de-regles-dans-le-backend.md) | Le moteur de règles de recommandation vit dans le backend, pas dans `ml/` |
|
||||
|
||||
@@ -146,7 +146,6 @@ Deux fichiers d'environnement, deux usages : `.env` à la racine alimente `docke
|
||||
| GET | `/api/v1/alerts` | Liste les alertes, filtrable par `site_id` et `severity`. `lecteur` | 401, 403, 422, 500 |
|
||||
| GET | `/api/v1/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 |
|
||||
@@ -195,21 +194,6 @@ plutôt qu'un statut inventé : le domaine `available`/`insufficient_data`/`erro
|
||||
LightGBM elle-même ; elle lit ce que le pipeline de scoring a déjà écrit, cf.
|
||||
[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
|
||||
|
||||
+331
-74
@@ -6,10 +6,14 @@ système qui en découle.
|
||||
|
||||
## Ce que couvre ce document
|
||||
|
||||
**Dix tables applicatives existent** : quatre pour l'authentification, six pour les données
|
||||
d'énergie, dont l'hypertable `reading`. Les sections marquées `Fait` relèvent le code. Celles
|
||||
marquées `Cible` décrivent ce qui n'est pas écrit, au premier rang desquelles la chaîne
|
||||
d'ingestion, les agrégats continus, la compression et la rétention.
|
||||
**Douze tables applicatives existent** : six pour l'authentification et six pour les données
|
||||
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.
|
||||
|
||||
## Trois emplacements, trois rôles
|
||||
|
||||
@@ -35,8 +39,16 @@ Statut : `Fait`.
|
||||
- `db/init/100-extensions.sql` crée l'extension `timescaledb`.
|
||||
- `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par
|
||||
`apps/backend/tests/conftest.py`.
|
||||
- Cinq révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle
|
||||
établit `alembic_version` et refuse de s'appliquer si l'extension manque :
|
||||
- Six révisions Alembic sont actuellement appliquées.
|
||||
- La première, `5353c0e4f094`, **ne crée aucune table** : elle établit `alembic_version`
|
||||
et refuse de s'appliquer si l'extension TimescaleDB manque.
|
||||
- Les révisions suivantes créent les tables liées à l'authentification :
|
||||
`app_user`, `login_attempt`, `audit_log` et `refresh_token`.
|
||||
- La révision `e6d2026091501` crée les six tables Data et déclare l'hypertable `reading`.
|
||||
- La révision `c0adab96238c` ajoute les tables `password_reset_attempt`
|
||||
et `password_reset_token`.
|
||||
|
||||
La garde de la première migration est :
|
||||
|
||||
```sql
|
||||
IF NOT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'timescaledb') THEN
|
||||
@@ -47,37 +59,58 @@ END IF;
|
||||
Cette garde forme paire avec le 503 de `/api/v1/health/ready`. Un bootstrap sauté ne se voit pas
|
||||
au démarrage de l'API : ces deux gardes le rendent visible tôt, des deux côtés.
|
||||
|
||||
Les trois suivantes créent les tables de l'authentification, décrites plus bas : `app_user`,
|
||||
puis `login_attempt` et `audit_log`, puis `refresh_token`. La cinquième, `e6d2026091501`, crée
|
||||
les six tables de données décrites en fin de document et déclare l'hypertable `reading`.
|
||||
|
||||
## Cycle de vie d'une mesure
|
||||
|
||||
Statut : `Cible`, sauf l'hypertable `reading` qui existe. Ni l'ingestion, ni les agrégats
|
||||
continus, ni la compression, ni la rétention ne sont écrits.
|
||||
Statut : `Partiellement fait`.
|
||||
|
||||
Les mécanismes d'ingestion sont maintenant implémentés pour les deux sources de données du MVP :
|
||||
|
||||
- le dataset historique CSV/JSON avec `historical_import.py` ;
|
||||
- l'API Mock avec `mock_api_import.py`.
|
||||
|
||||
Les traitements sont actuellement exécutables directement depuis le backend.
|
||||
|
||||
L'orchestration avec Apache Airflow reste une cible, tout comme les agrégats continus,
|
||||
la compression et les politiques de rétention.
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
src["Source de mesures"] -.-> ing["Ingestion Airflow"]
|
||||
ing -.-> hy[("Hypertable reading")]
|
||||
csv["CSV + JSON"] --> hist["historical_import.py"]
|
||||
mock["API Mock"] --> api["mock_api_import.py"]
|
||||
|
||||
hist --> hy[("Hypertable reading")]
|
||||
api --> hy
|
||||
|
||||
airflow["Airflow"] -.-> hist
|
||||
airflow -.-> api
|
||||
|
||||
hy -.-> agg[("Agrégat continu")]
|
||||
hy -.-> comp["Compression"]
|
||||
hy -.-> ret["Rétention"]
|
||||
agg -.-> api["API FastAPI"]
|
||||
|
||||
agg -.-> backend["API FastAPI"]
|
||||
agg -.-> graf["Grafana"]
|
||||
```
|
||||
|
||||
Les lectures de l'API et de Grafana visent 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.
|
||||
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.
|
||||
|
||||
## Tables d'authentification
|
||||
|
||||
Statut : `Fait`. Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les
|
||||
hypertables ; elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
||||
Statut : `Fait`.
|
||||
|
||||
Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les hypertables ;
|
||||
elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
||||
|
||||
```mermaid
|
||||
erDiagram
|
||||
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
||||
APP_USER ||--o{ PASSWORD_RESET_TOKEN : recoit
|
||||
|
||||
APP_USER {
|
||||
uuid id PK
|
||||
string email UK
|
||||
@@ -88,6 +121,7 @@ erDiagram
|
||||
bool must_change_password
|
||||
timestamptz credentials_changed_at
|
||||
}
|
||||
|
||||
REFRESH_TOKEN {
|
||||
uuid id PK
|
||||
uuid family_id
|
||||
@@ -99,6 +133,7 @@ erDiagram
|
||||
text revoked_reason
|
||||
uuid replaced_by
|
||||
}
|
||||
|
||||
LOGIN_ATTEMPT {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
@@ -106,6 +141,7 @@ erDiagram
|
||||
inet client_ip
|
||||
text outcome
|
||||
}
|
||||
|
||||
AUDIT_LOG {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
@@ -114,31 +150,47 @@ erDiagram
|
||||
text action
|
||||
jsonb detail
|
||||
}
|
||||
|
||||
PASSWORD_RESET_ATTEMPT {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
string email_tried
|
||||
inet client_ip
|
||||
}
|
||||
|
||||
PASSWORD_RESET_TOKEN {
|
||||
uuid id PK
|
||||
uuid user_id FK
|
||||
bytea token_hash UK
|
||||
timestamptz issued_at
|
||||
timestamptz expires_at
|
||||
timestamptz consumed_at
|
||||
inet client_ip
|
||||
text user_agent
|
||||
}
|
||||
```
|
||||
|
||||
Quatre choix de modélisation portent une intention et se défendent seuls :
|
||||
Plusieurs choix de modélisation portent une intention précise :
|
||||
|
||||
- **`app_user` et non `user`** : `user` est un mot réservé PostgreSQL, raccourci de
|
||||
`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).
|
||||
`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.
|
||||
|
||||
`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.
|
||||
`audit_log` porte des déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`.
|
||||
Elle n'est donc **pas** une hypertable.
|
||||
|
||||
## Gabarit de révision créant une hypertable
|
||||
|
||||
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision. La révision
|
||||
`e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision.
|
||||
|
||||
La révision `e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||
|
||||
```python
|
||||
def upgrade() -> None:
|
||||
@@ -182,24 +234,28 @@ colonne de temps : les index déclarés dans la révision le couvrent déjà.
|
||||
|
||||
## Questions ouvertes
|
||||
|
||||
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma est livré : ce qui suit
|
||||
porte sur son exploitation, plus sur sa forme.
|
||||
Elles portent maintenant principalement sur l'exploitation du schéma :
|
||||
|
||||
- **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure.
|
||||
- **Quels agrégats continus**, et sur quelles fenêtres.
|
||||
- **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse.
|
||||
- **Multi-tenant ou non** : un site appartient-il à un client, et faut-il cloisonner les lectures.
|
||||
- **Quelle granularité** conserver à long terme à l'ingestion : seconde, minute ou quart d'heure.
|
||||
- **Quels agrégats continus** créer et sur quelles fenêtres.
|
||||
- **Quelle profondeur de rétention** conserver en données brutes et à partir de quand compresser.
|
||||
- **Multi-tenant ou non** : un site appartient-il à un client et faut-il cloisonner les lectures.
|
||||
|
||||
## Modélisation détaillée des données
|
||||
|
||||
Cette modélisation prend en compte les fichiers CSV historiques,
|
||||
leurs métadonnées JSON et les données de l’API Mock.
|
||||
Elle comprend six tables, depuis le stockage des mesures
|
||||
jusqu’aux recommandations proposées à l’utilisateur.
|
||||
Cette modélisation prend en compte :
|
||||
|
||||
- les fichiers CSV historiques ;
|
||||
- leurs métadonnées JSON ;
|
||||
- les données de l'API Mock.
|
||||
|
||||
Elle comprend six tables Data, depuis le stockage des mesures jusqu'aux recommandations proposées
|
||||
à l'utilisateur.
|
||||
|
||||
### Schéma de données
|
||||
|
||||
Le diagramme ci-dessous présente les tables et leurs relations.
|
||||
|
||||
La révision `e6d2026091501` les crée.
|
||||
|
||||

|
||||
@@ -208,27 +264,20 @@ La révision `e6d2026091501` les crée.
|
||||
|
||||
### Description des tables
|
||||
|
||||
Chaque table remplit un rôle précis dans le traitement et l’exploitation
|
||||
des données.
|
||||
Chaque table remplit un rôle précis dans le traitement et l'exploitation des données.
|
||||
|
||||
| Table | Rôle | Origine des informations |
|
||||
|---|---|---|
|
||||
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l’import |
|
||||
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l'import |
|
||||
| `site` | Regrouper les informations des sites : identifiant, nom, type et caractéristiques disponibles | CSV et API Mock `/api/v1/sites` |
|
||||
| `reading` | Stocker les mesures, leur provenance, leur qualité et les éventuelles valeurs imputées | CSV et API Mock `/current` et `/readings` |
|
||||
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d’EnerVision |
|
||||
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d'EnerVision |
|
||||
| `alert` | Enregistrer les alertes, leur type, leur gravité et leur message | API Mock `/alerts` et détections EnerVision |
|
||||
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d’EnerVision |
|
||||
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d'EnerVision |
|
||||
|
||||
Les anomalies historiques décrites dans les JSON sont conservées
|
||||
dans `dataset.metadata`. Elles servent à l’analyse des données
|
||||
et ne sont pas considérées comme des alertes actuelles.
|
||||
Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`.
|
||||
|
||||
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.
|
||||
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
|
||||
|
||||
### Relations entre les tables
|
||||
|
||||
@@ -238,15 +287,21 @@ ligne.
|
||||
- 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
|
||||
|
||||
Le MVP EnerVision initialise les données énergétiques à partir du dataset fourni dans le cadre du projet.
|
||||
Statut : `Fait`.
|
||||
|
||||
Le dataset de référence contient 122 647 mesures issues de 7 sites et couvre la période du 1er janvier 2023 au 31 décembre 2024.
|
||||
Le MVP EnerVision initialise les données énergétiques à partir du dataset fourni dans le cadre
|
||||
du projet.
|
||||
|
||||
Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisation des données. Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
|
||||
Le dataset de référence contient 122 647 mesures issues de 7 sites et couvre la période
|
||||
du 1er janvier 2023 au 31 décembre 2024.
|
||||
|
||||
### Architecture du flux
|
||||
Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisation des données.
|
||||
|
||||
Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
|
||||
|
||||
## Architecture du flux historique
|
||||
|
||||
```text
|
||||
Dataset CSV + métadonnées JSON
|
||||
@@ -277,17 +332,26 @@ Dataset CSV + métadonnées JSON
|
||||
|
||||
Le pipeline est développé en Python.
|
||||
|
||||
Pandas est utilisé pour l'extraction, la validation et la préparation des données. SQLAlchemy Async assure le chargement transactionnel dans PostgreSQL/TimescaleDB.
|
||||
Pandas est utilisé pour l'extraction, la validation et la préparation des données.
|
||||
|
||||
SQLAlchemy Async assure le chargement transactionnel dans PostgreSQL/TimescaleDB.
|
||||
|
||||
Une empreinte SHA-256 permet d'identifier le dataset utilisé et d'assurer sa traçabilité.
|
||||
|
||||
Les valeurs manquantes sont conservées pendant l'ingestion afin de préserver les données sources. Aucune imputation n'est réalisée à cette étape.
|
||||
Les valeurs manquantes sont conservées pendant l'ingestion afin de préserver les données sources.
|
||||
|
||||
Aucune imputation n'est réalisée à cette étape.
|
||||
|
||||
Le chargement des mesures est effectué par batches de 1 000 lignes.
|
||||
|
||||
Les données provenant du dataset CSV sont identifiées par `source = "csv"` et associées à leur `dataset_id`.
|
||||
Les données provenant du dataset CSV sont identifiées par :
|
||||
|
||||
### Résultats validés
|
||||
```text
|
||||
source = "csv"
|
||||
dataset_id = identifiant du dataset
|
||||
```
|
||||
|
||||
## Résultats validés pour l'historique
|
||||
|
||||
Le chargement de référence a permis d'obtenir :
|
||||
|
||||
@@ -296,14 +360,207 @@ Le chargement de référence a permis d'obtenir :
|
||||
- 122 647 mesures ;
|
||||
- 0 doublon détecté dans le dataset source.
|
||||
|
||||
L'idempotence a également été vérifiée par une deuxième exécution du pipeline : aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à 122 647.
|
||||
L'idempotence a également été vérifiée par une deuxième exécution du pipeline :
|
||||
aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à 122 647.
|
||||
|
||||
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`.
|
||||
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`.
|
||||
|
||||
### Évolution prévue
|
||||
# Ingestion depuis l'API Mock
|
||||
|
||||
L'étape suivante consiste à orchestrer les traitements Data avec Apache Airflow.
|
||||
Statut : `Fait`.
|
||||
|
||||
L'orchestration réutilisera la logique ETL existante afin de séparer la logique de traitement de la planification, du suivi des exécutions et de la gestion des erreurs.
|
||||
La deuxième source du pipeline Data est l'API Mock EnerVision.
|
||||
|
||||
Le pipeline servira ensuite de base à la préparation des données nécessaires au modèle de Machine Learning.
|
||||
Le traitement est implémenté dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/mock_api_import.py
|
||||
```
|
||||
|
||||
## Endpoints utilisés
|
||||
|
||||
Le pipeline récupère les informations des sites depuis :
|
||||
|
||||
```text
|
||||
GET /api/v1/sites
|
||||
```
|
||||
|
||||
puis les mesures historiques simulées depuis :
|
||||
|
||||
```text
|
||||
GET /api/v1/readings
|
||||
```
|
||||
|
||||
Pour `/api/v1/readings`, les informations suivantes sont envoyées :
|
||||
|
||||
```text
|
||||
site_id
|
||||
start_time
|
||||
end_time
|
||||
limit
|
||||
```
|
||||
|
||||
Les paramètres de ligne de commande disponibles pour l'import sont :
|
||||
|
||||
```text
|
||||
--start-time
|
||||
--end-time
|
||||
--limit
|
||||
--dry-run
|
||||
```
|
||||
|
||||
## Flux d'ingestion API Mock
|
||||
|
||||
```text
|
||||
API Mock
|
||||
|
|
||||
+-----+------+
|
||||
| |
|
||||
v v
|
||||
/sites /readings
|
||||
| |
|
||||
+-----+------+
|
||||
|
|
||||
v
|
||||
mock_api_import.py
|
||||
|
|
||||
v
|
||||
Transformation
|
||||
+ qualité data
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
| |
|
||||
v v
|
||||
site reading
|
||||
```
|
||||
|
||||
Les informations des sites sont insérées ou mises à jour dans `site`.
|
||||
|
||||
Les mesures sont enregistrées dans l'hypertable `reading` avec :
|
||||
|
||||
```text
|
||||
source = "api_history"
|
||||
dataset_id = NULL
|
||||
```
|
||||
|
||||
Les données provenant de l'API Mock ne sont donc pas associées à un enregistrement de la table
|
||||
`dataset`.
|
||||
|
||||
La réponse source reçue depuis l'API est conservée dans :
|
||||
|
||||
```text
|
||||
raw_data
|
||||
```
|
||||
|
||||
## Qualité des données de l'API Mock
|
||||
|
||||
Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion.
|
||||
|
||||
Les informations suivantes fournies par l'API sont conservées :
|
||||
|
||||
```text
|
||||
data_quality
|
||||
null_reasons
|
||||
```
|
||||
|
||||
Cette conservation permet de distinguer une valeur manquante d'une valeur réelle égale à zéro
|
||||
et de garder les informations liées aux éventuelles défaillances de capteurs.
|
||||
|
||||
Aucune imputation n'est réalisée pendant cette phase :
|
||||
|
||||
```text
|
||||
imputed_values = NULL
|
||||
imputation_method = NULL
|
||||
```
|
||||
|
||||
## Validation de l'import API Mock
|
||||
|
||||
Un scénario de validation a été exécuté pour les 7 sites sur la période :
|
||||
|
||||
```text
|
||||
15/06/2024 12:00 UTC
|
||||
à
|
||||
15/06/2024 13:00 UTC
|
||||
```
|
||||
|
||||
avec :
|
||||
|
||||
```text
|
||||
limit = 60
|
||||
```
|
||||
|
||||
Résultat :
|
||||
|
||||
```text
|
||||
7 sites
|
||||
60 lectures par site
|
||||
420 lectures récupérées
|
||||
```
|
||||
|
||||
Les données ont été chargées dans PostgreSQL/TimescaleDB puis contrôlées directement en base.
|
||||
|
||||
Les contrôles ont confirmé :
|
||||
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` ;
|
||||
- la conservation de `null_reasons` ;
|
||||
- la conservation de `raw_data`.
|
||||
|
||||
L'idempotence a été vérifiée en rejouant le même import.
|
||||
|
||||
Une mesure déjà présente n'est pas ajoutée une seconde fois.
|
||||
|
||||
Les tests automatisés couvrent également :
|
||||
|
||||
- la récupération des sites ;
|
||||
- les paramètres envoyés à `/api/v1/readings` ;
|
||||
- les réponses HTTP en erreur ;
|
||||
- le format de la réponse ;
|
||||
- la transformation des mesures ;
|
||||
- les valeurs manquantes ;
|
||||
- la qualité des données ;
|
||||
- la conservation des données sources ;
|
||||
- l'idempotence en base.
|
||||
|
||||
# Évolution prévue
|
||||
|
||||
La prochaine étape consiste à orchestrer les deux mécanismes d'ingestion avec Apache Airflow.
|
||||
|
||||
```text
|
||||
CSV / JSON ----------------+
|
||||
|
|
||||
v
|
||||
+------------------+
|
||||
| Airflow |
|
||||
+------------------+
|
||||
|
|
||||
+----------------+----------------+
|
||||
| |
|
||||
v v
|
||||
historical_import.py mock_api_import.py
|
||||
| |
|
||||
+----------------+----------------+
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
```
|
||||
|
||||
Airflow servira à :
|
||||
|
||||
- planifier les traitements ;
|
||||
- définir leur ordre d'exécution ;
|
||||
- suivre leur état ;
|
||||
- gérer et remonter les erreurs ;
|
||||
- faciliter les exécutions récurrentes.
|
||||
|
||||
Airflow ne remplacera pas la logique ETL déjà implémentée.
|
||||
|
||||
Les scripts Python resteront responsables de l'extraction, de la validation, de la transformation
|
||||
et du chargement des données.
|
||||
|
||||
Le pipeline servira ensuite de base à la préparation des données nécessaires au modèle
|
||||
de Machine Learning.
|
||||
+307
-22
@@ -2,9 +2,10 @@
|
||||
|
||||
## Objectif
|
||||
|
||||
Le pipeline ETL EnerVision permet d'intégrer les données énergétiques historiques dans PostgreSQL/TimescaleDB.
|
||||
Le pipeline ETL EnerVision permet d'intégrer les données énergétiques dans PostgreSQL/TimescaleDB à partir de deux sources :
|
||||
|
||||
Cette première étape du pipeline Data permet de charger le dataset fourni dans le cadre du projet, contenant les mesures énergétiques de 7 sites sur la période du 1er janvier 2023 au 31 décembre 2024.
|
||||
- le dataset historique CSV/JSON fourni dans le cadre du projet ;
|
||||
- l'API Mock EnerVision.
|
||||
|
||||
Le pipeline assure :
|
||||
|
||||
@@ -12,12 +13,15 @@ Le pipeline assure :
|
||||
- la validation de leur structure et de leur cohérence ;
|
||||
- la normalisation des données nécessaires au stockage ;
|
||||
- le suivi de la qualité des données ;
|
||||
- la traçabilité du dataset importé ;
|
||||
- la traçabilité des données importées ;
|
||||
- le chargement des données dans PostgreSQL/TimescaleDB ;
|
||||
- la conservation des valeurs manquantes et des informations de qualité ;
|
||||
- l'idempotence du chargement afin d'éviter la création de doublons.
|
||||
|
||||
## Données sources
|
||||
|
||||
### Dataset historique
|
||||
|
||||
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
||||
|
||||
Il contient les deux fichiers suivants :
|
||||
@@ -29,7 +33,7 @@ dataset_metadata.json
|
||||
|
||||
Ces fichiers sont nécessaires une seule fois pour initialiser les données historiques de l'environnement.
|
||||
|
||||
Ils ne sont pas versionnés dans Git. Chaque membre de l'équipe récupère manuellement une fois les fichiers fournis par le formateur et les place dans :
|
||||
Ils ne sont pas versionnés dans Git. Chaque membre de l'équipe récupère manuellement les fichiers fournis par le formateur et les place dans :
|
||||
|
||||
```text
|
||||
data/raw/
|
||||
@@ -47,14 +51,26 @@ data/
|
||||
|
||||
Le fichier `.gitkeep` est versionné afin de conserver le répertoire `data/raw/` dans Git. Les fichiers CSV et JSON sont ignorés par Git.
|
||||
|
||||
### API Mock
|
||||
|
||||
La deuxième source est l'API Mock EnerVision.
|
||||
|
||||
Elle permet de récupérer :
|
||||
|
||||
- les informations des sites avec `GET /api/v1/sites` ;
|
||||
- les mesures simulées avec `GET /api/v1/readings`.
|
||||
|
||||
L'API Mock est utilisée pour compléter les données historiques avec des mesures simulées récupérées sur une période donnée.
|
||||
|
||||
## Technologies utilisées
|
||||
|
||||
| Technologie | Utilisation |
|
||||
|---|---|
|
||||
| Python | Développement du pipeline ETL |
|
||||
| Pandas | Lecture, validation et transformation des données |
|
||||
| JSON | Lecture des métadonnées du dataset |
|
||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset |
|
||||
| Pandas | Lecture, validation et transformation du dataset historique |
|
||||
| JSON | Lecture des métadonnées et conservation des données sources |
|
||||
| HTTPX | Appels HTTP asynchrones vers l'API Mock |
|
||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset historique |
|
||||
| SQLAlchemy Async | Connexion et chargement asynchrone en base |
|
||||
| PostgreSQL | Stockage relationnel |
|
||||
| TimescaleDB | Stockage des séries temporelles énergétiques |
|
||||
@@ -62,11 +78,14 @@ Le fichier `.gitkeep` est versionné afin de conserver le répertoire `data/raw/
|
||||
| Alembic | Gestion des migrations du schéma |
|
||||
| uv | Gestion et exécution de l'environnement Python |
|
||||
| Ruff | Contrôle de la qualité du code |
|
||||
| mypy | Vérification du typage |
|
||||
| Pytest | Tests automatisés |
|
||||
|
||||
## Fonctionnement du pipeline
|
||||
# Import du dataset historique
|
||||
|
||||
Le script principal d'import se trouve dans :
|
||||
## Fonctionnement du pipeline historique
|
||||
|
||||
Le script d'import se trouve dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/historical_import.py
|
||||
@@ -207,7 +226,7 @@ Valeurs manquantes identifiées :
|
||||
| `humidity_percent` | 3 423 |
|
||||
| `solar_irradiance_wm2` | 3 964 |
|
||||
|
||||
## Exécution en dry-run
|
||||
## Exécution historique en dry-run
|
||||
|
||||
Depuis le dossier :
|
||||
|
||||
@@ -227,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 réel
|
||||
## Chargement historique réel
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
@@ -249,7 +268,7 @@ Chargement : 2000/122647
|
||||
Chargement : 122647/122647
|
||||
```
|
||||
|
||||
## Résultats obtenus
|
||||
## Résultats obtenus pour le dataset historique
|
||||
|
||||
Après le chargement initial, les contrôles en base ont confirmé :
|
||||
|
||||
@@ -266,7 +285,7 @@ Le premier import a créé :
|
||||
nouvelles lectures : 122647
|
||||
```
|
||||
|
||||
## Idempotence
|
||||
## 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.
|
||||
|
||||
@@ -280,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
|
||||
## Vérifications SQL du dataset historique
|
||||
|
||||
Depuis la racine du projet, vérifier le nombre d'enregistrements avec :
|
||||
|
||||
@@ -302,21 +321,213 @@ Vérifier la source des mesures avec :
|
||||
docker compose exec db psql -U enervision -d enervision -c "SELECT source, COUNT(*) FROM reading GROUP BY source ORDER BY source;"
|
||||
```
|
||||
|
||||
Résultat attendu :
|
||||
Résultat attendu pour le dataset historique :
|
||||
|
||||
```text
|
||||
csv | 122647
|
||||
```
|
||||
|
||||
## Tests et qualité
|
||||
# Import depuis l'API Mock
|
||||
|
||||
Les tests automatisés du pipeline sont situés dans :
|
||||
## Fonctionnement
|
||||
|
||||
Le script d'import de l'API Mock se trouve dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/mock_api_import.py
|
||||
```
|
||||
|
||||
Le flux est le suivant :
|
||||
|
||||
```text
|
||||
API Mock
|
||||
|
|
||||
+-----+------+
|
||||
| |
|
||||
v v
|
||||
/sites /readings
|
||||
| |
|
||||
+-----+------+
|
||||
|
|
||||
v
|
||||
mock_api_import.py
|
||||
|
|
||||
v
|
||||
Transformation
|
||||
+ qualité data
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
| |
|
||||
v v
|
||||
site reading
|
||||
```
|
||||
|
||||
Le pipeline commence par récupérer les sites avec :
|
||||
|
||||
```text
|
||||
GET /api/v1/sites
|
||||
```
|
||||
|
||||
Il récupère ensuite les mesures de chaque site avec :
|
||||
|
||||
```text
|
||||
GET /api/v1/readings
|
||||
```
|
||||
|
||||
Les paramètres envoyés à `/api/v1/readings` sont :
|
||||
|
||||
```text
|
||||
site_id
|
||||
start_time
|
||||
end_time
|
||||
limit
|
||||
```
|
||||
|
||||
Le paramètre `limit` doit être compris entre 1 et 1000.
|
||||
|
||||
## Configuration de l'API Mock
|
||||
|
||||
La connexion à l'API Mock est configurée avec les variables d'environnement suivantes :
|
||||
|
||||
```text
|
||||
APP_MOCK_API_BASE_URL
|
||||
APP_MOCK_API_USERNAME
|
||||
APP_MOCK_API_PASSWORD
|
||||
APP_MOCK_API_TIMEOUT_SECONDS
|
||||
```
|
||||
|
||||
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
|
||||
|
||||
Les mesures provenant de l'API Mock sont enregistrées dans `reading` avec :
|
||||
|
||||
```text
|
||||
source = "api_history"
|
||||
dataset_id = NULL
|
||||
```
|
||||
|
||||
Les mesures provenant de l'API ne sont donc pas rattachées à un dataset historique.
|
||||
|
||||
Le timestamp reçu depuis l'API est converti en `datetime` avec timezone avant le chargement.
|
||||
|
||||
La réponse source est conservée dans :
|
||||
|
||||
```text
|
||||
raw_data
|
||||
```
|
||||
|
||||
afin de préserver la donnée reçue et faciliter la traçabilité.
|
||||
|
||||
## Qualité des données API
|
||||
|
||||
Les valeurs `NULL` fournies par l'API sont conservées telles quelles.
|
||||
|
||||
Une valeur manquante n'est pas transformée en zéro et la mesure n'est pas supprimée.
|
||||
|
||||
Le pipeline conserve également :
|
||||
|
||||
```text
|
||||
data_quality
|
||||
null_reasons
|
||||
```
|
||||
|
||||
Les niveaux de qualité possibles sont :
|
||||
|
||||
```text
|
||||
good
|
||||
partial
|
||||
degraded
|
||||
critical
|
||||
```
|
||||
|
||||
Aucune imputation n'est réalisée pendant l'ingestion :
|
||||
|
||||
```text
|
||||
imputed_values = NULL
|
||||
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
|
||||
|
||||
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.
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
```powershell
|
||||
uv run python -m app.etl.mock_api_import `
|
||||
--start-time "2024-06-15T12:00:00" `
|
||||
--end-time "2024-06-15T13:00:00" `
|
||||
--limit 60 `
|
||||
--dry-run
|
||||
```
|
||||
|
||||
## Chargement réel depuis l'API Mock
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
```powershell
|
||||
uv run python -m app.etl.mock_api_import `
|
||||
--start-time "2024-06-15T12:00:00" `
|
||||
--end-time "2024-06-15T13:00:00" `
|
||||
--limit 60
|
||||
```
|
||||
|
||||
## Résultat validé pour l'API Mock
|
||||
|
||||
Le scénario de validation utilisé couvre la période :
|
||||
|
||||
```text
|
||||
15/06/2024 12:00 UTC
|
||||
à
|
||||
15/06/2024 13:00 UTC
|
||||
```
|
||||
|
||||
avec une limite de 60 lectures par site.
|
||||
|
||||
Résultat obtenu :
|
||||
|
||||
```text
|
||||
sites récupérés : 7
|
||||
lectures par site : 60
|
||||
lectures récupérées : 420
|
||||
source : api_history
|
||||
dataset_id : NULL
|
||||
```
|
||||
|
||||
Les contrôles effectués directement dans PostgreSQL/TimescaleDB ont confirmé :
|
||||
|
||||
- l'enregistrement des mesures dans `reading` ;
|
||||
- la présence des 7 sites ;
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` ;
|
||||
- la conservation de `null_reasons` ;
|
||||
- la conservation de la donnée source dans `raw_data`.
|
||||
|
||||
## 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.
|
||||
|
||||
L'idempotence repose sur la contrainte d'unicité de la table `reading` et sur la gestion des conflits lors de l'insertion.
|
||||
|
||||
Un test d'intégration automatisé vérifie également ce comportement.
|
||||
|
||||
# Tests et qualité
|
||||
|
||||
Les tests automatisés des pipelines ETL sont situés dans :
|
||||
|
||||
```text
|
||||
apps/backend/tests/etl/
|
||||
```
|
||||
|
||||
Ils couvrent notamment :
|
||||
Les tests de l'import historique couvrent notamment :
|
||||
|
||||
- la validation du dataset ;
|
||||
- les colonnes obligatoires ;
|
||||
@@ -328,22 +539,96 @@ Ils couvrent notamment :
|
||||
- la construction des mesures destinées à la BDD ;
|
||||
- le respect des contraintes du modèle de données.
|
||||
|
||||
Les tests de l'import API Mock couvrent notamment :
|
||||
|
||||
- la récupération des sites ;
|
||||
- l'appel à `/api/v1/readings` ;
|
||||
- les paramètres `site_id`, `start_time`, `end_time` et `limit` ;
|
||||
- la gestion des erreurs HTTP ;
|
||||
- la validation du format de la réponse ;
|
||||
- la transformation des mesures ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` et `null_reasons` ;
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation de `raw_data` ;
|
||||
- l'idempotence du chargement.
|
||||
|
||||
Exécuter les tests ETL :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl -v
|
||||
```
|
||||
|
||||
Exécuter les tests unitaires de l'import API Mock :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl\test_mock_api_import.py -v
|
||||
```
|
||||
|
||||
Exécuter le test d'intégration de l'import API Mock :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl\test_mock_api_import.py -m integration -v
|
||||
```
|
||||
|
||||
Contrôler la qualité du code :
|
||||
|
||||
```powershell
|
||||
uv run ruff check app\etl tests\etl
|
||||
```
|
||||
|
||||
## Suite du pipeline Data
|
||||
Contrôler le typage :
|
||||
|
||||
L'import historique constitue la première brique du pipeline Data EnerVision.
|
||||
```powershell
|
||||
uv run mypy app
|
||||
```
|
||||
|
||||
La prochaine étape consiste à orchestrer les traitements ETL avec Apache Airflow, puis à préparer les données nécessaires à l'entraînement du modèle de Machine Learning.
|
||||
Exécuter la suite complète avec le seuil de couverture :
|
||||
|
||||
Airflow sera utilisé comme orchestrateur des traitements existants et ne remplacera pas la logique métier déjà implémentée dans le pipeline ETL.
|
||||
```powershell
|
||||
uv run pytest --cov-fail-under=85
|
||||
```
|
||||
|
||||
Lors de la validation de l'import API Mock :
|
||||
|
||||
```text
|
||||
8 tests unitaires passés
|
||||
1 test d'intégration passé
|
||||
```
|
||||
|
||||
La suite backend complète a également été validée avec une couverture supérieure au seuil de 85 %.
|
||||
|
||||
# Suite du pipeline Data
|
||||
|
||||
Deux sources de données sont maintenant prises en charge :
|
||||
|
||||
```text
|
||||
Dataset CSV/JSON
|
||||
|
|
||||
v
|
||||
historical_import.py
|
||||
|
|
||||
+-----------------+
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
^
|
||||
|
|
||||
+-----------------+
|
||||
|
|
||||
mock_api_import.py
|
||||
^
|
||||
|
|
||||
API Mock
|
||||
```
|
||||
|
||||
La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP.
|
||||
|
||||
La prochaine étape consiste à orchestrer ces traitements avec Apache Airflow.
|
||||
|
||||
Airflow permettra de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs.
|
||||
|
||||
Airflow ne remplacera pas la logique ETL Python existante. Les scripts actuels resteront responsables de l'extraction, de la validation, de la transformation et du chargement.
|
||||
|
||||
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
|
||||
Reference in New Issue
Block a user