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_SECRET_KEY=change_me
|
||||||
APP_CORS_ORIGINS=http://localhost:4200
|
APP_CORS_ORIGINS=http://localhost:4200
|
||||||
BACKEND_PORT=8000
|
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 \
|
.PHONY: help install install-backend install-frontend install-ml dev dev-backend dev-frontend \
|
||||||
lint format typecheck test test-cov test-integration check \
|
lint format typecheck test test-cov test-integration check \
|
||||||
openapi docker-build db-up db-down db-reset db-logs db-psql migrate bootstrap-admin \
|
openapi docker-build db-up db-down db-reset db-logs db-psql migrate bootstrap-admin \
|
||||||
ml-lint ml-typecheck ml-test ml-check ml-train ml-score recommendations
|
ml-lint ml-typecheck ml-test ml-check ml-train ml-score
|
||||||
|
|
||||||
help: ## Liste les cibles disponibles
|
help: ## Liste les cibles disponibles
|
||||||
@grep -E '^[a-zA-Z_-]+:.*?## .*$$' $(MAKEFILE_LIST) | awk 'BEGIN {FS = ":.*?## "}; {printf " \033[36m%-16s\033[0m %s\n", $$1, $$2}'
|
@grep -E '^[a-zA-Z_-]+:.*?## .*$$' $(MAKEFILE_LIST) | awk 'BEGIN {FS = ":.*?## "}; {printf " \033[36m%-16s\033[0m %s\n", $$1, $$2}'
|
||||||
@@ -77,9 +77,6 @@ ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DAT
|
|||||||
ml-score: ## Score le prochain pas horaire et l'ecrit dans `prediction`. CSV=chemin optionnel
|
ml-score: ## Score le prochain pas horaire et l'ecrit dans `prediction`. CSV=chemin optionnel
|
||||||
cd $(ML) && uv run python -m enervision_ml.score $(if $(CSV),--csv $(CSV),)
|
cd $(ML) && uv run python -m enervision_ml.score $(if $(CSV),--csv $(CSV),)
|
||||||
|
|
||||||
recommendations: ## Genere les recommandations depuis les alertes en base. SITE=identifiant optionnel
|
|
||||||
cd $(BACKEND) && uv run python -m app.cli generate-recommendations $(if $(SITE),--site-id $(SITE),)
|
|
||||||
|
|
||||||
docker-build: ## Construit l'image du backend
|
docker-build: ## Construit l'image du backend
|
||||||
docker build -t enervision-backend:local $(BACKEND)
|
docker build -t enervision-backend:local $(BACKEND)
|
||||||
|
|
||||||
|
|||||||
@@ -18,3 +18,7 @@ APP_SMTP_HOST=localhost
|
|||||||
APP_SMTP_PORT=1025
|
APP_SMTP_PORT=1025
|
||||||
APP_SMTP_USE_TLS=false
|
APP_SMTP_USE_TLS=false
|
||||||
APP_SMTP_FROM_ADDRESS=no-reply@enervision.fr
|
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/sites/{site_id}` | Décrit un site | `lecteur` |
|
||||||
| `/api/v1/recommendations` | Liste les recommandations | `lecteur` |
|
| `/api/v1/recommendations` | Liste les recommandations | `lecteur` |
|
||||||
| `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation | `lecteur` |
|
| `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation | `lecteur` |
|
||||||
| `/api/v1/recommendations/generate` | Génère les recommandations depuis les alertes (POST) | `admin` |
|
|
||||||
| `/metrics` | Métriques au format Prometheus | jeton si `APP_METRICS_TOKEN` |
|
| `/metrics` | Métriques au format Prometheus | jeton si `APP_METRICS_TOKEN` |
|
||||||
| `/docs`, `/openapi.json` | Documentation, fermée en `staging` et `prod` | public sinon |
|
| `/docs`, `/openapi.json` | Documentation, fermée en `staging` et `prod` | public sinon |
|
||||||
|
|
||||||
|
|||||||
@@ -192,11 +192,7 @@ AlertServiceDep = Annotated[AlertService, Depends(get_alert_service)]
|
|||||||
|
|
||||||
|
|
||||||
def get_recommendation_service(session: SessionDep) -> RecommendationService:
|
def get_recommendation_service(session: SessionDep) -> RecommendationService:
|
||||||
return RecommendationService(
|
return RecommendationService(recommendations=RecommendationRepository(session))
|
||||||
recommendations=RecommendationRepository(session),
|
|
||||||
alerts=AlertRepository(session),
|
|
||||||
transaction=session,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
RecommendationServiceDep = Annotated[RecommendationService, Depends(get_recommendation_service)]
|
RecommendationServiceDep = Annotated[RecommendationService, Depends(get_recommendation_service)]
|
||||||
|
|||||||
@@ -63,7 +63,7 @@ TAGS: Final[list[dict[str, Any]]] = [
|
|||||||
"name": "recommendations",
|
"name": "recommendations",
|
||||||
"description": (
|
"description": (
|
||||||
"Consultation des recommandations issues des alertes. Accessible à partir du rôle "
|
"Consultation des recommandations issues des alertes. Accessible à partir du rôle "
|
||||||
"`lecteur`. Leur génération par le moteur de règles est réservée au rôle `admin`."
|
"`lecteur`."
|
||||||
),
|
),
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -1,18 +1,13 @@
|
|||||||
from fastapi import APIRouter, HTTPException, status
|
from fastapi import APIRouter, HTTPException, status
|
||||||
|
|
||||||
from app.api.deps import AdminDep, LecteurDep, RecommendationServiceDep
|
from app.api.deps import LecteurDep, RecommendationServiceDep
|
||||||
from app.api.openapi import REPONSE_VALIDATION, REPONSES_ADMIN, Reponses
|
from app.api.openapi import REPONSE_VALIDATION, Reponses
|
||||||
from app.schemas.errors import ErrorResponse
|
from app.schemas.errors import ErrorResponse
|
||||||
from app.schemas.recommendation import (
|
from app.schemas.recommendation import RecommendationResponse
|
||||||
RecommendationGenerationResponse,
|
|
||||||
RecommendationResponse,
|
|
||||||
)
|
|
||||||
from app.services.recommendation import RecommendationNotFoundError
|
from app.services.recommendation import RecommendationNotFoundError
|
||||||
|
|
||||||
router = APIRouter()
|
router = APIRouter()
|
||||||
|
|
||||||
REPONSES_GENERATION: Reponses = {**REPONSES_ADMIN, **REPONSE_VALIDATION}
|
|
||||||
|
|
||||||
REPONSES_INTROUVABLE: Reponses = {
|
REPONSES_INTROUVABLE: Reponses = {
|
||||||
**REPONSE_VALIDATION,
|
**REPONSE_VALIDATION,
|
||||||
404: {"model": ErrorResponse, "description": "Aucune recommandation ne porte cet identifiant."},
|
404: {"model": ErrorResponse, "description": "Aucune recommandation ne porte cet identifiant."},
|
||||||
@@ -43,22 +38,3 @@ async def get_recommendation(
|
|||||||
status_code=status.HTTP_404_NOT_FOUND, detail="Recommandation introuvable"
|
status_code=status.HTTP_404_NOT_FOUND, detail="Recommandation introuvable"
|
||||||
) from erreur
|
) from erreur
|
||||||
return RecommendationResponse.model_validate(recommendation)
|
return RecommendationResponse.model_validate(recommendation)
|
||||||
|
|
||||||
|
|
||||||
@router.post(
|
|
||||||
"/generate",
|
|
||||||
response_model=RecommendationGenerationResponse,
|
|
||||||
summary="Génère les recommandations à partir des alertes",
|
|
||||||
responses=REPONSES_GENERATION,
|
|
||||||
)
|
|
||||||
async def generate_recommendations(
|
|
||||||
_: AdminDep,
|
|
||||||
service: RecommendationServiceDep,
|
|
||||||
site_id: str | None = None,
|
|
||||||
) -> RecommendationGenerationResponse:
|
|
||||||
rapport = await service.generate(site_id=site_id)
|
|
||||||
return RecommendationGenerationResponse(
|
|
||||||
alerts_examined=rapport.alertes_examinees,
|
|
||||||
recommendations_created=rapport.recommandations_creees,
|
|
||||||
already_present=rapport.deja_presentes,
|
|
||||||
)
|
|
||||||
|
|||||||
@@ -22,11 +22,8 @@ from app.core.hashing import build_hasher
|
|||||||
from app.core.roles import Role
|
from app.core.roles import Role
|
||||||
from app.db.session import get_session_factory
|
from app.db.session import get_session_factory
|
||||||
from app.main import create_app
|
from app.main import create_app
|
||||||
from app.repositories.alert import AlertRepository
|
|
||||||
from app.repositories.recommendation import RecommendationRepository
|
|
||||||
from app.repositories.user import UserRepository
|
from app.repositories.user import UserRepository
|
||||||
from app.schemas.auth import PASSWORD_MIN_LENGTH, SPECIAL_CHARACTERS, valide_complexite
|
from app.schemas.auth import PASSWORD_MIN_LENGTH, SPECIAL_CHARACTERS, valide_complexite
|
||||||
from app.services.recommendation import RecommendationService
|
|
||||||
|
|
||||||
LONGUEUR_MOT_DE_PASSE_GENERE = 24
|
LONGUEUR_MOT_DE_PASSE_GENERE = 24
|
||||||
CHEMIN_CONTRAT = Path(__file__).resolve().parent.parent / "openapi.json"
|
CHEMIN_CONTRAT = Path(__file__).resolve().parent.parent / "openapi.json"
|
||||||
@@ -66,22 +63,6 @@ async def create_admin(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
async def generate_recommendations(*, site_id: str | None) -> str:
|
|
||||||
async with get_session_factory()() as session:
|
|
||||||
service = RecommendationService(
|
|
||||||
recommendations=RecommendationRepository(session),
|
|
||||||
alerts=AlertRepository(session),
|
|
||||||
transaction=session,
|
|
||||||
)
|
|
||||||
rapport = await service.generate(site_id=site_id)
|
|
||||||
|
|
||||||
return (
|
|
||||||
f"{rapport.alertes_examinees} alerte(s) examinée(s), "
|
|
||||||
f"{rapport.recommandations_creees} recommandation(s) créée(s), "
|
|
||||||
f"{rapport.deja_presentes} déjà présente(s)"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
# Piège : le schéma ne doit dépendre ni du `.env` du poste ni des variables `APP_*`, sinon le
|
# Piège : le schéma ne doit dépendre ni du `.env` du poste ni des variables `APP_*`, sinon le
|
||||||
# fichier versionné changerait de machine en machine et le test de dérive deviendrait un oracle
|
# fichier versionné changerait de machine en machine et le test de dérive deviendrait un oracle
|
||||||
# de configuration locale. Tout ce qui atteint le schéma est donc posé ici, `_env_file` compris.
|
# de configuration locale. Tout ce qui atteint le schéma est donc posé ici, `_env_file` compris.
|
||||||
@@ -128,14 +109,6 @@ def build_parser() -> argparse.ArgumentParser:
|
|||||||
"export-openapi", help="Écrit le contrat OpenAPI sur disque"
|
"export-openapi", help="Écrit le contrat OpenAPI sur disque"
|
||||||
)
|
)
|
||||||
contrat.add_argument("--output", default=str(CHEMIN_CONTRAT))
|
contrat.add_argument("--output", default=str(CHEMIN_CONTRAT))
|
||||||
|
|
||||||
recommandations = sous_commandes.add_parser(
|
|
||||||
"generate-recommendations",
|
|
||||||
help="Applique le moteur de règles aux alertes en base",
|
|
||||||
)
|
|
||||||
recommandations.add_argument(
|
|
||||||
"--site-id", default=None, help="Limite le traitement aux alertes d'un site"
|
|
||||||
)
|
|
||||||
return parser
|
return parser
|
||||||
|
|
||||||
|
|
||||||
@@ -179,10 +152,6 @@ def main(argv: list[str] | None = None) -> int:
|
|||||||
print(export_openapi(Path(arguments.output)))
|
print(export_openapi(Path(arguments.output)))
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
if arguments.commande == "generate-recommendations":
|
|
||||||
print(asyncio.run(generate_recommendations(site_id=arguments.site_id)))
|
|
||||||
return 0
|
|
||||||
|
|
||||||
mot_de_passe = read_password(generate=arguments.generate)
|
mot_de_passe = read_password(generate=arguments.generate)
|
||||||
|
|
||||||
succes, message = asyncio.run(
|
succes, message = asyncio.run(
|
||||||
|
|||||||
@@ -34,6 +34,11 @@ class Settings(BaseSettings):
|
|||||||
database_pool_size: int = 5
|
database_pool_size: int = 5
|
||||||
database_max_overflow: int = 10
|
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_issuer: str = "enervision-api"
|
||||||
jwt_audience: str = "enervision-web"
|
jwt_audience: str = "enervision-web"
|
||||||
access_token_ttl_seconds: int = Field(default=900, ge=60, le=3600)
|
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 collections.abc import Sequence
|
||||||
from dataclasses import asdict, dataclass
|
|
||||||
|
|
||||||
from sqlalchemy import select
|
from sqlalchemy import select
|
||||||
from sqlalchemy.dialects.postgresql import insert
|
|
||||||
from sqlalchemy.ext.asyncio import AsyncSession
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
from app.models.energy import Recommendation
|
from app.models.energy import Recommendation
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True, slots=True)
|
|
||||||
class NouvelleRecommandation:
|
|
||||||
alert_id: int
|
|
||||||
action: str
|
|
||||||
explanation: str
|
|
||||||
rule_reference: str
|
|
||||||
|
|
||||||
|
|
||||||
TAILLE_DE_LOT = 1000
|
|
||||||
|
|
||||||
|
|
||||||
class RecommendationRepository:
|
class RecommendationRepository:
|
||||||
def __init__(self, session: AsyncSession) -> None:
|
def __init__(self, session: AsyncSession) -> None:
|
||||||
self._session = session
|
self._session = session
|
||||||
@@ -33,19 +20,3 @@ class RecommendationRepository:
|
|||||||
)
|
)
|
||||||
recommendation: Recommendation | None = await self._session.scalar(requete)
|
recommendation: Recommendation | None = await self._session.scalar(requete)
|
||||||
return recommendation
|
return recommendation
|
||||||
|
|
||||||
# Pourquoi : l'idempotence est déléguée à `uq_recommendation_alert_rule` plutôt qu'à une
|
|
||||||
# lecture préalable, qui laisserait une fenêtre entre le contrôle et l'insertion.
|
|
||||||
async def create_missing(self, nouvelles: Sequence[NouvelleRecommandation]) -> int:
|
|
||||||
creees = 0
|
|
||||||
# Piège : asyncpg plafonne une requête à 32 767 paramètres, soit 8 191 lignes de quatre
|
|
||||||
# colonnes. Au-delà de ce seuil un `INSERT` d'un seul tenant échouerait.
|
|
||||||
for debut in range(0, len(nouvelles), TAILLE_DE_LOT):
|
|
||||||
requete = (
|
|
||||||
insert(Recommendation)
|
|
||||||
.values([asdict(nouvelle) for nouvelle in nouvelles[debut : debut + TAILLE_DE_LOT]])
|
|
||||||
.on_conflict_do_nothing(constraint="uq_recommendation_alert_rule")
|
|
||||||
.returning(Recommendation.recommendation_id)
|
|
||||||
)
|
|
||||||
creees += len((await self._session.scalars(requete)).all())
|
|
||||||
return creees
|
|
||||||
|
|||||||
@@ -12,9 +12,3 @@ class RecommendationResponse(BaseModel):
|
|||||||
explanation: str
|
explanation: str
|
||||||
rule_reference: str
|
rule_reference: str
|
||||||
created_at: datetime
|
created_at: datetime
|
||||||
|
|
||||||
|
|
||||||
class RecommendationGenerationResponse(BaseModel):
|
|
||||||
alerts_examined: int
|
|
||||||
recommendations_created: int
|
|
||||||
already_present: int
|
|
||||||
|
|||||||
@@ -1,15 +1,7 @@
|
|||||||
from collections.abc import Sequence
|
from collections.abc import Sequence
|
||||||
from dataclasses import dataclass
|
|
||||||
from typing import Protocol
|
|
||||||
|
|
||||||
from app.models.energy import Recommendation
|
from app.models.energy import Recommendation
|
||||||
from app.repositories.alert import AlertRepository
|
|
||||||
from app.repositories.recommendation import RecommendationRepository
|
from app.repositories.recommendation import RecommendationRepository
|
||||||
from app.services.recommendation_rules import applique_les_regles
|
|
||||||
|
|
||||||
|
|
||||||
class Transaction(Protocol):
|
|
||||||
async def commit(self) -> None: ...
|
|
||||||
|
|
||||||
|
|
||||||
class RecommendationError(Exception):
|
class RecommendationError(Exception):
|
||||||
@@ -20,24 +12,9 @@ class RecommendationNotFoundError(RecommendationError):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True, slots=True)
|
|
||||||
class RapportGeneration:
|
|
||||||
alertes_examinees: int
|
|
||||||
recommandations_creees: int
|
|
||||||
deja_presentes: int
|
|
||||||
|
|
||||||
|
|
||||||
class RecommendationService:
|
class RecommendationService:
|
||||||
def __init__(
|
def __init__(self, *, recommendations: RecommendationRepository) -> None:
|
||||||
self,
|
|
||||||
*,
|
|
||||||
recommendations: RecommendationRepository,
|
|
||||||
alerts: AlertRepository,
|
|
||||||
transaction: Transaction,
|
|
||||||
) -> None:
|
|
||||||
self._recommendations = recommendations
|
self._recommendations = recommendations
|
||||||
self._alerts = alerts
|
|
||||||
self._transaction = transaction
|
|
||||||
|
|
||||||
async def list_all(self) -> Sequence[Recommendation]:
|
async def list_all(self) -> Sequence[Recommendation]:
|
||||||
return await self._recommendations.list_all()
|
return await self._recommendations.list_all()
|
||||||
@@ -47,16 +24,3 @@ class RecommendationService:
|
|||||||
if recommendation is None:
|
if recommendation is None:
|
||||||
raise RecommendationNotFoundError(recommendation_id)
|
raise RecommendationNotFoundError(recommendation_id)
|
||||||
return recommendation
|
return recommendation
|
||||||
|
|
||||||
async def generate(self, *, site_id: str | None = None) -> RapportGeneration:
|
|
||||||
alertes = await self._alerts.list_all(site_id=site_id)
|
|
||||||
nouvelles = [nouvelle for alerte in alertes for nouvelle in applique_les_regles(alerte)]
|
|
||||||
|
|
||||||
creees = await self._recommendations.create_missing(nouvelles)
|
|
||||||
await self._transaction.commit()
|
|
||||||
|
|
||||||
return RapportGeneration(
|
|
||||||
alertes_examinees=len(alertes),
|
|
||||||
recommandations_creees=creees,
|
|
||||||
deja_presentes=len(nouvelles) - creees,
|
|
||||||
)
|
|
||||||
|
|||||||
@@ -1,117 +0,0 @@
|
|||||||
# Piège : `rule_reference` est la clé d'idempotence en base, portée par la contrainte
|
|
||||||
# `uq_recommendation_alert_rule`. Renommer une référence déjà livrée ne remplace pas les
|
|
||||||
# recommandations existantes, il en crée de nouvelles à côté. Une règle qui change de sens
|
|
||||||
# prend donc une référence suffixée `-v2` - REGLES.
|
|
||||||
|
|
||||||
from collections.abc import Callable
|
|
||||||
from dataclasses import dataclass
|
|
||||||
from typing import Final
|
|
||||||
|
|
||||||
from app.models.energy import Alert
|
|
||||||
from app.repositories.recommendation import NouvelleRecommandation
|
|
||||||
from app.schemas.alert import AlertSeverity, AlertType
|
|
||||||
|
|
||||||
FACTEUR_DEPASSEMENT_MAJEUR: Final = 1.2
|
|
||||||
POURCENTAGE_DEPASSEMENT_MAJEUR: Final = round((FACTEUR_DEPASSEMENT_MAJEUR - 1) * 100)
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True, slots=True)
|
|
||||||
class Regle:
|
|
||||||
reference: str
|
|
||||||
action: str
|
|
||||||
declencheur: Callable[[Alert], bool]
|
|
||||||
motif: Callable[[Alert], str]
|
|
||||||
|
|
||||||
|
|
||||||
def _du_type(attendu: AlertType) -> Callable[[Alert], bool]:
|
|
||||||
return lambda alerte: alerte.type == attendu
|
|
||||||
|
|
||||||
|
|
||||||
def _de_severite(attendue: AlertSeverity) -> Callable[[Alert], bool]:
|
|
||||||
return lambda alerte: alerte.severity == attendue
|
|
||||||
|
|
||||||
|
|
||||||
# Un seuil nul ou négatif rendrait le rapport `value / threshold` arbitraire : l'alerte ne
|
|
||||||
# renseigne alors aucun dépassement exploitable, et la règle ne se déclenche pas.
|
|
||||||
def _depasse_largement_le_seuil(alerte: Alert) -> bool:
|
|
||||||
if alerte.value is None or alerte.threshold is None or alerte.threshold <= 0:
|
|
||||||
return False
|
|
||||||
return alerte.value >= alerte.threshold * FACTEUR_DEPASSEMENT_MAJEUR
|
|
||||||
|
|
||||||
|
|
||||||
REGLES: Final[tuple[Regle, ...]] = (
|
|
||||||
Regle(
|
|
||||||
reference="spike-delestage-v1",
|
|
||||||
action="Délester les équipements non prioritaires sur le créneau du pic",
|
|
||||||
declencheur=_du_type(AlertType.SPIKE),
|
|
||||||
motif=lambda alerte: f"Pic de consommation signalé sur le site {alerte.site_id}",
|
|
||||||
),
|
|
||||||
Regle(
|
|
||||||
reference="threshold-reduction-v1",
|
|
||||||
action="Ramener la puissance appelée sous le seuil contractuel",
|
|
||||||
declencheur=_du_type(AlertType.THRESHOLD),
|
|
||||||
motif=lambda alerte: f"Seuil de consommation dépassé sur le site {alerte.site_id}",
|
|
||||||
),
|
|
||||||
Regle(
|
|
||||||
reference="outage-secours-v1",
|
|
||||||
action="Basculer sur l'alimentation de secours et prévenir l'exploitant",
|
|
||||||
declencheur=_du_type(AlertType.OUTAGE),
|
|
||||||
motif=lambda alerte: (
|
|
||||||
f"Risque de surcharge ou de coupure imminente sur le site {alerte.site_id}"
|
|
||||||
),
|
|
||||||
),
|
|
||||||
Regle(
|
|
||||||
reference="sensor-maintenance-v1",
|
|
||||||
action="Planifier une intervention de maintenance sur le capteur",
|
|
||||||
declencheur=_du_type(AlertType.SENSOR),
|
|
||||||
motif=lambda alerte: (
|
|
||||||
f"Capteur défaillant sur le site {alerte.site_id}, les mesures ne sont plus fiables"
|
|
||||||
),
|
|
||||||
),
|
|
||||||
Regle(
|
|
||||||
reference="anomaly-verification-v1",
|
|
||||||
action="Confronter la mesure à la prévision et vérifier le paramétrage du site",
|
|
||||||
declencheur=_du_type(AlertType.ANOMALY),
|
|
||||||
motif=lambda alerte: (
|
|
||||||
f"Écart anormal entre la mesure et le comportement attendu du site {alerte.site_id}"
|
|
||||||
),
|
|
||||||
),
|
|
||||||
Regle(
|
|
||||||
reference="escalade-astreinte-v1",
|
|
||||||
action="Escalader à l'astreinte sous une heure",
|
|
||||||
declencheur=_de_severite(AlertSeverity.CRITICAL),
|
|
||||||
motif=lambda alerte: f"Alerte de sévérité critique sur le site {alerte.site_id}",
|
|
||||||
),
|
|
||||||
Regle(
|
|
||||||
reference="contrat-puissance-v1",
|
|
||||||
action="Réévaluer la puissance souscrite au contrat",
|
|
||||||
declencheur=_depasse_largement_le_seuil,
|
|
||||||
motif=lambda alerte: (
|
|
||||||
f"Dépassement d'au moins {POURCENTAGE_DEPASSEMENT_MAJEUR} % du seuil "
|
|
||||||
f"sur le site {alerte.site_id}"
|
|
||||||
),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def applique_les_regles(alerte: Alert) -> list[NouvelleRecommandation]:
|
|
||||||
contexte = _contexte_de_mesure(alerte)
|
|
||||||
return [
|
|
||||||
NouvelleRecommandation(
|
|
||||||
alert_id=alerte.alert_id,
|
|
||||||
action=regle.action,
|
|
||||||
explanation=f"{regle.motif(alerte)}{contexte}.",
|
|
||||||
rule_reference=regle.reference,
|
|
||||||
)
|
|
||||||
for regle in REGLES
|
|
||||||
if regle.declencheur(alerte)
|
|
||||||
]
|
|
||||||
|
|
||||||
|
|
||||||
def _contexte_de_mesure(alerte: Alert) -> str:
|
|
||||||
if alerte.value is None:
|
|
||||||
return ""
|
|
||||||
grandeur = alerte.metric or "valeur"
|
|
||||||
if alerte.threshold is None:
|
|
||||||
return f" ({grandeur} mesurée à {alerte.value})"
|
|
||||||
return f" ({grandeur} mesurée à {alerte.value}, seuil {alerte.threshold})"
|
|
||||||
+1
-108
@@ -1453,90 +1453,6 @@
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"/api/v1/recommendations/generate": {
|
|
||||||
"post": {
|
|
||||||
"tags": [
|
|
||||||
"recommendations"
|
|
||||||
],
|
|
||||||
"summary": "Génère les recommandations à partir des alertes",
|
|
||||||
"operationId": "generate_recommendations_api_v1_recommendations_generate_post",
|
|
||||||
"security": [
|
|
||||||
{
|
|
||||||
"Jeton d'accès": []
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"parameters": [
|
|
||||||
{
|
|
||||||
"name": "site_id",
|
|
||||||
"in": "query",
|
|
||||||
"required": false,
|
|
||||||
"schema": {
|
|
||||||
"anyOf": [
|
|
||||||
{
|
|
||||||
"type": "string"
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"type": "null"
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"title": "Site Id"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"responses": {
|
|
||||||
"200": {
|
|
||||||
"description": "Successful Response",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"$ref": "#/components/schemas/RecommendationGenerationResponse"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"500": {
|
|
||||||
"description": "Erreur interne. `correlation` identifie la trace côté serveur, qui n'est pas renvoyée au client.",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"$ref": "#/components/schemas/InternalErrorResponse"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"401": {
|
|
||||||
"description": "Jeton absent, illisible, périmé, ou rendu caduc par un changement de rôle ou une désactivation. L'en-tête `WWW-Authenticate` porte la cause dans `error=`.",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"$ref": "#/components/schemas/ErrorResponse"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"403": {
|
|
||||||
"description": "Droits insuffisants, ou mot de passe provisoire à changer quand `detail` vaut `password_change_required`.",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"$ref": "#/components/schemas/ErrorResponse"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"422": {
|
|
||||||
"description": "Corps invalide. Le détail nomme le champ fautif et le type d'erreur, jamais la valeur envoyée.",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"$ref": "#/components/schemas/ValidationErrorResponse"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"/api/v1/stats/summary": {
|
"/api/v1/stats/summary": {
|
||||||
"get": {
|
"get": {
|
||||||
"tags": [
|
"tags": [
|
||||||
@@ -2428,29 +2344,6 @@
|
|||||||
],
|
],
|
||||||
"title": "ReadingSource"
|
"title": "ReadingSource"
|
||||||
},
|
},
|
||||||
"RecommendationGenerationResponse": {
|
|
||||||
"properties": {
|
|
||||||
"alerts_examined": {
|
|
||||||
"type": "integer",
|
|
||||||
"title": "Alerts Examined"
|
|
||||||
},
|
|
||||||
"recommendations_created": {
|
|
||||||
"type": "integer",
|
|
||||||
"title": "Recommendations Created"
|
|
||||||
},
|
|
||||||
"already_present": {
|
|
||||||
"type": "integer",
|
|
||||||
"title": "Already Present"
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"type": "object",
|
|
||||||
"required": [
|
|
||||||
"alerts_examined",
|
|
||||||
"recommendations_created",
|
|
||||||
"already_present"
|
|
||||||
],
|
|
||||||
"title": "RecommendationGenerationResponse"
|
|
||||||
},
|
|
||||||
"RecommendationResponse": {
|
"RecommendationResponse": {
|
||||||
"properties": {
|
"properties": {
|
||||||
"recommendation_id": {
|
"recommendation_id": {
|
||||||
@@ -3260,7 +3153,7 @@
|
|||||||
},
|
},
|
||||||
{
|
{
|
||||||
"name": "recommendations",
|
"name": "recommendations",
|
||||||
"description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`. Leur génération par le moteur de règles est réservée au rôle `admin`."
|
"description": "Consultation des recommandations issues des alertes. Accessible à partir du rôle `lecteur`."
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"name": "stats",
|
"name": "stats",
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ dependencies = [
|
|||||||
"argon2-cffi>=23.1",
|
"argon2-cffi>=23.1",
|
||||||
"anyio>=4.0",
|
"anyio>=4.0",
|
||||||
"aiosmtplib>=5.1.3",
|
"aiosmtplib>=5.1.3",
|
||||||
|
"httpx>=0.28.1",
|
||||||
"pandas>=3.0.5",
|
"pandas>=3.0.5",
|
||||||
]
|
]
|
||||||
|
|
||||||
@@ -27,7 +28,6 @@ dev = [
|
|||||||
"pytest>=9.1.1",
|
"pytest>=9.1.1",
|
||||||
"pytest-asyncio>=1.4.0",
|
"pytest-asyncio>=1.4.0",
|
||||||
"pytest-cov>=7.1.0",
|
"pytest-cov>=7.1.0",
|
||||||
"httpx>=0.28.1",
|
|
||||||
"pandas-stubs>=3.0.5.260914",
|
"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/alerts"): Role.LECTEUR,
|
||||||
("GET", "/api/v1/recommendations"): Role.LECTEUR,
|
("GET", "/api/v1/recommendations"): Role.LECTEUR,
|
||||||
("GET", "/api/v1/recommendations/{recommendation_id}"): Role.LECTEUR,
|
("GET", "/api/v1/recommendations/{recommendation_id}"): Role.LECTEUR,
|
||||||
("POST", "/api/v1/recommendations/generate"): Role.ADMIN,
|
|
||||||
("GET", "/api/v1/stats/summary"): Role.LECTEUR,
|
("GET", "/api/v1/stats/summary"): Role.LECTEUR,
|
||||||
("GET", "/api/v1/readings"): Role.LECTEUR,
|
("GET", "/api/v1/readings"): Role.LECTEUR,
|
||||||
("GET", "/api/v1/predictions"): Role.LECTEUR,
|
("GET", "/api/v1/predictions"): Role.LECTEUR,
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ from app.api.deps import get_current_principal, get_recommendation_service
|
|||||||
from app.core.principal import Principal
|
from app.core.principal import Principal
|
||||||
from app.core.roles import AccountKind, Role
|
from app.core.roles import AccountKind, Role
|
||||||
from app.models.energy import Recommendation
|
from app.models.energy import Recommendation
|
||||||
from app.services.recommendation import RapportGeneration, RecommendationNotFoundError
|
from app.services.recommendation import RecommendationNotFoundError
|
||||||
|
|
||||||
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
||||||
|
|
||||||
@@ -40,7 +40,6 @@ class FauxService:
|
|||||||
def __init__(self, erreur: Exception | None = None) -> None:
|
def __init__(self, erreur: Exception | None = None) -> None:
|
||||||
self._erreur = erreur
|
self._erreur = erreur
|
||||||
self.recommendation = recommendation()
|
self.recommendation = recommendation()
|
||||||
self.site_demande: str | None = None
|
|
||||||
|
|
||||||
async def list_all(self) -> list[Recommendation]:
|
async def list_all(self) -> list[Recommendation]:
|
||||||
return [self.recommendation]
|
return [self.recommendation]
|
||||||
@@ -50,10 +49,6 @@ class FauxService:
|
|||||||
raise self._erreur
|
raise self._erreur
|
||||||
return self.recommendation
|
return self.recommendation
|
||||||
|
|
||||||
async def generate(self, *, site_id: str | None = None) -> RapportGeneration:
|
|
||||||
self.site_demande = site_id
|
|
||||||
return RapportGeneration(alertes_examinees=2, recommandations_creees=3, deja_presentes=1)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
def lecteur_connecte(app: FastAPI) -> Iterator[None]:
|
def lecteur_connecte(app: FastAPI) -> Iterator[None]:
|
||||||
@@ -147,56 +142,3 @@ async def test_get_recommendation_returns_404_when_the_session_finds_nothing(
|
|||||||
response = await client.get("/api/v1/recommendations/404")
|
response = await client.get("/api/v1/recommendations/404")
|
||||||
|
|
||||||
assert response.status_code == 404
|
assert response.status_code == 404
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
|
||||||
def admin_connecte(app: FastAPI) -> Iterator[None]:
|
|
||||||
app.dependency_overrides[get_current_principal] = lambda: principal(Role.ADMIN)
|
|
||||||
yield
|
|
||||||
app.dependency_overrides.pop(get_current_principal, None)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
|
||||||
def servi_en_admin(app: FastAPI, admin_connecte: None) -> Iterator[Callable[[], FauxService]]:
|
|
||||||
def installe() -> FauxService:
|
|
||||||
service = FauxService()
|
|
||||||
app.dependency_overrides[get_recommendation_service] = lambda: service
|
|
||||||
return service
|
|
||||||
|
|
||||||
yield installe
|
|
||||||
app.dependency_overrides.pop(get_recommendation_service, None)
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_recommendations_returns_the_generation_report(
|
|
||||||
servi_en_admin: Callable[[], FauxService], client: AsyncClient
|
|
||||||
) -> None:
|
|
||||||
servi_en_admin()
|
|
||||||
|
|
||||||
response = await client.post("/api/v1/recommendations/generate")
|
|
||||||
|
|
||||||
assert response.status_code == 200
|
|
||||||
assert response.json() == {
|
|
||||||
"alerts_examined": 2,
|
|
||||||
"recommendations_created": 3,
|
|
||||||
"already_present": 1,
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_recommendations_forwards_the_requested_site(
|
|
||||||
servi_en_admin: Callable[[], FauxService], client: AsyncClient
|
|
||||||
) -> None:
|
|
||||||
service = servi_en_admin()
|
|
||||||
|
|
||||||
await client.post("/api/v1/recommendations/generate", params={"site_id": "SITE002"})
|
|
||||||
|
|
||||||
assert service.site_demande == "SITE002"
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_recommendations_refuses_a_reader(
|
|
||||||
servi: Callable[..., FauxService], client: AsyncClient
|
|
||||||
) -> None:
|
|
||||||
servi()
|
|
||||||
|
|
||||||
response = await client.post("/api/v1/recommendations/generate")
|
|
||||||
|
|
||||||
assert response.status_code == 403
|
|
||||||
|
|||||||
@@ -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 sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
from app.models.energy import Alert, Recommendation, Site
|
from app.models.energy import Alert, Recommendation, Site
|
||||||
from app.repositories import recommendation as module_recommendation
|
from app.repositories.recommendation import RecommendationRepository
|
||||||
from app.repositories.recommendation import NouvelleRecommandation, RecommendationRepository
|
|
||||||
|
|
||||||
pytestmark = pytest.mark.integration
|
pytestmark = pytest.mark.integration
|
||||||
|
|
||||||
@@ -84,59 +83,3 @@ async def test_list_all_returns_the_recommendations_sorted_by_identifier(
|
|||||||
await session.rollback()
|
await session.rollback()
|
||||||
|
|
||||||
assert identifiants == sorted(identifiants)
|
assert identifiants == sorted(identifiants)
|
||||||
|
|
||||||
|
|
||||||
def nouvelle(alert_id: int, reference: str = "spike-delestage-v1") -> NouvelleRecommandation:
|
|
||||||
return NouvelleRecommandation(
|
|
||||||
alert_id=alert_id,
|
|
||||||
action="Délester les équipements non prioritaires",
|
|
||||||
explanation="Pic de consommation signalé.",
|
|
||||||
rule_reference=reference,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
async def test_create_missing_inserts_the_proposals(session: AsyncSession) -> None:
|
|
||||||
depot = RecommendationRepository(session)
|
|
||||||
alert_id = await creer_alerte(session)
|
|
||||||
|
|
||||||
creees = await depot.create_missing(
|
|
||||||
[nouvelle(alert_id), nouvelle(alert_id, "escalade-astreinte-v1")]
|
|
||||||
)
|
|
||||||
await session.rollback()
|
|
||||||
|
|
||||||
assert creees == 2
|
|
||||||
|
|
||||||
|
|
||||||
async def test_create_missing_ignores_a_rule_already_held_for_the_alert(
|
|
||||||
session: AsyncSession,
|
|
||||||
) -> None:
|
|
||||||
depot = RecommendationRepository(session)
|
|
||||||
alert_id = await creer_alerte(session)
|
|
||||||
await depot.create_missing([nouvelle(alert_id)])
|
|
||||||
|
|
||||||
creees = await depot.create_missing([nouvelle(alert_id)])
|
|
||||||
await session.rollback()
|
|
||||||
|
|
||||||
assert creees == 0
|
|
||||||
|
|
||||||
|
|
||||||
async def test_create_missing_returns_zero_without_any_proposal(session: AsyncSession) -> None:
|
|
||||||
creees = await RecommendationRepository(session).create_missing([])
|
|
||||||
|
|
||||||
assert creees == 0
|
|
||||||
|
|
||||||
|
|
||||||
async def test_create_missing_inserts_every_proposal_across_several_batches(
|
|
||||||
session: AsyncSession, monkeypatch: pytest.MonkeyPatch
|
|
||||||
) -> None:
|
|
||||||
monkeypatch.setattr(module_recommendation, "TAILLE_DE_LOT", 2)
|
|
||||||
depot = RecommendationRepository(session)
|
|
||||||
alert_id = await creer_alerte(session)
|
|
||||||
propositions = [nouvelle(alert_id, f"regle-{index}-v1") for index in range(5)]
|
|
||||||
|
|
||||||
creees = await depot.create_missing(propositions)
|
|
||||||
enregistrees = [r for r in await depot.list_all() if r.alert_id == alert_id]
|
|
||||||
await session.rollback()
|
|
||||||
|
|
||||||
assert creees == 5
|
|
||||||
assert len(enregistrees) == 5
|
|
||||||
|
|||||||
@@ -1,14 +1,10 @@
|
|||||||
from collections.abc import Sequence
|
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from app.models.energy import Alert, Recommendation
|
from app.models.energy import Recommendation
|
||||||
from app.repositories.recommendation import NouvelleRecommandation
|
|
||||||
from app.services.recommendation import RecommendationNotFoundError, RecommendationService
|
from app.services.recommendation import RecommendationNotFoundError, RecommendationService
|
||||||
|
|
||||||
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
|
||||||
|
|
||||||
|
|
||||||
def recommendation(recommendation_id: int = 1) -> Recommendation:
|
def recommendation(recommendation_id: int = 1) -> Recommendation:
|
||||||
return Recommendation(
|
return Recommendation(
|
||||||
@@ -17,33 +13,13 @@ def recommendation(recommendation_id: int = 1) -> Recommendation:
|
|||||||
action="Vérifier la consommation",
|
action="Vérifier la consommation",
|
||||||
explanation="Pic détecté",
|
explanation="Pic détecté",
|
||||||
rule_reference="spike-v1",
|
rule_reference="spike-v1",
|
||||||
created_at=MOMENT,
|
created_at=datetime(2024, 1, 1, tzinfo=UTC),
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def alerte(alert_id: int = 1, site_id: str = "SITE001", severity: str = "high") -> Alert:
|
|
||||||
return Alert(
|
|
||||||
alert_id=alert_id,
|
|
||||||
source_alert_id=f"ALR-{alert_id}",
|
|
||||||
site_id=site_id,
|
|
||||||
source="api_mock",
|
|
||||||
timestamp=MOMENT,
|
|
||||||
type="spike",
|
|
||||||
severity=severity,
|
|
||||||
message="Pic de consommation",
|
|
||||||
value=None,
|
|
||||||
threshold=None,
|
|
||||||
metric=None,
|
|
||||||
prediction_id=None,
|
|
||||||
raw_data={},
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
class FakeRepository:
|
class FakeRepository:
|
||||||
def __init__(self, recommendations: list[Recommendation], creees: int | None = None) -> None:
|
def __init__(self, recommendations: list[Recommendation]) -> None:
|
||||||
self._recommendations = recommendations
|
self._recommendations = recommendations
|
||||||
self._creees = creees
|
|
||||||
self.recues: list[NouvelleRecommandation] = []
|
|
||||||
|
|
||||||
async def list_all(self) -> list[Recommendation]:
|
async def list_all(self) -> list[Recommendation]:
|
||||||
return self._recommendations
|
return self._recommendations
|
||||||
@@ -53,111 +29,27 @@ class FakeRepository:
|
|||||||
(r for r in self._recommendations if r.recommendation_id == recommendation_id), None
|
(r for r in self._recommendations if r.recommendation_id == recommendation_id), None
|
||||||
)
|
)
|
||||||
|
|
||||||
async def create_missing(self, nouvelles: Sequence[NouvelleRecommandation]) -> int:
|
|
||||||
self.recues = list(nouvelles)
|
|
||||||
return len(self.recues) if self._creees is None else self._creees
|
|
||||||
|
|
||||||
|
|
||||||
class FakeAlertRepository:
|
|
||||||
def __init__(self, alertes: list[Alert]) -> None:
|
|
||||||
self._alertes = alertes
|
|
||||||
self.site_demande: str | None = None
|
|
||||||
|
|
||||||
async def list_all(
|
|
||||||
self, *, site_id: str | None = None, severity: str | None = None
|
|
||||||
) -> list[Alert]:
|
|
||||||
self.site_demande = site_id
|
|
||||||
if site_id is None:
|
|
||||||
return self._alertes
|
|
||||||
return [a for a in self._alertes if a.site_id == site_id]
|
|
||||||
|
|
||||||
|
|
||||||
class FakeTransaction:
|
|
||||||
def __init__(self) -> None:
|
|
||||||
self.commits = 0
|
|
||||||
|
|
||||||
async def commit(self) -> None:
|
|
||||||
self.commits += 1
|
|
||||||
|
|
||||||
|
|
||||||
def service(
|
|
||||||
recommendations: FakeRepository | None = None,
|
|
||||||
alerts: FakeAlertRepository | None = None,
|
|
||||||
transaction: FakeTransaction | None = None,
|
|
||||||
) -> RecommendationService:
|
|
||||||
return RecommendationService(
|
|
||||||
recommendations=recommendations or FakeRepository([]),
|
|
||||||
alerts=alerts or FakeAlertRepository([]),
|
|
||||||
transaction=transaction or FakeTransaction(),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
async def test_list_all_returns_the_repository_recommendations() -> None:
|
async def test_list_all_returns_the_repository_recommendations() -> None:
|
||||||
depot = FakeRepository([recommendation(1), recommendation(2)])
|
service = RecommendationService(
|
||||||
|
recommendations=FakeRepository([recommendation(1), recommendation(2)])
|
||||||
|
)
|
||||||
|
|
||||||
recommendations = await service(recommendations=depot).list_all()
|
recommendations = await service.list_all()
|
||||||
|
|
||||||
assert [r.recommendation_id for r in recommendations] == [1, 2]
|
assert [r.recommendation_id for r in recommendations] == [1, 2]
|
||||||
|
|
||||||
|
|
||||||
async def test_get_by_id_returns_the_matching_recommendation() -> None:
|
async def test_get_by_id_returns_the_matching_recommendation() -> None:
|
||||||
trouve = await service(recommendations=FakeRepository([recommendation(1)])).get_by_id(1)
|
service = RecommendationService(recommendations=FakeRepository([recommendation(1)]))
|
||||||
|
|
||||||
|
trouve = await service.get_by_id(1)
|
||||||
|
|
||||||
assert trouve.recommendation_id == 1
|
assert trouve.recommendation_id == 1
|
||||||
|
|
||||||
|
|
||||||
async def test_get_by_id_raises_when_the_recommendation_is_unknown() -> None:
|
async def test_get_by_id_raises_when_the_recommendation_is_unknown() -> None:
|
||||||
|
service = RecommendationService(recommendations=FakeRepository([]))
|
||||||
|
|
||||||
with pytest.raises(RecommendationNotFoundError):
|
with pytest.raises(RecommendationNotFoundError):
|
||||||
await service().get_by_id(404)
|
await service.get_by_id(404)
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_persists_one_proposal_per_triggered_rule() -> None:
|
|
||||||
depot = FakeRepository([])
|
|
||||||
|
|
||||||
rapport = await service(
|
|
||||||
recommendations=depot, alerts=FakeAlertRepository([alerte(severity="critical")])
|
|
||||||
).generate()
|
|
||||||
|
|
||||||
assert {n.rule_reference for n in depot.recues} == {
|
|
||||||
"spike-delestage-v1",
|
|
||||||
"escalade-astreinte-v1",
|
|
||||||
}
|
|
||||||
assert rapport.recommandations_creees == 2
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_commits_once() -> None:
|
|
||||||
transaction = FakeTransaction()
|
|
||||||
|
|
||||||
await service(alerts=FakeAlertRepository([alerte()]), transaction=transaction).generate()
|
|
||||||
|
|
||||||
assert transaction.commits == 1
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_restricts_the_alerts_to_the_requested_site() -> None:
|
|
||||||
alertes = FakeAlertRepository([alerte(1, site_id="SITE001"), alerte(2, site_id="SITE002")])
|
|
||||||
depot = FakeRepository([])
|
|
||||||
|
|
||||||
rapport = await service(recommendations=depot, alerts=alertes).generate(site_id="SITE002")
|
|
||||||
|
|
||||||
assert alertes.site_demande == "SITE002"
|
|
||||||
assert rapport.alertes_examinees == 1
|
|
||||||
assert {n.alert_id for n in depot.recues} == {2}
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_reports_nothing_when_no_alert_matches() -> None:
|
|
||||||
rapport = await service().generate()
|
|
||||||
|
|
||||||
assert rapport.alertes_examinees == 0
|
|
||||||
assert rapport.recommandations_creees == 0
|
|
||||||
assert rapport.deja_presentes == 0
|
|
||||||
|
|
||||||
|
|
||||||
async def test_generate_counts_the_proposals_the_database_already_held() -> None:
|
|
||||||
depot = FakeRepository([], creees=0)
|
|
||||||
|
|
||||||
rapport = await service(
|
|
||||||
recommendations=depot, alerts=FakeAlertRepository([alerte()])
|
|
||||||
).generate()
|
|
||||||
|
|
||||||
assert rapport.recommandations_creees == 0
|
|
||||||
assert rapport.deja_presentes == 1
|
|
||||||
|
|||||||
@@ -1,142 +0,0 @@
|
|||||||
from datetime import UTC, datetime
|
|
||||||
|
|
||||||
import pytest
|
|
||||||
|
|
||||||
from app.models.energy import Alert
|
|
||||||
from app.services.recommendation_rules import FACTEUR_DEPASSEMENT_MAJEUR, applique_les_regles
|
|
||||||
|
|
||||||
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
|
|
||||||
|
|
||||||
|
|
||||||
def alerte(
|
|
||||||
*,
|
|
||||||
alert_id: int = 1,
|
|
||||||
type_alerte: str = "spike",
|
|
||||||
severity: str = "high",
|
|
||||||
value: float | None = None,
|
|
||||||
threshold: float | None = None,
|
|
||||||
metric: str | None = None,
|
|
||||||
site_id: str = "SITE001",
|
|
||||||
) -> Alert:
|
|
||||||
return Alert(
|
|
||||||
alert_id=alert_id,
|
|
||||||
source_alert_id=f"ALR-{alert_id}",
|
|
||||||
site_id=site_id,
|
|
||||||
source="api_mock",
|
|
||||||
timestamp=MOMENT,
|
|
||||||
type=type_alerte,
|
|
||||||
severity=severity,
|
|
||||||
message="Alerte de test",
|
|
||||||
value=value,
|
|
||||||
threshold=threshold,
|
|
||||||
metric=metric,
|
|
||||||
prediction_id=None,
|
|
||||||
raw_data={},
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize(
|
|
||||||
("type_alerte", "attendue"),
|
|
||||||
[
|
|
||||||
("spike", "spike-delestage-v1"),
|
|
||||||
("threshold", "threshold-reduction-v1"),
|
|
||||||
("outage", "outage-secours-v1"),
|
|
||||||
("sensor", "sensor-maintenance-v1"),
|
|
||||||
("anomaly", "anomaly-verification-v1"),
|
|
||||||
],
|
|
||||||
ids=["pic", "seuil", "coupure", "capteur", "anomalie"],
|
|
||||||
)
|
|
||||||
def test_each_alert_type_yields_its_own_rule(type_alerte: str, attendue: str) -> None:
|
|
||||||
proposees = applique_les_regles(alerte(type_alerte=type_alerte))
|
|
||||||
|
|
||||||
assert [p.rule_reference for p in proposees] == [attendue]
|
|
||||||
|
|
||||||
|
|
||||||
def test_a_critical_alert_adds_the_escalation_rule() -> None:
|
|
||||||
proposees = applique_les_regles(alerte(severity="critical"))
|
|
||||||
|
|
||||||
assert "escalade-astreinte-v1" in {p.rule_reference for p in proposees}
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize("severity", ["low", "medium", "high"], ids=["faible", "moyenne", "haute"])
|
|
||||||
def test_a_non_critical_alert_does_not_escalate(severity: str) -> None:
|
|
||||||
proposees = applique_les_regles(alerte(severity=severity))
|
|
||||||
|
|
||||||
assert "escalade-astreinte-v1" not in {p.rule_reference for p in proposees}
|
|
||||||
|
|
||||||
|
|
||||||
def test_a_large_overshoot_adds_the_contract_rule() -> None:
|
|
||||||
proposees = applique_les_regles(
|
|
||||||
alerte(value=720.0 * FACTEUR_DEPASSEMENT_MAJEUR, threshold=720.0)
|
|
||||||
)
|
|
||||||
|
|
||||||
assert "contrat-puissance-v1" in {p.rule_reference for p in proposees}
|
|
||||||
|
|
||||||
|
|
||||||
def test_an_overshoot_below_the_factor_does_not_add_the_contract_rule() -> None:
|
|
||||||
proposees = applique_les_regles(alerte(value=800.0, threshold=720.0))
|
|
||||||
|
|
||||||
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize(
|
|
||||||
("value", "threshold"),
|
|
||||||
[(None, 720.0), (900.0, None), (900.0, 0.0), (900.0, -10.0)],
|
|
||||||
ids=["sans mesure", "sans seuil", "seuil nul", "seuil negatif"],
|
|
||||||
)
|
|
||||||
def test_the_contract_rule_stays_silent_without_an_exploitable_threshold(
|
|
||||||
value: float | None, threshold: float | None
|
|
||||||
) -> None:
|
|
||||||
proposees = applique_les_regles(alerte(value=value, threshold=threshold))
|
|
||||||
|
|
||||||
assert "contrat-puissance-v1" not in {p.rule_reference for p in proposees}
|
|
||||||
|
|
||||||
|
|
||||||
def test_the_explanation_quotes_the_measure_and_the_threshold() -> None:
|
|
||||||
proposees = applique_les_regles(alerte(value=812.5, threshold=720.0, metric="consumption_kw"))
|
|
||||||
|
|
||||||
assert "(consumption_kw mesurée à 812.5, seuil 720.0)" in proposees[0].explanation
|
|
||||||
|
|
||||||
|
|
||||||
def test_the_explanation_quotes_the_measure_alone_when_no_threshold_is_known() -> None:
|
|
||||||
proposees = applique_les_regles(alerte(value=812.5, metric="consumption_kw"))
|
|
||||||
|
|
||||||
assert "(consumption_kw mesurée à 812.5)" in proposees[0].explanation
|
|
||||||
|
|
||||||
|
|
||||||
def test_the_explanation_omits_the_measure_when_the_alert_carries_none() -> None:
|
|
||||||
proposees = applique_les_regles(alerte())
|
|
||||||
|
|
||||||
assert "(" not in proposees[0].explanation
|
|
||||||
|
|
||||||
|
|
||||||
def test_the_explanation_names_the_site() -> None:
|
|
||||||
proposees = applique_les_regles(alerte(site_id="SITE042"))
|
|
||||||
|
|
||||||
assert "SITE042" in proposees[0].explanation
|
|
||||||
|
|
||||||
|
|
||||||
def test_every_proposal_carries_the_alert_identifier() -> None:
|
|
||||||
proposees = applique_les_regles(alerte(alert_id=77, severity="critical"))
|
|
||||||
|
|
||||||
assert {p.alert_id for p in proposees} == {77}
|
|
||||||
|
|
||||||
|
|
||||||
def test_an_alert_never_yields_the_same_rule_twice() -> None:
|
|
||||||
proposees = applique_les_regles(
|
|
||||||
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
|
|
||||||
)
|
|
||||||
|
|
||||||
assert len(proposees) == len({p.rule_reference for p in proposees})
|
|
||||||
|
|
||||||
|
|
||||||
def test_a_critical_alert_over_the_threshold_yields_the_three_rules() -> None:
|
|
||||||
proposees = applique_les_regles(
|
|
||||||
alerte(severity="critical", value=900.0, threshold=720.0, metric="consumption_kw")
|
|
||||||
)
|
|
||||||
|
|
||||||
assert {p.rule_reference for p in proposees} == {
|
|
||||||
"spike-delestage-v1",
|
|
||||||
"escalade-astreinte-v1",
|
|
||||||
"contrat-puissance-v1",
|
|
||||||
}
|
|
||||||
@@ -118,30 +118,3 @@ def test_main_exports_the_contract_without_asking_for_a_password(
|
|||||||
assert code == 0
|
assert code == 0
|
||||||
assert destination.exists()
|
assert destination.exists()
|
||||||
assert str(destination) in capsys.readouterr().out
|
assert str(destination) in capsys.readouterr().out
|
||||||
|
|
||||||
|
|
||||||
def test_build_parser_reads_the_generate_recommendations_arguments() -> None:
|
|
||||||
arguments = cli.build_parser().parse_args(["generate-recommendations", "--site-id", "SITE002"])
|
|
||||||
|
|
||||||
assert arguments.commande == "generate-recommendations"
|
|
||||||
assert arguments.site_id == "SITE002"
|
|
||||||
|
|
||||||
|
|
||||||
def test_build_parser_defaults_the_generation_to_every_site() -> None:
|
|
||||||
arguments = cli.build_parser().parse_args(["generate-recommendations"])
|
|
||||||
|
|
||||||
assert arguments.site_id is None
|
|
||||||
|
|
||||||
|
|
||||||
def test_main_generates_the_recommendations_without_asking_for_a_password(
|
|
||||||
monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str]
|
|
||||||
) -> None:
|
|
||||||
async def fausse_generation(*, site_id: str | None) -> str:
|
|
||||||
return f"génération lancée pour {site_id}"
|
|
||||||
|
|
||||||
monkeypatch.setattr(cli, "generate_recommendations", fausse_generation)
|
|
||||||
|
|
||||||
code = cli.main(["generate-recommendations", "--site-id", "SITE002"])
|
|
||||||
|
|
||||||
assert code == 0
|
|
||||||
assert "SITE002" in capsys.readouterr().out
|
|
||||||
|
|||||||
Generated
+2
-2
@@ -326,6 +326,7 @@ dependencies = [
|
|||||||
{ name = "argon2-cffi" },
|
{ name = "argon2-cffi" },
|
||||||
{ name = "asyncpg" },
|
{ name = "asyncpg" },
|
||||||
{ name = "fastapi" },
|
{ name = "fastapi" },
|
||||||
|
{ name = "httpx" },
|
||||||
{ name = "pandas" },
|
{ name = "pandas" },
|
||||||
{ name = "prometheus-fastapi-instrumentator" },
|
{ name = "prometheus-fastapi-instrumentator" },
|
||||||
{ name = "pydantic", extra = ["email"] },
|
{ name = "pydantic", extra = ["email"] },
|
||||||
@@ -338,7 +339,6 @@ dependencies = [
|
|||||||
|
|
||||||
[package.dev-dependencies]
|
[package.dev-dependencies]
|
||||||
dev = [
|
dev = [
|
||||||
{ name = "httpx" },
|
|
||||||
{ name = "mypy" },
|
{ name = "mypy" },
|
||||||
{ name = "pandas-stubs" },
|
{ name = "pandas-stubs" },
|
||||||
{ name = "pytest" },
|
{ name = "pytest" },
|
||||||
@@ -355,6 +355,7 @@ requires-dist = [
|
|||||||
{ name = "argon2-cffi", specifier = ">=23.1" },
|
{ name = "argon2-cffi", specifier = ">=23.1" },
|
||||||
{ name = "asyncpg", specifier = ">=0.31.0" },
|
{ name = "asyncpg", specifier = ">=0.31.0" },
|
||||||
{ name = "fastapi", specifier = ">=0.141.1" },
|
{ name = "fastapi", specifier = ">=0.141.1" },
|
||||||
|
{ name = "httpx", specifier = ">=0.28.1" },
|
||||||
{ name = "pandas", specifier = ">=3.0.5" },
|
{ name = "pandas", specifier = ">=3.0.5" },
|
||||||
{ name = "prometheus-fastapi-instrumentator", specifier = ">=8.1.0" },
|
{ name = "prometheus-fastapi-instrumentator", specifier = ">=8.1.0" },
|
||||||
{ name = "pydantic", extras = ["email"], specifier = ">=2.13.5" },
|
{ name = "pydantic", extras = ["email"], specifier = ">=2.13.5" },
|
||||||
@@ -367,7 +368,6 @@ requires-dist = [
|
|||||||
|
|
||||||
[package.metadata.requires-dev]
|
[package.metadata.requires-dev]
|
||||||
dev = [
|
dev = [
|
||||||
{ name = "httpx", specifier = ">=0.28.1" },
|
|
||||||
{ name = "mypy", specifier = ">=2.3.1" },
|
{ name = "mypy", specifier = ">=2.3.1" },
|
||||||
{ name = "pandas-stubs", specifier = ">=3.0.5.260914" },
|
{ name = "pandas-stubs", specifier = ">=3.0.5.260914" },
|
||||||
{ name = "pytest", specifier = ">=9.1.1" },
|
{ name = "pytest", specifier = ">=9.1.1" },
|
||||||
|
|||||||
@@ -50,6 +50,12 @@ services:
|
|||||||
APP_SECRET_KEY: ${APP_SECRET_KEY:?}
|
APP_SECRET_KEY: ${APP_SECRET_KEY:?}
|
||||||
APP_CORS_ORIGINS: ${APP_CORS_ORIGINS:-http://localhost:4200}
|
APP_CORS_ORIGINS: ${APP_CORS_ORIGINS:-http://localhost:4200}
|
||||||
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
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_FRONTEND_RESET_PASSWORD_URL: ${APP_FRONTEND_RESET_PASSWORD_URL:-http://localhost:4200/reset-password}
|
||||||
APP_SMTP_HOST: mailpit
|
APP_SMTP_HOST: mailpit
|
||||||
APP_SMTP_PORT: "1025"
|
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 |
|
| ADR | Objet |
|
||||||
|---|---|
|
|---|---|
|
||||||
| [0001](../adr/0001-postgresql-timescaledb.md) | PostgreSQL 17 avec l'extension TimescaleDB, et la frontière `db/` vs `alembic/` |
|
| [0001](../adr/0001-postgresql-timescaledb.md) | PostgreSQL 17 avec l'extension TimescaleDB, et la frontière `db/` vs `alembic/` |
|
||||||
| [0002](../adr/0002-authentification-jwt-et-refresh-opaque.md) | Authentification par JWT d'accès et jeton de rafraîchissement opaque |
|
|
||||||
| [0003](../adr/0003-autorisation-rbac-a-trois-roles.md) | Autorisation RBAC à trois rôles, avec relecture du compte à chaque requête |
|
|
||||||
| [0004](../adr/0004-journal-d-audit-en-ajout-seul.md) | Journal d'audit en ajout seul, garanti par PostgreSQL |
|
|
||||||
| [0005](../adr/0005-modele-prediction-lightgbm.md) | Modèle de prédiction de consommation : LightGBM |
|
|
||||||
| [0006](../adr/0006-moteur-de-regles-dans-le-backend.md) | Le moteur de règles de recommandation vit dans le backend, pas dans `ml/` |
|
|
||||||
|
|||||||
@@ -146,7 +146,6 @@ Deux fichiers d'environnement, deux usages : `.env` à la racine alimente `docke
|
|||||||
| GET | `/api/v1/alerts` | Liste les alertes, filtrable par `site_id` et `severity`. `lecteur` | 401, 403, 422, 500 |
|
| GET | `/api/v1/alerts` | Liste les alertes, filtrable par `site_id` et `severity`. `lecteur` | 401, 403, 422, 500 |
|
||||||
| GET | `/api/v1/recommendations` | Liste les recommandations. `lecteur` | 401, 403, 500 |
|
| GET | `/api/v1/recommendations` | Liste les recommandations. `lecteur` | 401, 403, 500 |
|
||||||
| GET | `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation. `lecteur` | 401, 403, 404, 422, 500 |
|
| GET | `/api/v1/recommendations/{recommendation_id}` | Décrit une recommandation. `lecteur` | 401, 403, 404, 422, 500 |
|
||||||
| POST | `/api/v1/recommendations/generate` | Applique le moteur de règles aux alertes, filtrable par `site_id`. `admin` | 401, 403, 422, 500 |
|
|
||||||
| GET | `/api/v1/stats/summary` | Résume la consommation instantanée du parc. `lecteur` | 401, 403, 500 |
|
| GET | `/api/v1/stats/summary` | Résume la consommation instantanée du parc. `lecteur` | 401, 403, 500 |
|
||||||
| GET | `/api/v1/readings` | Historique des lectures, filtrable par `site_id`, fenêtre `start`/`end` (24h par défaut, 90 jours maximum) et paginé par `limit`/`offset`. `lecteur` | 400, 401, 403, 422, 500 |
|
| GET | `/api/v1/readings` | Historique des lectures, filtrable par `site_id`, fenêtre `start`/`end` (24h par défaut, 90 jours maximum) et paginé par `limit`/`offset`. `lecteur` | 400, 401, 403, 422, 500 |
|
||||||
| GET | `/api/v1/sensors/status` | État de santé des capteurs par site, dérivé de la dernière lecture. `admin` | 401, 403, 500 |
|
| GET | `/api/v1/sensors/status` | État de santé des capteurs par site, dérivé de la dernière lecture. `admin` | 401, 403, 500 |
|
||||||
@@ -195,21 +194,6 @@ plutôt qu'un statut inventé : le domaine `available`/`insufficient_data`/`erro
|
|||||||
LightGBM elle-même ; elle lit ce que le pipeline de scoring a déjà écrit, cf.
|
LightGBM elle-même ; elle lit ce que le pipeline de scoring a déjà écrit, cf.
|
||||||
[ML-START.md](../../ML-START.md) section 3.
|
[ML-START.md](../../ML-START.md) section 3.
|
||||||
|
|
||||||
`POST /recommendations/generate` est la seule route d'écriture métier du contrat. Elle applique
|
|
||||||
le moteur de règles d'`app/services/recommendation_rules.py` aux lignes d'`alert`, sans modèle ni
|
|
||||||
feature ML : le catalogue `REGLES` associe à chaque type et à chaque gravité d'alerte une action et
|
|
||||||
son explication, et une même alerte peut en déclencher plusieurs, comme le prévoit
|
|
||||||
[40-data.md](40-data.md). L'idempotence est portée par la base, pas par le service :
|
|
||||||
`RecommendationRepository.create_missing()` insère en `ON CONFLICT DO NOTHING` sur
|
|
||||||
`uq_recommendation_alert_rule`, donc rejouer la génération sur les mêmes alertes ne crée rien et
|
|
||||||
le rapport rendu distingue `recommendations_created` de `already_present`. Le même traitement est
|
|
||||||
disponible hors HTTP par `python -m app.cli generate-recommendations` (cible `make
|
|
||||||
recommendations`), sur le patron de `make ml-score`. Le choix de loger le moteur dans le backend
|
|
||||||
plutôt que dans `ml/` est justifié par l'[ADR 0006](../adr/0006-moteur-de-regles-dans-le-backend.md).
|
|
||||||
Les alertes traitées sont celles qu'écrit la détection interne (#104, section ci-dessous) : la
|
|
||||||
génération ne rend donc de recommandations qu'une fois la détection passée. L'insertion est
|
|
||||||
découpée en lots de `TAILLE_DE_LOT` lignes, asyncpg plafonnant une requête à 32 767 paramètres.
|
|
||||||
|
|
||||||
`GET /readings` reprend le même gabarit mais s'en écarte sur un point : `reading` est l'hypertable,
|
`GET /readings` reprend le même gabarit mais s'en écarte sur un point : `reading` est l'hypertable,
|
||||||
donc la seule table métier pouvant porter des années d'historique, ce que `docs/architecture/
|
donc la seule table métier pouvant porter des années d'historique, ce que `docs/architecture/
|
||||||
owasp-traceabilite.md` documentait comme un risque ouvert (API4, aucune pagination plafonnée ni
|
owasp-traceabilite.md` documentait comme un risque ouvert (API4, aucune pagination plafonnée ni
|
||||||
|
|||||||
+331
-74
@@ -6,10 +6,14 @@ système qui en découle.
|
|||||||
|
|
||||||
## Ce que couvre ce document
|
## Ce que couvre ce document
|
||||||
|
|
||||||
**Dix tables applicatives existent** : quatre pour l'authentification, six pour les données
|
**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 le code. Celles
|
d'énergie, dont l'hypertable `reading`.
|
||||||
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.
|
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
|
## Trois emplacements, trois rôles
|
||||||
|
|
||||||
@@ -35,8 +39,16 @@ Statut : `Fait`.
|
|||||||
- `db/init/100-extensions.sql` crée l'extension `timescaledb`.
|
- `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
|
- `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par
|
||||||
`apps/backend/tests/conftest.py`.
|
`apps/backend/tests/conftest.py`.
|
||||||
- Cinq révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle
|
- Six révisions Alembic sont actuellement appliquées.
|
||||||
établit `alembic_version` et refuse de s'appliquer si l'extension manque :
|
- 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
|
```sql
|
||||||
IF NOT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'timescaledb') THEN
|
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
|
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.
|
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
|
## Cycle de vie d'une mesure
|
||||||
|
|
||||||
Statut : `Cible`, sauf l'hypertable `reading` qui existe. Ni l'ingestion, ni les agrégats
|
Statut : `Partiellement fait`.
|
||||||
continus, ni la compression, ni la rétention ne sont écrits.
|
|
||||||
|
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
|
```mermaid
|
||||||
flowchart LR
|
flowchart LR
|
||||||
src["Source de mesures"] -.-> ing["Ingestion Airflow"]
|
csv["CSV + JSON"] --> hist["historical_import.py"]
|
||||||
ing -.-> hy[("Hypertable reading")]
|
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 -.-> agg[("Agrégat continu")]
|
||||||
hy -.-> comp["Compression"]
|
hy -.-> comp["Compression"]
|
||||||
hy -.-> ret["Rétention"]
|
hy -.-> ret["Rétention"]
|
||||||
agg -.-> api["API FastAPI"]
|
|
||||||
|
agg -.-> backend["API FastAPI"]
|
||||||
agg -.-> graf["Grafana"]
|
agg -.-> graf["Grafana"]
|
||||||
```
|
```
|
||||||
|
|
||||||
Les lectures de l'API et de Grafana visent l'agrégat continu, pas la table brute : c'est tout
|
Les flèches pleines représentent les traitements actuellement implémentés.
|
||||||
l'intérêt de TimescaleDB, et cela doit rester vrai quand les volumes augmenteront.
|
|
||||||
|
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
|
## Tables d'authentification
|
||||||
|
|
||||||
Statut : `Fait`. Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les
|
Statut : `Fait`.
|
||||||
hypertables ; elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
|
||||||
|
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
|
```mermaid
|
||||||
erDiagram
|
erDiagram
|
||||||
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
||||||
|
APP_USER ||--o{ PASSWORD_RESET_TOKEN : recoit
|
||||||
|
|
||||||
APP_USER {
|
APP_USER {
|
||||||
uuid id PK
|
uuid id PK
|
||||||
string email UK
|
string email UK
|
||||||
@@ -88,6 +121,7 @@ erDiagram
|
|||||||
bool must_change_password
|
bool must_change_password
|
||||||
timestamptz credentials_changed_at
|
timestamptz credentials_changed_at
|
||||||
}
|
}
|
||||||
|
|
||||||
REFRESH_TOKEN {
|
REFRESH_TOKEN {
|
||||||
uuid id PK
|
uuid id PK
|
||||||
uuid family_id
|
uuid family_id
|
||||||
@@ -99,6 +133,7 @@ erDiagram
|
|||||||
text revoked_reason
|
text revoked_reason
|
||||||
uuid replaced_by
|
uuid replaced_by
|
||||||
}
|
}
|
||||||
|
|
||||||
LOGIN_ATTEMPT {
|
LOGIN_ATTEMPT {
|
||||||
bigint id PK
|
bigint id PK
|
||||||
timestamptz occurred_at
|
timestamptz occurred_at
|
||||||
@@ -106,6 +141,7 @@ erDiagram
|
|||||||
inet client_ip
|
inet client_ip
|
||||||
text outcome
|
text outcome
|
||||||
}
|
}
|
||||||
|
|
||||||
AUDIT_LOG {
|
AUDIT_LOG {
|
||||||
bigint id PK
|
bigint id PK
|
||||||
timestamptz occurred_at
|
timestamptz occurred_at
|
||||||
@@ -114,31 +150,47 @@ erDiagram
|
|||||||
text action
|
text action
|
||||||
jsonb detail
|
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
|
- **`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
|
`CURRENT_USER`. Le nom rappelle aussi qu'il s'agit d'un compte applicatif.
|
||||||
au rôle PostgreSQL qui portera le cantonnement de l'ETL.
|
- **`credentials_changed_at`, une seule colonne**, couvre notamment le changement de mot de passe,
|
||||||
- **`credentials_changed_at`, une seule colonne**, couvre le changement de mot de passe, le
|
le changement de rôle et la désactivation.
|
||||||
changement de rôle et la désactivation. Un compteur de version ne dirait rien à un humain qui
|
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation.
|
||||||
lit un audit.
|
- **`audit_log.actor_id` n'a aucune clé étrangère** afin de conserver les informations d'audit
|
||||||
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation. S'il
|
même si l'entité d'origine évolue.
|
||||||
glissait, la promesse de sept jours serait fictive et une session active ne finirait jamais.
|
- `password_reset_token` ne stocke que l'empreinte du jeton et jamais sa valeur directement.
|
||||||
- **`audit_log.actor_id` n'a aucune clé étrangère**, et `actor_email` comme `actor_role` sont
|
- `password_reset_attempt` est séparée de `audit_log`, car son volume peut être piloté
|
||||||
dénormalisés. Une contrainte `ON DELETE SET NULL` déclencherait un `UPDATE` que le déclencheur
|
par des demandes externes répétées.
|
||||||
d'ajout seul refuserait. Voir l'[ADR 0004](../adr/0004-journal-d-audit-en-ajout-seul.md).
|
|
||||||
|
|
||||||
`audit_log` porte deux déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`. Elle n'est
|
`audit_log` porte des déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`.
|
||||||
donc **pas** une hypertable : une politique de rétention émettrait des `DELETE` qu'ils
|
Elle n'est donc **pas** une hypertable.
|
||||||
refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est
|
|
||||||
piloté par l'attaquant.
|
|
||||||
|
|
||||||
## Gabarit de révision créant une hypertable
|
## 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
|
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision.
|
||||||
`e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
|
||||||
|
La révision `e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
def upgrade() -> None:
|
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
|
## Questions ouvertes
|
||||||
|
|
||||||
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma est livré : ce qui suit
|
Elles portent maintenant principalement sur l'exploitation du schéma :
|
||||||
porte sur son exploitation, plus sur sa forme.
|
|
||||||
|
|
||||||
- **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure.
|
- **Quelle granularité** conserver à long terme à l'ingestion : seconde, minute ou quart d'heure.
|
||||||
- **Quels agrégats continus**, et sur quelles fenêtres.
|
- **Quels agrégats continus** créer et sur quelles fenêtres.
|
||||||
- **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse.
|
- **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.
|
- **Multi-tenant ou non** : un site appartient-il à un client et faut-il cloisonner les lectures.
|
||||||
|
|
||||||
## Modélisation détaillée des données
|
## Modélisation détaillée des données
|
||||||
|
|
||||||
Cette modélisation prend en compte les fichiers CSV historiques,
|
Cette modélisation prend en compte :
|
||||||
leurs métadonnées JSON et les données de l’API Mock.
|
|
||||||
Elle comprend six tables, depuis le stockage des mesures
|
- les fichiers CSV historiques ;
|
||||||
jusqu’aux recommandations proposées à l’utilisateur.
|
- 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
|
### Schéma de données
|
||||||
|
|
||||||
Le diagramme ci-dessous présente les tables et leurs relations.
|
Le diagramme ci-dessous présente les tables et leurs relations.
|
||||||
|
|
||||||
La révision `e6d2026091501` les crée.
|
La révision `e6d2026091501` les crée.
|
||||||
|
|
||||||

|

|
||||||
@@ -208,27 +264,20 @@ La révision `e6d2026091501` les crée.
|
|||||||
|
|
||||||
### Description des tables
|
### Description des tables
|
||||||
|
|
||||||
Chaque table remplit un rôle précis dans le traitement et l’exploitation
|
Chaque table remplit un rôle précis dans le traitement et l'exploitation des données.
|
||||||
des données.
|
|
||||||
|
|
||||||
| Table | Rôle | Origine des informations |
|
| 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` |
|
| `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` |
|
| `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 |
|
| `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
|
Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`.
|
||||||
dans `dataset.metadata`. Elles servent à l’analyse des données
|
|
||||||
et ne sont pas considérées comme des alertes actuelles.
|
|
||||||
|
|
||||||
Les lignes de `recommendation` sont écrites par le moteur de règles du backend
|
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
|
||||||
(`app/services/recommendation_rules.py`), déclenché par `POST /api/v1/recommendations/generate`
|
|
||||||
ou par `make recommendations`, à partir des alertes déjà en base. Le couple
|
|
||||||
`(alert_id, rule_reference)` est unique : rejouer le moteur sur les mêmes alertes n'ajoute aucune
|
|
||||||
ligne.
|
|
||||||
|
|
||||||
### Relations entre les tables
|
### Relations entre les tables
|
||||||
|
|
||||||
@@ -238,15 +287,21 @@ ligne.
|
|||||||
- Une alerte peut être associée à une prévision du même site.
|
- Une alerte peut être associée à une prévision du même site.
|
||||||
- Une alerte peut donner lieu à plusieurs recommandations.
|
- 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
|
```text
|
||||||
Dataset CSV + métadonnées JSON
|
Dataset CSV + métadonnées JSON
|
||||||
@@ -277,17 +332,26 @@ Dataset CSV + métadonnées JSON
|
|||||||
|
|
||||||
Le pipeline est développé en Python.
|
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é.
|
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.
|
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 :
|
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 ;
|
- 122 647 mesures ;
|
||||||
- 0 doublon détecté dans le dataset source.
|
- 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
|
## 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 :
|
Le pipeline assure :
|
||||||
|
|
||||||
@@ -12,12 +13,15 @@ Le pipeline assure :
|
|||||||
- la validation de leur structure et de leur cohérence ;
|
- la validation de leur structure et de leur cohérence ;
|
||||||
- la normalisation des données nécessaires au stockage ;
|
- la normalisation des données nécessaires au stockage ;
|
||||||
- le suivi de la qualité des données ;
|
- 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 ;
|
- 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.
|
- l'idempotence du chargement afin d'éviter la création de doublons.
|
||||||
|
|
||||||
## Données sources
|
## Données sources
|
||||||
|
|
||||||
|
### Dataset historique
|
||||||
|
|
||||||
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
||||||
|
|
||||||
Il contient les deux fichiers suivants :
|
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.
|
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
|
```text
|
||||||
data/raw/
|
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.
|
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
|
## Technologies utilisées
|
||||||
|
|
||||||
| Technologie | Utilisation |
|
| Technologie | Utilisation |
|
||||||
|---|---|
|
|---|---|
|
||||||
| Python | Développement du pipeline ETL |
|
| Python | Développement du pipeline ETL |
|
||||||
| Pandas | Lecture, validation et transformation des données |
|
| Pandas | Lecture, validation et transformation du dataset historique |
|
||||||
| JSON | Lecture des métadonnées du dataset |
|
| JSON | Lecture des métadonnées et conservation des données sources |
|
||||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset |
|
| 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 |
|
| SQLAlchemy Async | Connexion et chargement asynchrone en base |
|
||||||
| PostgreSQL | Stockage relationnel |
|
| PostgreSQL | Stockage relationnel |
|
||||||
| TimescaleDB | Stockage des séries temporelles énergétiques |
|
| 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 |
|
| Alembic | Gestion des migrations du schéma |
|
||||||
| uv | Gestion et exécution de l'environnement Python |
|
| uv | Gestion et exécution de l'environnement Python |
|
||||||
| Ruff | Contrôle de la qualité du code |
|
| Ruff | Contrôle de la qualité du code |
|
||||||
|
| mypy | Vérification du typage |
|
||||||
| Pytest | Tests automatisés |
|
| 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
|
```text
|
||||||
apps/backend/app/etl/historical_import.py
|
apps/backend/app/etl/historical_import.py
|
||||||
@@ -207,7 +226,7 @@ Valeurs manquantes identifiées :
|
|||||||
| `humidity_percent` | 3 423 |
|
| `humidity_percent` | 3 423 |
|
||||||
| `solar_irradiance_wm2` | 3 964 |
|
| `solar_irradiance_wm2` | 3 964 |
|
||||||
|
|
||||||
## Exécution en dry-run
|
## Exécution historique en dry-run
|
||||||
|
|
||||||
Depuis le dossier :
|
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.
|
Aucune donnée n'est écrite dans la base pendant cette exécution.
|
||||||
|
|
||||||
## Chargement réel
|
## Chargement historique réel
|
||||||
|
|
||||||
Depuis `apps/backend/` :
|
Depuis `apps/backend/` :
|
||||||
|
|
||||||
@@ -249,7 +268,7 @@ Chargement : 2000/122647
|
|||||||
Chargement : 122647/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é :
|
Après le chargement initial, les contrôles en base ont confirmé :
|
||||||
|
|
||||||
@@ -266,7 +285,7 @@ Le premier import a créé :
|
|||||||
nouvelles lectures : 122647
|
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.
|
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é.
|
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 :
|
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;"
|
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
|
```text
|
||||||
csv | 122647
|
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
|
```text
|
||||||
apps/backend/tests/etl/
|
apps/backend/tests/etl/
|
||||||
```
|
```
|
||||||
|
|
||||||
Ils couvrent notamment :
|
Les tests de l'import historique couvrent notamment :
|
||||||
|
|
||||||
- la validation du dataset ;
|
- la validation du dataset ;
|
||||||
- les colonnes obligatoires ;
|
- les colonnes obligatoires ;
|
||||||
@@ -328,22 +539,96 @@ Ils couvrent notamment :
|
|||||||
- la construction des mesures destinées à la BDD ;
|
- la construction des mesures destinées à la BDD ;
|
||||||
- le respect des contraintes du modèle de données.
|
- 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 :
|
Exécuter les tests ETL :
|
||||||
|
|
||||||
```powershell
|
```powershell
|
||||||
uv run pytest tests\etl -v
|
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 :
|
Contrôler la qualité du code :
|
||||||
|
|
||||||
```powershell
|
```powershell
|
||||||
uv run ruff check app\etl tests\etl
|
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