Compare commits

..
51 changed files with 609 additions and 3563 deletions
+325
View File
@@ -0,0 +1,325 @@
name: DAST
# Scan dynamique OWASP ZAP de l'API (issue #41). Il attaque une API qui tourne : le job démarre
# la base et le backend sur le runner, sème un site et quelques relevés (sans ça le scan ne
# frappe que des gestionnaires d'erreur), crée un compte `lecteur` jetable
# (scripts/dast-token.sh), puis lance ZAP sur le contrat OpenAPI avec le jeton de ce compte.
#
# Non bloquant pour l'instant sur les alertes (`continue-on-error` sur la seule étape du scan) :
# le volume d'un premier passage trié est inconnu. Deux étapes suivantes, elles, bloquent si le
# scan n'a rien testé (import du contrat, absence de toute réponse de succès) : un job vert doit
# vouloir dire qu'un scan a eu lieu.
#
# Piège : ce scan tape la configuration par défaut du backend (`APP_ENV=local`, pas de TLS, pas
# de reverse proxy). Il ne dit rien des en-têtes ni du TLS posés par le proxy en production, et
# remontera des alertes (HSTS absent...) qui n'existent pas derrière lui.
on:
workflow_dispatch:
schedule:
# Un scan actif est long : hebdomadaire plutôt qu'à chaque PR.
- cron: "0 3 * * 1"
pull_request:
# Ne se lance sur une PR que si le scan lui-même change.
paths:
- ".github/workflows/dast.yml"
- "scripts/dast-token.sh"
permissions:
contents: read
concurrency:
group: dast-${{ github.ref }}
cancel-in-progress: true
jobs:
zap:
name: Scan OWASP ZAP de l'API
runs-on: ubuntu-latest
# Généreux face aux ~2 minutes observées de bout en bout : le vrai plafond est
# `scanner.maxScanDurationInMins` (étape Scan ZAP), sous le TTL du jeton. Une annulation par
# ce timeout-ci n'exécute pas les étapes `always()` : mieux vaut ne jamais l'atteindre.
timeout-minutes: 30
# Même image que docker-compose.yml : la première migration refuse de s'appliquer sans
# l'extension TimescaleDB (cf. backend.yml).
services:
db:
image: timescale/timescaledb-ha:pg17
env:
POSTGRES_USER: enervision
POSTGRES_PASSWORD: change_me
POSTGRES_DB: enervision_dast
ports:
- "5433:5432"
options: >-
--health-cmd "pg_isready -U enervision -d enervision_dast"
--health-interval 10s
--health-timeout 5s
--health-retries 12
--health-start-period 40s
env:
# Base jetable : ZAP y écrira et le script y crée deux comptes.
DATABASE_URL: postgresql+asyncpg://enervision:change_me@localhost:5433/enervision_dast
APP_SECRET_KEY: secret-de-scan-assez-long-pour-le-validateur
APP_ENV: local
# Le jeton du lecteur doit survivre à toute la durée du scan (15 minutes par défaut).
# 3600 est le plafond accepté par la configuration ; `scanner.maxScanDurationInMins`
# (étape Scan ZAP) reste très en dessous, marge comprise pour les étapes qui l'entourent.
APP_ACCESS_TOKEN_TTL_SECONDS: "3600"
PGPASSWORD: change_me
steps:
- name: Récupère le dépôt
uses: actions/checkout@v7
- name: Installe uv
# Épinglé sur le commit du tag v7 (règle Sonar githubactions:S7637 : dépendance tierce,
# contrairement à actions/checkout ou actions/upload-artifact, premières parties).
uses: astral-sh/setup-uv@37802adc94f370d6bfd71619e3f0bf239e1f3b78 # v7
with:
enable-cache: true
cache-dependency-glob: apps/backend/uv.lock
# `prune-cache` vaut `true` par défaut (encore sur ce commit) : l'étape de post-job
# « Pruning cache » est restée bloquée 5 minutes avant d'échouer (exit code 2) sur un
# run où les 16 étapes précédentes passaient, sans lien avec le scan. Le prune n'est
# qu'une optimisation de taille de cache entre deux runs, pas une garantie : le
# désactiver retire le blocage sans rien changer au comportement du job.
prune-cache: false
- name: Installe l'interpréteur déclaré par .python-version
run: uv python install
working-directory: apps/backend
# `--no-build` : aucune dépendance n'est construite depuis ses sources, donc aucun script de
# build exécuté (règle Sonar S8541). Le projet lui-même n'est pas installé : il tourne depuis
# `apps/backend`, comme dans son Dockerfile. Les `uv run` suivants portent `--frozen
# --no-sync` pour ne rien résoudre ni reconstruire (règle S8544).
- name: Synchronise les dépendances sans dévier du verrou
run: uv sync --frozen --no-dev --no-install-project --no-build
working-directory: apps/backend
- name: Active TimescaleDB sur la base du scan
run: psql -h localhost -p 5433 -U enervision -d enervision_dast -c "CREATE EXTENSION IF NOT EXISTS timescaledb"
- name: Applique les migrations
run: uv run --frozen --no-sync --no-build alembic upgrade head
working-directory: apps/backend
# Sans données, `GET /sites` rend `[]`, chaque `/{site_id}` rend 404 et le scan actif ne
# frappe que des gestionnaires d'erreur plutôt que la logique métier. `db/seeds/` est vide
# (pas encore d'outillage de jeu de données pour la CI) : un site et deux relevés à la main,
# juste assez pour que les routes de lecture aient quelque chose à rendre.
- name: Insère un site et des relevés minimaux pour le scan
run: |
psql -h localhost -p 5433 -U enervision -d enervision_dast <<'SQL'
INSERT INTO site (site_id, site_name, site_type, location, capacity_kw, status)
VALUES ('dast-site', 'Site du scan DAST', 'bureau', 'CI', 50, 'actif')
ON CONFLICT (site_id) DO NOTHING;
INSERT INTO reading (site_id, timestamp, source, consumption_kw, consumption_kwh, is_working_hours, data_quality, raw_data)
VALUES
('dast-site', now() - interval '2 hours', 'api_current', 12.5, 12.5, true, 'good', '{}'),
('dast-site', now() - interval '1 hour', 'api_current', 13.0, 13.0, true, 'good', '{}')
ON CONFLICT DO NOTHING;
SQL
- name: Démarre l'API
run: |
nohup uv run --frozen --no-sync --no-build uvicorn app.main:create_app --factory \
--host 0.0.0.0 --port 8000 > "$RUNNER_TEMP/api.log" 2>&1 &
for _ in $(seq 1 30); do
curl -fsS http://localhost:8000/api/v1/health/ready >/dev/null 2>&1 && exit 0
sleep 2
done
echo "L'API ne répond pas sur /health/ready" >&2
cat "$RUNNER_TEMP/api.log" >&2
exit 1
working-directory: apps/backend
- name: Crée le compte lecteur du scan
id: jeton
run: |
jeton="$(../../scripts/dast-token.sh)"
echo "::add-mask::$jeton"
echo "jeton=$jeton" >> "$GITHUB_OUTPUT"
working-directory: apps/backend
# Étape distincte du scan lui-même, et sans `continue-on-error` : un `curl` qui échoue ici
# (API tombée juste après la sonde de readiness, par exemple) doit rester un échec visible,
# pas se travestir en « ZAP n'a importé aucune URL » à l'étape de garde suivante.
- name: Prépare le contrat pour ZAP
run: |
mkdir -p zap-out zap-logs
curl -fsS http://localhost:8000/openapi.json -o zap-out/openapi.json
# Le dossier passe à l'uid 1000 (utilisateur du conteneur ZAP) : le runner n'y écrit
# plus après ce chown, d'où `zap-logs/` (uid du runner) pour les journaux ci-dessous.
# Pas de `chmod 777` (règle Sonar S2612).
sudo chown -R 1000:1000 zap-out
# `--network host` : ZAP atteint l'API sur le localhost du runner.
#
# Piège vécu : la clé du nom d'en-tête est `matchstr`, pas `matchstring`. ZAP accepte
# n'importe quelle clé `-config` sans erreur ; avec la mauvaise, il ajoutait à TOUTES les
# requêtes un en-tête au nom vide (`: Bearer <jeton>`), qu'uvicorn refuse par un 400
# (« Invalid HTTP request received »), y compris sur les routes publiques.
#
# Le jeton ne passe ni par `${{ }}` dans ce script (il finirait en clair dans le fichier de
# commande que GitHub écrit sur le disque du runner pour toute la durée de l'étape), ni par
# l'argv de `docker run` (visible par `ps aux` et par `docker inspect zap` tant que le
# conteneur existe) : il est écrit dans un fichier de configuration ZAP séparé, monté en
# lecture seule hors de `/zap/wrk` pour ne jamais atterrir dans l'artefact publié.
#
# Les routes d'authentification qui changent l'état du compte du scan sont exclues : un
# scan actif y déclencherait la limitation de débit du login, la réinitialisation de mots de
# passe et la fermeture des sessions, sans rien apprendre de plus.
#
# `scanner.maxScanDurationInMins`/`maxRuleDurationInMins` bornent le scan actif, que `-T` ne
# couvre pas (il ne borne que le démarrage et le scan passif) : sans ça, une règle qui
# traîne peut dépasser le TTL du jeton (401 muets en fin de scan) ou le timeout du job (qui
# annule sans exécuter les étapes `always()`, rapport et journaux perdus).
- name: Scan ZAP
id: zap
continue-on-error: true
env:
JETON: ${{ steps.jeton.outputs.jeton }}
run: |
set -o pipefail
printf 'replacer.full_list(0).description=auth\nreplacer.full_list(0).enabled=true\nreplacer.full_list(0).matchtype=REQ_HEADER\nreplacer.full_list(0).matchstr=Authorization\nreplacer.full_list(0).regex=false\nreplacer.full_list(0).replacement=Bearer %s\n' "$JETON" > "$RUNNER_TEMP/zap-auth.conf"
# Piège vécu : `chmod 600` seul rend le fichier illisible pour le conteneur, qui lit un
# montage bind avec son propre uid (1000), distinct de celui du runner qui l'a écrit.
# ZAP échoue alors dès le lancement (« File not readable: /zap/auth.conf »), et
# `zap-api-scan.py` attend `-T` minutes complètes avant d'abandonner : dix minutes qui
# ressemblent à un scan actif, pour un daemon mort depuis le début.
#
# Piège vécu (numéro deux) : une fois le fichier passé à l'uid 1000 par `sudo chown`,
# l'utilisateur du runner n'en est plus propriétaire et un `chmod` sans `sudo` échoue
# (« Operation not permitted »). Avec le `-e` implicite de bash sur les étapes GitHub
# Actions, cette erreur arrêtait toute l'étape avant même `docker run` : scan « réussi »
# en une fraction de seconde, sans le moindre journal ni rapport produit.
sudo chown 1000:1000 "$RUNNER_TEMP/zap-auth.conf"
sudo chmod 644 "$RUNNER_TEMP/zap-auth.conf"
docker run --name zap --network host \
-v "$PWD/zap-out:/zap/wrk:rw" \
-v "$RUNNER_TEMP/zap-auth.conf:/zap/auth.conf:ro" \
ghcr.io/zaproxy/zaproxy:stable zap-api-scan.py \
-t /zap/wrk/openapi.json -f openapi -O http://localhost:8000 \
-T 10 \
-r zap-report.html -J zap-report.json -w zap-report.md \
-z "-configfile /zap/auth.conf \
-config globalexcludeurl.url_list.url(0).description=auth-etat \
-config globalexcludeurl.url_list.url(0).enabled=true \
-config globalexcludeurl.url_list.url(0).regex='.*/api/v1/auth/(login|password|logout-all|forgot-password|reset-password).*' \
-config scanner.maxScanDurationInMins=15 \
-config scanner.maxRuleDurationInMins=5" \
2>&1 | tee "$RUNNER_TEMP/zap-stdout.log"
- name: Récupère les journaux de ZAP
if: always()
run: |
mkdir -p zap-logs
# ZAP journalise la valeur de chaque `-config`/`-configfile` chargé, y compris le jeton,
# à un niveau visible sans `-d` : les copies publiées en artefact sont donc caviardées,
# même si `::add-mask::` (posé à la création du jeton) protège déjà le journal du job.
masque() { sed -E 's/(Bearer )[A-Za-z0-9._-]+/\1[MASQUE]/Ig'; }
[ -f "$RUNNER_TEMP/zap-stdout.log" ] && masque < "$RUNNER_TEMP/zap-stdout.log" > zap-logs/zap-stdout.log
docker cp zap:/home/zap/.ZAP/zap.log "$RUNNER_TEMP/zap-internal.log" 2>/dev/null || true
[ -f "$RUNNER_TEMP/zap-internal.log" ] && masque < "$RUNNER_TEMP/zap-internal.log" > zap-logs/zap.log
[ -f "$RUNNER_TEMP/api.log" ] && masque < "$RUNNER_TEMP/api.log" > zap-logs/api.log
rm -f "$RUNNER_TEMP/zap-auth.conf"
docker rm -f zap >/dev/null 2>&1 || true
# `continue-on-error` sur le scan ne doit pas faire passer pour vert un scan qui n'a rien
# testé. Constaté une première fois : 2 URL importées sur 26 opérations, ZAP n'avait envoyé
# que des requêtes vouées au 404. Le seuil est dérivé du contrat plutôt que d'un nombre fixe
# : un contrat qui grossit ne doit pas rendre la garde plus permissive qu'elle ne l'était.
- name: Vérifie que le contrat a bien été importé
run: |
attendu="$(python3 -c "
import json
d = json.load(open('zap-out/openapi.json'))
methodes = ('get', 'post', 'put', 'patch', 'delete', 'head', 'options')
print(sum(1 for chemin in d['paths'].values() for m in chemin if m in methodes))
")"
minimum=$((attendu * 80 / 100))
importees="$(sed -n 's/.*Number of Imported URLs: \([0-9]*\).*/\1/p' "$RUNNER_TEMP/zap-stdout.log" | tail -1)"
echo "URL importées depuis le contrat OpenAPI : ${importees:-aucune} (contrat : $attendu opérations, minimum accepté : $minimum)"
if [ "${importees:-0}" -lt "$minimum" ]; then
echo "::error::ZAP n'a importé que ${importees:-0} URL sur $attendu opérations du contrat OpenAPI (minimum attendu : $minimum, soit 80%). Le scan n'a pas testé l'API, voir zap-logs/zap.log dans l'artefact zap-report."
exit 1
fi
# Deuxième garde-fou : le contrat peut être importé et ZAP n'obtenir que des erreurs
# (constaté : base sans données, toutes les routes de site répondaient 404).
#
# Piège de conception, trouvé en répétant ce job en local avant de l'écrire ici : borner le
# pourcentage de 4xx ne marche pas. Un scan actif fuzze délibérément un grand nombre
# d'entrées invalides (identifiants inventés, méthodes non supportées...), donc même un scan
# sain, contre l'API seedée juste au-dessus, reste à 98% de 4xx avec seulement 1% de 2xx :
# c'est la forme normale d'un scan actif, pas un signe d'échec. Le signal qui distingue
# vraiment un scan cassé (0% de 2xx, `insight.code.2xx` absent du rapport dans le premier
# incident) d'un scan sain (2xx non nul, aussi faible soit-il) est donc l'absence de succès,
# pas la part d'échecs. Dérivé de `zap-report.json` (champ structuré `insights[]`) plutôt
# que du texte libre du rapport Markdown, qui aurait le même défaut de conception en plus
# d'être fragile au format.
- name: Vérifie que le scan a obtenu au moins une réponse de succès
run: |
python3 - <<'PY'
import json
import sys
try:
rapport = json.load(open("zap-out/zap-report.json"))
except FileNotFoundError:
print("::error::Aucun rapport ZAP produit : le scan n'a rien testé.")
sys.exit(1)
pourcentage_2xx = 0.0
for insight in rapport.get("insights", []):
if insight.get("key") == "insight.code.2xx":
pourcentage_2xx = float(insight.get("statistic", 0))
break
print(f"Pourcentage de réponses 2xx : {pourcentage_2xx}%")
if pourcentage_2xx <= 0:
print(
"::error::Aucune réponse 2xx (succès) reçue : le scan n'a atteint aucune route "
"réelle de l'API. Voir zap-logs/api.log et zap-logs/zap.log dans l'artefact "
"zap-report."
)
sys.exit(1)
PY
# Uniquement la synthèse (jusqu'à « Alert Detail » exclu) : `$GITHUB_STEP_SUMMARY` est
# limité à 1 Mio, et cette étape tourne sous `always()` - son échec ferait échouer le job
# après le passage des deux garde-fous, pour une simple raison de mise en forme. Le rapport
# complet reste dans l'artefact `zap-report`.
- name: Publie le résumé
if: always()
run: |
if [ -f zap-out/zap-report.md ]; then
awk '/^## Alert Detail/{exit} {print}' zap-out/zap-report.md >> "$GITHUB_STEP_SUMMARY"
echo "" >> "$GITHUB_STEP_SUMMARY"
echo "Rapport complet (HTML/JSON/Markdown) dans l'artefact \`zap-report\`." >> "$GITHUB_STEP_SUMMARY"
else
echo "Aucun rapport ZAP produit, voir le journal du job." >> "$GITHUB_STEP_SUMMARY"
fi
- name: Publie les rapports
if: always()
uses: actions/upload-artifact@v7
with:
name: zap-report
path: |
zap-out/
zap-logs/
if-no-files-found: warn
# Diagnostic de dernier recours : les journaux de l'API sont déjà dans l'artefact
# (zap-logs/api.log) via l'étape « Récupère les journaux de ZAP » (always()), mais les
# afficher directement dans le journal du job évite d'avoir à le télécharger pour un échec
# évident (l'API n'a jamais démarré, par exemple).
- name: Journal de l'API en cas d'échec
if: failure() || steps.zap.outcome == 'failure'
run: cat "$RUNNER_TEMP/api.log" || true
+2 -100
View File
@@ -8,28 +8,10 @@ on:
paths: paths:
- "ml/**" - "ml/**"
- ".github/workflows/ml.yml" - ".github/workflows/ml.yml"
# Piege : `backend.yml` ne joue jamais `-m chaine`, ce job est le seul. Le test de chaine
# traverse tout `apps/backend/app` jusqu'a GET /predictions : restreindre le filtre aux
# migrations et aux modeles le laisserait muet sur la PR meme qui le casse. Meme raisonnement
# que le filtre d'airflow.yml, qui inclut deja des chemins de ml/ et de apps/backend/.
- "apps/backend/alembic/**"
- "apps/backend/app/**"
- "apps/backend/tests/test_chaine_ml_api.py"
- "apps/backend/pyproject.toml"
- "apps/backend/uv.lock"
pull_request: pull_request:
paths: paths:
- "ml/**" - "ml/**"
- ".github/workflows/ml.yml" - ".github/workflows/ml.yml"
# Piege : `backend.yml` ne joue jamais `-m chaine`, ce job est le seul. Le test de chaine
# traverse tout `apps/backend/app` jusqu'a GET /predictions : restreindre le filtre aux
# migrations et aux modeles le laisserait muet sur la PR meme qui le casse. Meme raisonnement
# que le filtre d'airflow.yml, qui inclut deja des chemins de ml/ et de apps/backend/.
- "apps/backend/alembic/**"
- "apps/backend/app/**"
- "apps/backend/tests/test_chaine_ml_api.py"
- "apps/backend/pyproject.toml"
- "apps/backend/uv.lock"
permissions: permissions:
contents: read contents: read
@@ -71,91 +53,11 @@ jobs:
- name: Typage - name: Typage
run: uv run mypy enervision_ml tests run: uv run mypy enervision_ml tests
# Les tests exigeant une base portent le marqueur `integration`, ecarte par defaut et # Aucun test ne touche PostgreSQL ni MLflow distant : tout tourne sur donnees
# joue par le job `integration` ci-dessous. # synthetiques ou un magasin SQLite local jetable (cf. ml/tests/test_train.py).
- name: Tests - name: Tests
run: uv run pytest run: uv run pytest
# Le seul job du depot qui dispose a la fois des deux environnements uv et d'une base. Piege :
# le schema de la base ML est celui du backend (apps/backend/alembic, proprietaire du schema).
# Le reconstruire ici a la main rendrait ce job vert sur une base qui n'est pas la notre.
integration:
name: ML - DB et chaîne ML - DB - API
runs-on: ubuntu-latest
services:
db:
image: timescale/timescaledb-ha:pg17
env:
POSTGRES_USER: enervision
POSTGRES_PASSWORD: change_me
POSTGRES_DB: enervision_test
ports:
- "5433:5432"
options: >-
--health-cmd "pg_isready -U enervision -d enervision_test"
--health-interval 10s
--health-timeout 5s
--health-retries 12
--health-start-period 40s
env:
# Deux variables, deux dialectes : Alembic et l'API parlent asyncpg, le pipeline ML parle
# psycopg en synchrone. Cf. docs/ML-START.md, section 1.
DATABASE_URL: postgresql+asyncpg://enervision:change_me@localhost:5433/enervision_test
ML_DATABASE_URL: postgresql+psycopg://enervision:change_me@localhost:5433/enervision_test
APP_SECRET_KEY: secret-de-test-assez-long-pour-le-validateur
PGPASSWORD: change_me
steps:
- name: Récupère le dépôt
uses: actions/checkout@v7
- name: Installe uv
uses: astral-sh/setup-uv@v7
with:
enable-cache: true
cache-dependency-glob: |
ml/uv.lock
apps/backend/uv.lock
- name: Installe l'interpréteur déclaré par .python-version
working-directory: ml
run: uv python install
- name: Synchronise le pipeline ML sans dévier du verrou
working-directory: ml
run: uv sync --all-groups --frozen
# Le backend est installé ici parce qu'il porte les migrations, seule source du schéma, et
# le test de chaîne, qui interroge l'API.
- name: Synchronise le backend sans dévier du verrou
working-directory: apps/backend
run: uv sync --all-groups --frozen
# db/init/110-test-database.sql n'est pas monté ici, et sans l'extension la première
# révision Alembic refuse de s'appliquer.
- name: Active TimescaleDB sur la base de test
run: psql -h localhost -p 5433 -U enervision -d enervision_test -c "CREATE EXTENSION IF NOT EXISTS timescaledb"
- name: Applique les migrations du backend, propriétaire du schéma
working-directory: apps/backend
run: uv run alembic upgrade head
# `-m` en ligne de commande écrase celui d'addopts. Couverture désactivée : ce job ne joue
# qu'une partie de la suite, son taux n'aurait pas de sens (même raison que backend.yml).
- name: Tests ML exigeant une base
working-directory: ml
run: uv run pytest -m integration --no-cov
# Lance les vrais binaires enervision_ml.train et .score en sous-processus, comme les DAGs
# ml_train et ml_score, puis relit le résultat par GET /api/v1/predictions.
- name: Chaîne complète ML vers DB vers API
working-directory: apps/backend
env:
ML_PYTHON: ${{ github.workspace }}/ml/.venv/bin/python
run: uv run pytest -m chaine --no-cov
sast: sast:
name: Analyse statique de sécurité name: Analyse statique de sécurité
runs-on: ubuntu-latest runs-on: ubuntu-latest
+2 -23
View File
@@ -25,13 +25,6 @@ MAILPIT_UI_PORT := $(or $(strip $(call env-val,MAILPIT_UI_PORT)),8025)
ML_DATABASE_URL ?= postgresql+psycopg://$(PG_USER):$(PG_PASSWORD)@localhost:$(PG_PORT)/$(PG_DB) ML_DATABASE_URL ?= postgresql+psycopg://$(PG_USER):$(PG_PASSWORD)@localhost:$(PG_PORT)/$(PG_DB)
export ML_DATABASE_URL export ML_DATABASE_URL
# Piege : la base des tests d'integration n'est pas la base de developpement. Ces tests ecrivent
# et suppriment des lignes, et leurs fixtures refusent de demarrer ailleurs que sur
# `enervision_test` (garde sur le nom, cf. ml/tests/conftest.py).
PG_TEST_DB ?= enervision_test
TEST_DATABASE_URL ?= postgresql+asyncpg://$(PG_USER):$(PG_PASSWORD)@localhost:$(PG_PORT)/$(PG_TEST_DB)
ML_TEST_DATABASE_URL ?= postgresql+psycopg://$(PG_USER):$(PG_PASSWORD)@localhost:$(PG_PORT)/$(PG_TEST_DB)
# Le jeu historique s'arrete au 31/12/2024 : score et detection ancres a l'horloge reelle ne # Le jeu historique s'arrete au 31/12/2024 : score et detection ancres a l'horloge reelle ne
# verraient qu'un parc muet depuis des mois. Cf. `--now` de enervision_ml.score. # verraient qu'un parc muet depuis des mois. Cf. `--now` de enervision_ml.score.
DEMO_NOW ?= 2024-12-31T00:00:00Z DEMO_NOW ?= 2024-12-31T00:00:00Z
@@ -39,10 +32,9 @@ DEMO_NOW ?= 2024-12-31T00:00:00Z
.DEFAULT_GOAL := help .DEFAULT_GOAL := help
.PHONY: help install install-backend install-frontend install-ml install-airflow \ .PHONY: help install install-backend install-frontend install-ml install-airflow \
dev dev-backend dev-frontend \ dev dev-backend dev-frontend \
lint format typecheck test test-cov test-integration ml-test-integration \ lint format typecheck test test-cov test-integration check \
test-chaine check \
openapi docker-build db-up db-down db-reset db-logs db-psql db-wait db-ensure-airflow \ openapi docker-build db-up db-down db-reset db-logs db-psql db-wait db-ensure-airflow \
migrate migrate-test bootstrap-admin services-up demo-data demo-data-force \ migrate bootstrap-admin services-up demo-data demo-data-force \
ml-lint ml-typecheck ml-test ml-check ml-train ml-score detect-alerts recommendations \ ml-lint ml-typecheck ml-test ml-check ml-train ml-score detect-alerts recommendations \
airflow-lint airflow-test airflow-check airflow-up airflow-down airflow-logs \ airflow-lint airflow-test airflow-check airflow-up airflow-down airflow-logs \
tls-selfsigned tls-acme tls-renew stack-up stack-down stack-logs tls-selfsigned tls-acme tls-renew stack-up stack-down stack-logs
@@ -120,16 +112,6 @@ ml-test: ## Exécute les tests du pipeline ML (donnees synthetiques, sans base n
ml-check: ml-lint ml-typecheck ml-test ## Chaîne de vérification complète du pipeline ML ml-check: ml-lint ml-typecheck ml-test ## Chaîne de vérification complète du pipeline ML
# La cible surcharge ML_DATABASE_URL, que ce Makefile exporte vers la base de développement : la
# garde du conftest ferait échouer la cible sans cette surcharge.
ml-test-integration: ML_DATABASE_URL := $(ML_TEST_DATABASE_URL)
ml-test-integration: ## Tests ML exigeant une base migrée. Faire `make db-up migrate-test` avant
cd $(ML) && uv run pytest -m integration --no-cov
test-chaine: ## Chaîne ML -> DB -> API, vrais binaires. Exige les deux environnements uv
cd $(BACKEND) && DATABASE_URL=$(TEST_DATABASE_URL) ML_PYTHON=$(CURDIR)/$(ML)/.venv/bin/python \
uv run pytest -m chaine --no-cov
ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DATABASE_URL ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DATABASE_URL
cd $(ML) && uv run python -m enervision_ml.train $(if $(CSV),--csv $(CSV),) cd $(ML) && uv run python -m enervision_ml.train $(if $(CSV),--csv $(CSV),)
@@ -227,9 +209,6 @@ db-ensure-airflow: ## Crée la base de métadonnées Airflow si le volume pgdata
migrate: ## Applique les migrations Alembic migrate: ## Applique les migrations Alembic
cd $(BACKEND) && uv run alembic upgrade head cd $(BACKEND) && uv run alembic upgrade head
migrate-test: ## Applique les migrations sur enervision_test, la base des tests d'intégration
cd $(BACKEND) && DATABASE_URL=$(TEST_DATABASE_URL) uv run alembic upgrade head
bootstrap-admin: ## Crée le premier administrateur, mot de passe saisi au clavier bootstrap-admin: ## Crée le premier administrateur, mot de passe saisi au clavier
cd $(BACKEND) && uv run python -m app.cli create-admin --email $${EMAIL:?EMAIL=... requis} cd $(BACKEND) && uv run python -m app.cli create-admin --email $${EMAIL:?EMAIL=... requis}
+2 -2
View File
@@ -23,7 +23,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m
| Backend | FastAPI, Python 3.14 | `apps/backend` | En place | | Backend | FastAPI, Python 3.14 | `apps/backend` | En place |
| Frontend | Angular 22, Node 26 | `apps/frontend` | En place | | Frontend | Angular 22, Node 26 | `apps/frontend` | En place |
| Base | PostgreSQL 17 + TimescaleDB | `db` | En place | | Base | PostgreSQL 17 + TimescaleDB | `db` | En place |
| ETL | Apache Airflow | `etl/airflow` | Cinq DAGs | | ETL | Apache Airflow | `etl/airflow` | Quatre DAGs |
| Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise | | Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise |
| Reverse proxy | Nginx, TLS | `infra/proxy` | En place | | Reverse proxy | Nginx, TLS | `infra/proxy` | En place |
| CI/CD | GitHub Actions | `.github/workflows` | En place | | CI/CD | GitHub Actions | `.github/workflows` | En place |
@@ -50,7 +50,7 @@ L'etat detaille de chaque brique et les vues d'architecture sont dans
│ ├── migrations/ Migrations SQL versionnees │ ├── migrations/ Migrations SQL versionnees
│ └── seeds/ Jeux de donnees de reference │ └── seeds/ Jeux de donnees de reference
├── etl/airflow/ ├── etl/airflow/
│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, import, dérive) │ ├── dags/ DAGs d'orchestration (pipeline ML, alertes, import historique)
│ ├── plugins/ Operateurs et hooks maison │ ├── plugins/ Operateurs et hooks maison
│ ├── include/ Requetes SQL et ressources des DAGs │ ├── include/ Requetes SQL et ressources des DAGs
│ └── tests/ Tests d'integrite des DAGs │ └── tests/ Tests d'integrite des DAGs
@@ -1,81 +0,0 @@
"""rapports de derive du modele de prevision
Revision ID: d3f1a2b7c904
Revises: c0adab96238c
Create Date: 2026-09-22 14:40:00.000000
`site_id` est nullable, et c'est le coeur du schema : une ligne par site, plus une ligne
globale tous sites confondus, que `NULL` designe. Un seul site qui derive est invisible dans
une moyenne d'ensemble, et une derive d'ensemble sans rupture par site signale un changement
de modele ou de saison, pas une panne.
L'unicite passe par un index a `coalesce` et non par une `UniqueConstraint` : deux lignes
globales successives ont toutes deux `site_id` a NULL, et NULL n'est egal a aucune valeur, pas
meme a lui-meme. Meme forme que `uq_reading_source`.
Les trois `CHECK` sont portees par la base, comme `ck_prediction_status` : un verdict sans
motif, ou un statut inconnu, ne doit pas dependre de la vigilance de l'appelant.
"""
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql
revision: str = "d3f1a2b7c904"
down_revision: str | Sequence[str] | None = "c0adab96238c"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.create_table(
"drift_report",
sa.Column("drift_report_id", sa.BigInteger(), autoincrement=True, nullable=False),
sa.Column(
"computed_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
sa.Column("site_id", sa.Text(), nullable=True),
sa.Column("window_start", sa.DateTime(timezone=True), nullable=False),
sa.Column("window_end", sa.DateTime(timezone=True), nullable=False),
sa.Column("reference_start", sa.DateTime(timezone=True), nullable=True),
sa.Column("reference_end", sa.DateTime(timezone=True), nullable=True),
sa.Column("n_observations", sa.Integer(), nullable=False),
sa.Column("mae", sa.Double(), nullable=True),
sa.Column("mape", sa.Double(), nullable=True),
sa.Column("bias", sa.Double(), nullable=True),
sa.Column("reference_mae", sa.Double(), nullable=True),
sa.Column("coverage_ratio", sa.Double(), nullable=True),
sa.Column("insufficient_data_ratio", sa.Double(), nullable=True),
sa.Column("model_references", postgresql.ARRAY(sa.Text()), nullable=False),
sa.Column("status", sa.Text(), nullable=False),
sa.Column("reason", sa.Text(), nullable=True),
sa.CheckConstraint(
"status IN ('stable', 'derive', 'indetermine')", name="ck_drift_report_status"
),
sa.CheckConstraint(
"status = 'stable' OR reason IS NOT NULL", name="ck_drift_report_reason"
),
sa.CheckConstraint("n_observations >= 0", name="ck_drift_report_observations"),
sa.ForeignKeyConstraint(
["site_id"], ["site.site_id"], name="fk_drift_report_site", ondelete="RESTRICT"
),
sa.PrimaryKeyConstraint("drift_report_id"),
)
op.create_index(
"ix_drift_report_site_computed", "drift_report", ["site_id", "computed_at"], unique=False
)
op.create_index(
"uq_drift_report_window",
"drift_report",
["window_end", sa.literal_column("coalesce(site_id, '')")],
unique=True,
)
def downgrade() -> None:
op.drop_table("drift_report")
+3 -10
View File
@@ -24,7 +24,6 @@ from app.core.security import decode_access_token as decode_token
from app.db.session import get_session from app.db.session import get_session
from app.repositories.alert import AlertRepository from app.repositories.alert import AlertRepository
from app.repositories.audit_log import AuditLogRepository from app.repositories.audit_log import AuditLogRepository
from app.repositories.drift import DriftRepository
from app.repositories.login_attempt import LoginAttemptRepository from app.repositories.login_attempt import LoginAttemptRepository
from app.repositories.password_reset_attempt import PasswordResetAttemptRepository from app.repositories.password_reset_attempt import PasswordResetAttemptRepository
from app.repositories.password_reset_token import PasswordResetTokenRepository from app.repositories.password_reset_token import PasswordResetTokenRepository
@@ -36,7 +35,6 @@ from app.repositories.site import SiteRepository
from app.repositories.user import UserRepository from app.repositories.user import UserRepository
from app.services.alert import AlertService from app.services.alert import AlertService
from app.services.auth import AuthService, LoginPolicy, PasswordResetPolicy from app.services.auth import AuthService, LoginPolicy, PasswordResetPolicy
from app.services.drift import DriftService
from app.services.prediction import PredictionService from app.services.prediction import PredictionService
from app.services.reading import ReadingService from app.services.reading import ReadingService
from app.services.recommendation import RecommendationService from app.services.recommendation import RecommendationService
@@ -50,7 +48,9 @@ SettingsDep = Annotated[Settings, Depends(get_settings)]
CODE_CHANGEMENT_REQUIS = "password_change_required" CODE_CHANGEMENT_REQUIS = "password_change_required"
_porteur = HTTPBearer(auto_error=False, scheme_name="Jeton d'accès") # Nom ASCII : un outillage tiers (ZAP, cf. .github/workflows/dast.yml) peut mal analyser un nom
# de schéma accentué dans le contrat OpenAPI. Piège vécu, pas anticipé.
_porteur = HTTPBearer(auto_error=False, scheme_name="JetonAcces")
CredentialsDep = Annotated[HTTPAuthorizationCredentials | None, Depends(_porteur)] CredentialsDep = Annotated[HTTPAuthorizationCredentials | None, Depends(_porteur)]
@@ -234,13 +234,6 @@ def get_prediction_service(session: SessionDep) -> PredictionService:
PredictionServiceDep = Annotated[PredictionService, Depends(get_prediction_service)] PredictionServiceDep = Annotated[PredictionService, Depends(get_prediction_service)]
def get_drift_service(session: SessionDep) -> DriftService:
return DriftService(DriftRepository(session))
DriftServiceDep = Annotated[DriftService, Depends(get_drift_service)]
async def get_current_principal( async def get_current_principal(
credentials: CredentialsDep, credentials: CredentialsDep,
session: SessionDep, session: SessionDep,
+3 -20
View File
@@ -90,19 +90,13 @@ TAGS: Final[list[dict[str, Any]]] = [
"de scoring (`ml/`) et simplement lue ici. Accessible à partir du rôle `lecteur`." "de scoring (`ml/`) et simplement lue ici. Accessible à partir du rôle `lecteur`."
), ),
}, },
{
"name": "monitoring",
"description": (
"Surveillance de la dérive du modèle : écart entre les prévisions déjà écrites et "
"les lectures réellement arrivées, par site et tous sites confondus. Réservé à "
"partir du rôle `operateur`, qui agit sur un pipeline dégradé."
),
},
] ]
cookie_de_rafraichissement = APIKeyCookie( cookie_de_rafraichissement = APIKeyCookie(
name=REFRESH_COOKIE_DEFAUT, name=REFRESH_COOKIE_DEFAUT,
scheme_name="Cookie de rafraîchissement", # Nom ASCII : un outillage tiers (ZAP, cf. .github/workflows/dast.yml) peut mal analyser un
# nom de schéma accentué dans le contrat OpenAPI. Piège vécu, pas anticipé.
scheme_name="CookieRafraichissement",
description=( description=(
"Cookie `HttpOnly` posé par `/auth/login` et tourné par `/auth/refresh`. Il prend le " "Cookie `HttpOnly` posé par `/auth/login` et tourné par `/auth/refresh`. Il prend le "
"préfixe `__Secure-` dès que l'API tourne derrière TLS, et n'est émis que vers " "préfixe `__Secure-` dès que l'API tourne derrière TLS, et n'est émis que vers "
@@ -161,17 +155,6 @@ REPONSES_ADMIN: Final[Reponses] = {
}, },
} }
REPONSES_OPERATEUR: Final[Reponses] = {
**REPONSES_AUTHENTIFIEES,
403: {
"model": ErrorResponse,
"description": (
"Droits insuffisants, ou mot de passe provisoire à changer quand `detail` vaut "
"`password_change_required`."
),
},
}
# `lecteur` est le rôle minimum : `require_role` n'y refuse jamais un 403 pour droits # `lecteur` est le rôle minimum : `require_role` n'y refuse jamais un 403 pour droits
# insuffisants, seulement pour le mot de passe provisoire. # insuffisants, seulement pour le mot de passe provisoire.
REPONSES_LECTEUR: Final[Reponses] = { REPONSES_LECTEUR: Final[Reponses] = {
@@ -1,20 +0,0 @@
from fastapi import APIRouter
from app.api.deps import DriftServiceDep, OperateurDep
from app.api.openapi import REPONSE_VALIDATION
from app.schemas.drift import DriftReportResponse
router = APIRouter()
@router.get(
"/drift",
response_model=list[DriftReportResponse],
summary="Dernier rapport de dérive par site, plus la ligne globale",
responses=REPONSE_VALIDATION,
)
async def get_drift(
_: OperateurDep, service: DriftServiceDep, site_id: str | None = None
) -> list[DriftReportResponse]:
rapports = await service.derniers(site_id=site_id)
return [DriftReportResponse.model_validate(rapport) for rapport in rapports]
+1 -10
View File
@@ -1,16 +1,10 @@
from fastapi import APIRouter from fastapi import APIRouter
from app.api.openapi import ( from app.api.openapi import REPONSE_SERVEUR, REPONSES_ADMIN, REPONSES_LECTEUR
REPONSE_SERVEUR,
REPONSES_ADMIN,
REPONSES_LECTEUR,
REPONSES_OPERATEUR,
)
from app.api.v1.endpoints import ( from app.api.v1.endpoints import (
alerts, alerts,
auth, auth,
health, health,
monitoring,
predictions, predictions,
readings, readings,
recommendations, recommendations,
@@ -44,6 +38,3 @@ api_router.include_router(
api_router.include_router( api_router.include_router(
predictions.router, prefix="/predictions", tags=["predictions"], responses=REPONSES_LECTEUR predictions.router, prefix="/predictions", tags=["predictions"], responses=REPONSES_LECTEUR
) )
api_router.include_router(
monitoring.router, prefix="/monitoring", tags=["monitoring"], responses=REPONSES_OPERATEUR
)
+1 -10
View File
@@ -2,15 +2,7 @@
# --autogenerate`, qui générerait alors un drop de sa table. # --autogenerate`, qui générerait alors un drop de sa table.
from app.models.audit_log import AuditLog from app.models.audit_log import AuditLog
from app.models.energy import ( from app.models.energy import Alert, Dataset, Prediction, Reading, Recommendation, Site
Alert,
Dataset,
DriftReport,
Prediction,
Reading,
Recommendation,
Site,
)
from app.models.login_attempt import LoginAttempt from app.models.login_attempt import LoginAttempt
from app.models.password_reset_attempt import PasswordResetAttempt from app.models.password_reset_attempt import PasswordResetAttempt
from app.models.password_reset_token import PasswordResetToken from app.models.password_reset_token import PasswordResetToken
@@ -22,7 +14,6 @@ __all__ = [
"AppUser", "AppUser",
"AuditLog", "AuditLog",
"Dataset", "Dataset",
"DriftReport",
"LoginAttempt", "LoginAttempt",
"PasswordResetAttempt", "PasswordResetAttempt",
"PasswordResetToken", "PasswordResetToken",
-46
View File
@@ -208,49 +208,3 @@ class Recommendation(Base):
explanation: Mapped[str] = mapped_column(Text) explanation: Mapped[str] = mapped_column(Text)
rule_reference: Mapped[str] = mapped_column(Text) rule_reference: Mapped[str] = mapped_column(Text)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
class DriftReport(Base):
__tablename__ = "drift_report"
__table_args__ = (
CheckConstraint(
"status IN ('stable', 'derive', 'indetermine')", name="ck_drift_report_status"
),
CheckConstraint("status = 'stable' OR reason IS NOT NULL", name="ck_drift_report_reason"),
CheckConstraint("n_observations >= 0", name="ck_drift_report_observations"),
Index("ix_drift_report_site_computed", "site_id", "computed_at"),
)
drift_report_id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
computed_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now()
)
# `NULL` porte la ligne globale, tous sites confondus : une derive d'ensemble et la derive
# d'un seul site ne se lisent pas dans le meme chiffre.
site_id: Mapped[str | None] = mapped_column(
Text, ForeignKey("site.site_id", name="fk_drift_report_site", ondelete="RESTRICT")
)
window_start: Mapped[datetime] = mapped_column(DateTime(timezone=True))
window_end: Mapped[datetime] = mapped_column(DateTime(timezone=True))
reference_start: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
reference_end: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
n_observations: Mapped[int] = mapped_column(Integer)
mae: Mapped[float | None] = mapped_column(Double)
mape: Mapped[float | None] = mapped_column(Double)
bias: Mapped[float | None] = mapped_column(Double)
reference_mae: Mapped[float | None] = mapped_column(Double)
coverage_ratio: Mapped[float | None] = mapped_column(Double)
insufficient_data_ratio: Mapped[float | None] = mapped_column(Double)
model_references: Mapped[list[str]] = mapped_column(ARRAY(Text))
status: Mapped[str] = mapped_column(Text)
reason: Mapped[str | None] = mapped_column(Text)
# Piège : une `UniqueConstraint` ne dédoublonnerait pas les lignes globales, dont `site_id` est
# NULL et qu'aucune n'est égale à une autre. Même forme que `uq_reading_source`.
Index(
"uq_drift_report_window",
DriftReport.window_end,
func.coalesce(DriftReport.site_id, text("''")),
unique=True,
)
-105
View File
@@ -1,105 +0,0 @@
# Surveillance de dérive du modèle de prévision (EC06, issue #45) : même gabarit que
# `app.detection.internal_alerts`, ordonnancé par le DAG `derive`.
from __future__ import annotations
import argparse
import asyncio
import sys
from datetime import UTC, datetime, timedelta
from app.core.config import get_settings
from app.db.session import get_session_factory
from app.repositories.drift import DriftRepository, NouveauRapportDerive
from app.services.drift import STATUT_DERIVE, DriftService, Seuils
async def run_drift(
*, now: datetime | None = None, site_id: str | None = None, seuils: Seuils | None = None
) -> list[NouveauRapportDerive]:
"""Calcule les rapports de la fenêtre et les enregistre. Rend ce qui a été calculé, que la
ligne ait été écrite ou ignorée par l'index d'idempotence."""
async with get_session_factory()() as session:
depot = DriftRepository(session)
rapports = await DriftService(depot, seuils=seuils).evaluate(now=now, site_id=site_id)
await depot.enregistre(rapports)
await session.commit()
return rapports
def _parse_instant(valeur: str) -> datetime:
instant = datetime.fromisoformat(valeur)
return instant if instant.tzinfo is not None else instant.replace(tzinfo=UTC)
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
defauts = Seuils()
parser = argparse.ArgumentParser(
prog="python -m app.monitoring.drift",
description="Surveillance de dérive du modèle de prévision EnerVision",
)
parser.add_argument("--site-id", default=None, help="Limite le calcul à un seul site.")
parser.add_argument(
"--now",
type=_parse_instant,
default=None,
help=(
"Instant de référence (ISO 8601, UTC si le fuseau est omis). Défaut : l'heure courante."
),
)
parser.add_argument(
"--window-hours",
type=int,
default=int(defauts.fenetre.total_seconds() // 3600),
help="Durée de la fenêtre récente, et de la fenêtre de référence qui la précède.",
)
parser.add_argument(
"--grace-hours",
type=int,
default=int(defauts.grace.total_seconds() // 3600),
help="Délai laissé à l'ingestion avant qu'une prévision soit jugée vérifiable.",
)
parser.add_argument(
"--min-observations",
type=int,
default=defauts.min_observations,
help="En deçà, le verdict est `indetermine` plutôt qu'un chiffre trompeur.",
)
parser.add_argument(
"--fail-on-drift",
action="store_true",
help="Sort en code non nul si une dérive est constatée, pour que la tâche rougisse.",
)
return parser.parse_args(argv)
def seuils_depuis(args: argparse.Namespace) -> Seuils:
return Seuils(
fenetre=timedelta(hours=args.window_hours),
grace=timedelta(hours=args.grace_hours),
min_observations=args.min_observations,
)
def main(argv: list[str] | None = None) -> int:
args = parse_args(argv)
# Échoue tôt si `APP_SECRET_KEY`/`DATABASE_URL` manquent, avant toute requête à la base.
get_settings()
rapports = asyncio.run(
run_drift(now=args.now, site_id=args.site_id, seuils=seuils_depuis(args))
)
for rapport in rapports:
cible = rapport.site_id or "TOUS SITES"
mae = f"{rapport.mae:.2f}" if rapport.mae is not None else "-"
print(
f"{cible} : {rapport.status}, MAE {mae} kWh sur {rapport.n_observations} prévision(s)"
f"{' : ' + rapport.reason if rapport.reason else ''}"
)
derive = any(rapport.status == STATUT_DERIVE for rapport in rapports)
return 1 if derive and args.fail_on_drift else 0
if __name__ == "__main__": # pragma: no cover
sys.exit(main())
-188
View File
@@ -1,188 +0,0 @@
"""Piège : deux dédoublonnages, pas un - DriftRepository.paires()
`prediction` n'a pas d'unicité sur `(site_id, target_at)` : chaque run de scoring empile une
ligne de plus. `uq_reading_source` autorise de son côté deux lectures au même instant quand la
`source` diffère. Joindre les deux tables sans `DISTINCT ON` des deux côtés compterait donc la
même heure plusieurs fois, et la moyenne d'erreur pèserait ces sites en double.
On retient la prédiction du run le plus récent, celle que sert `GET /api/v1/predictions`, avec
`prediction_id` en départage : `created_at` vaut l'heure de début de transaction et ne
distingue pas deux lignes du même run.
Côté lectures, le départage est `reading_id` décroissant, la règle même de
`GET /api/v1/sites/{site_id}/current`. La dérive se mesure donc contre le réalisé que l'API
affiche, et non contre une source élue ici et nulle part ailleurs.
"""
from collections.abc import Sequence
from dataclasses import asdict, dataclass
from datetime import datetime
from sqlalchemy import Subquery, func, select
from sqlalchemy.dialects.postgresql import insert
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.energy import DriftReport, Prediction, Reading
TARGET_METRIC = "consumption_kwh"
STATUT_DISPONIBLE = "available"
@dataclass(frozen=True, slots=True)
class PaireDerive:
site_id: str
target_at: datetime
predicted_value: float
actual_value: float
model_reference: str
@dataclass(frozen=True, slots=True)
class NouveauRapportDerive:
site_id: str | None
window_start: datetime
window_end: datetime
reference_start: datetime | None
reference_end: datetime | None
n_observations: int
mae: float | None
mape: float | None
bias: float | None
reference_mae: float | None
coverage_ratio: float | None
insufficient_data_ratio: float | None
model_references: list[str]
status: str
reason: str | None
@dataclass(frozen=True, slots=True)
class ComptageStatut:
site_id: str
status: str
nombre: int
def _predictions_retenues(*, debut: datetime, fin: datetime, site_id: str | None) -> Subquery:
requete = (
select(
Prediction.site_id,
Prediction.target_at,
Prediction.predicted_value,
Prediction.model_reference,
Prediction.status,
)
.distinct(Prediction.site_id, Prediction.target_at)
.where(
Prediction.target_metric == TARGET_METRIC,
Prediction.target_at >= debut,
Prediction.target_at < fin,
)
.order_by(Prediction.site_id, Prediction.target_at, Prediction.prediction_id.desc())
)
if site_id is not None:
requete = requete.where(Prediction.site_id == site_id)
return requete.subquery()
def _lectures_retenues(*, debut: datetime, fin: datetime, site_id: str | None) -> Subquery:
requete = (
select(Reading.site_id, Reading.timestamp, Reading.consumption_kwh)
.distinct(Reading.site_id, Reading.timestamp)
.where(
Reading.timestamp >= debut,
Reading.timestamp < fin,
Reading.consumption_kwh.is_not(None),
)
.order_by(Reading.site_id, Reading.timestamp, Reading.reading_id.desc())
)
if site_id is not None:
requete = requete.where(Reading.site_id == site_id)
return requete.subquery()
class DriftRepository:
def __init__(self, session: AsyncSession) -> None:
self._session = session
async def paires(
self, *, debut: datetime, fin: datetime, site_id: str | None = None
) -> Sequence[PaireDerive]:
predictions = _predictions_retenues(debut=debut, fin=fin, site_id=site_id)
lectures = _lectures_retenues(debut=debut, fin=fin, site_id=site_id)
requete = (
select(
predictions.c.site_id,
predictions.c.target_at,
predictions.c.predicted_value,
lectures.c.consumption_kwh,
predictions.c.model_reference,
)
.select_from(predictions)
.join(
lectures,
(lectures.c.site_id == predictions.c.site_id)
& (lectures.c.timestamp == predictions.c.target_at),
)
.where(predictions.c.status == STATUT_DISPONIBLE)
.order_by(predictions.c.site_id, predictions.c.target_at)
)
lignes = await self._session.execute(requete)
return [
PaireDerive(
site_id=ligne[0],
target_at=ligne[1],
predicted_value=ligne[2],
actual_value=ligne[3],
model_reference=ligne[4],
)
for ligne in lignes
]
async def comptages(
self, *, debut: datetime, fin: datetime, site_id: str | None = None
) -> Sequence[ComptageStatut]:
predictions = _predictions_retenues(debut=debut, fin=fin, site_id=site_id)
requete = (
select(predictions.c.site_id, predictions.c.status, func.count())
.select_from(predictions)
.group_by(predictions.c.site_id, predictions.c.status)
)
lignes = await self._session.execute(requete)
return [
ComptageStatut(site_id=ligne[0], status=ligne[1], nombre=ligne[2]) for ligne in lignes
]
# Pourquoi : l'idempotence est déléguée à `uq_drift_report_window` plutôt qu'à une lecture
# préalable, comme pour les recommandations. Rejouer la commande sur la même fenêtre ne
# duplique donc rien.
async def enregistre(self, rapports: Sequence[NouveauRapportDerive]) -> int:
if not rapports:
return 0
valeurs = [asdict(rapport) for rapport in rapports]
requete = (
insert(DriftReport)
.values(valeurs)
.on_conflict_do_nothing(
index_elements=[DriftReport.window_end, func.coalesce(DriftReport.site_id, "")]
)
.returning(DriftReport.drift_report_id)
)
return len((await self._session.scalars(requete)).all())
async def derniers(self, *, site_id: str | None = None) -> Sequence[DriftReport]:
requete = (
select(DriftReport)
.distinct(DriftReport.site_id)
.order_by(
DriftReport.site_id,
DriftReport.computed_at.desc(),
DriftReport.drift_report_id.desc(),
)
)
if site_id is not None:
requete = requete.where(DriftReport.site_id == site_id)
return (await self._session.scalars(requete)).all()
-31
View File
@@ -1,31 +0,0 @@
from datetime import datetime
from enum import StrEnum
from pydantic import BaseModel, ConfigDict
class DriftStatus(StrEnum):
STABLE = "stable"
DERIVE = "derive"
INDETERMINE = "indetermine"
class DriftReportResponse(BaseModel):
model_config = ConfigDict(from_attributes=True)
site_id: str | None
computed_at: datetime
window_start: datetime
window_end: datetime
reference_start: datetime | None
reference_end: datetime | None
n_observations: int
mae: float | None
mape: float | None
bias: float | None
reference_mae: float | None
coverage_ratio: float | None
insufficient_data_ratio: float | None
model_references: list[str]
status: DriftStatus
reason: str | None
-245
View File
@@ -1,245 +0,0 @@
"""Contrainte : la dérive se mesure sur ce qui a déjà eu lieu - DriftService.evaluate()
Une prévision ne devient vérifiable que quand la lecture de son instant cible est ingérée. La
fenêtre est donc fermée à droite par un délai de grâce : sans lui, la dernière heure ferait
chuter le taux de couverture à chaque exécution, et le verdict dirait « dérive » alors que
seule l'ingestion n'avait pas fini son tour.
La comparaison se fait entre deux fenêtres vives de même durée, pas contre la métrique de
référence du modèle journalisée à l'entraînement. Ce ne sont pas les mêmes grandeurs :
l'entraînement mesure un backtest où la météo de l'heure cible est connue, le scoring prévoit
une heure future dont la météo ne l'est pas. Les comparer classerait le modèle « en dérive »
dès le premier jour, ce qui ne prouverait rien.
"""
from collections.abc import Sequence
from dataclasses import dataclass, replace
from datetime import UTC, datetime, timedelta
from app.models.energy import DriftReport
from app.repositories.drift import (
ComptageStatut,
DriftRepository,
NouveauRapportDerive,
PaireDerive,
)
STATUT_STABLE = "stable"
STATUT_DERIVE = "derive"
STATUT_INDETERMINE = "indetermine"
STATUT_INSUFFISANT = "insufficient_data"
STATUT_DISPONIBLE = "available"
@dataclass(frozen=True, slots=True)
class Seuils:
# 168 h, la saisonnalité hebdomadaire que le modèle apprend par son lag principal : une
# fenêtre plus courte comparerait un week-end à une semaine ouvrée.
fenetre: timedelta = timedelta(hours=168)
grace: timedelta = timedelta(hours=2)
min_observations: int = 24
ratio_derive: float = 1.25
mae_plancher: float = 0.0
seuil_biais: float = 0.0
seuil_couverture: float = 0.8
@dataclass(frozen=True, slots=True)
class Metriques:
n_observations: int
mae: float | None
mape: float | None
bias: float | None
model_references: list[str]
def mesure(paires: Sequence[PaireDerive]) -> Metriques:
if not paires:
return Metriques(n_observations=0, mae=None, mape=None, bias=None, model_references=[])
ecarts = [paire.predicted_value - paire.actual_value for paire in paires]
# Le MAPE diverge sur une consommation nulle : les sites à l'arrêt sortent de ce seul
# rapport, jamais des autres métriques.
ratios = [
abs(ecart / paire.actual_value)
for ecart, paire in zip(ecarts, paires, strict=True)
if paire.actual_value != 0
]
return Metriques(
n_observations=len(paires),
mae=sum(abs(ecart) for ecart in ecarts) / len(ecarts),
mape=(sum(ratios) / len(ratios) * 100) if ratios else None,
bias=sum(ecarts) / len(ecarts),
model_references=sorted({paire.model_reference for paire in paires}),
)
@dataclass(frozen=True, slots=True)
class Verdict:
status: str
reason: str | None
class DriftService:
def __init__(self, depot: DriftRepository, *, seuils: Seuils | None = None) -> None:
self._depot = depot
self._seuils = seuils or Seuils()
async def derniers(self, *, site_id: str | None = None) -> Sequence[DriftReport]:
"""Ce que sert l'API : le dernier rapport de chaque site, plus la ligne globale."""
return await self._depot.derniers(site_id=site_id)
async def evaluate(
self, *, now: datetime | None = None, site_id: str | None = None
) -> list[NouveauRapportDerive]:
"""Une ligne par site, plus une ligne globale dont le `site_id` est nul."""
# Piège : `window_end` est la clé d'idempotence de `uq_drift_report_window`. Sans
# troncature à l'heure, deux exécutions ne collident jamais, aux microsecondes près.
instant = (now or datetime.now(UTC)).replace(minute=0, second=0, microsecond=0)
fin = instant - self._seuils.grace
debut = fin - self._seuils.fenetre
reference_fin = debut
reference_debut = reference_fin - self._seuils.fenetre
recentes = await self._depot.paires(debut=debut, fin=fin, site_id=site_id)
anciennes = await self._depot.paires(
debut=reference_debut, fin=reference_fin, site_id=site_id
)
comptages = await self._depot.comptages(debut=debut, fin=fin, site_id=site_id)
gabarit = NouveauRapportDerive(
site_id=None,
window_start=debut,
window_end=fin,
reference_start=reference_debut,
reference_end=reference_fin,
n_observations=0,
mae=None,
mape=None,
bias=None,
reference_mae=None,
coverage_ratio=None,
insufficient_data_ratio=None,
model_references=[],
status=STATUT_INDETERMINE,
reason=None,
)
rapports = [
self._rapport(
gabarit,
site=site,
recentes=[p for p in recentes if p.site_id == site],
anciennes=[p for p in anciennes if p.site_id == site],
comptages=[c for c in comptages if c.site_id == site],
)
# Piège : un site présent dans la référence et absent de la fenêtre récente a cessé
# d'être scoré. C'est la panne à crier, pas une ligne à omettre.
for site in sorted(
{paire.site_id for paire in recentes}
| {paire.site_id for paire in anciennes}
| {c.site_id for c in comptages}
)
]
rapports.append(
self._rapport(
gabarit, site=None, recentes=recentes, anciennes=anciennes, comptages=comptages
)
)
return rapports
def _rapport(
self,
gabarit: NouveauRapportDerive,
*,
site: str | None,
recentes: Sequence[PaireDerive],
anciennes: Sequence[PaireDerive],
comptages: Sequence[ComptageStatut],
) -> NouveauRapportDerive:
metriques = mesure(recentes)
reference = mesure(anciennes)
couverture = _couverture(len(recentes), comptages)
verdict = self._verdict(metriques, reference_mae=reference.mae, couverture=couverture)
return replace(
gabarit,
site_id=site,
n_observations=metriques.n_observations,
mae=metriques.mae,
mape=metriques.mape,
bias=metriques.bias,
reference_mae=reference.mae,
coverage_ratio=couverture,
insufficient_data_ratio=_part_insuffisante(comptages),
model_references=metriques.model_references,
status=verdict.status,
reason=verdict.reason,
)
def _verdict(
self, metriques: Metriques, *, reference_mae: float | None, couverture: float | None
) -> Verdict:
seuils = self._seuils
if metriques.n_observations < seuils.min_observations:
return Verdict(
STATUT_INDETERMINE,
f"{metriques.n_observations} prévision(s) vérifiée(s) sur la fenêtre, "
f"minimum {seuils.min_observations}.",
)
if couverture is not None and couverture < seuils.seuil_couverture:
return Verdict(
STATUT_DERIVE,
f"Couverture de {couverture:.0%}, sous le seuil de {seuils.seuil_couverture:.0%} : "
"le pipeline, pas le modèle.",
)
plafond = _plafond(reference_mae, ratio=seuils.ratio_derive, plancher=seuils.mae_plancher)
if plafond is None:
return Verdict(
STATUT_INDETERMINE,
"Fenêtre de référence sans paire vérifiée, et aucun plancher de MAE : "
"rien à quoi comparer cette fenêtre.",
)
if metriques.mae is not None and metriques.mae > plafond:
return Verdict(
STATUT_DERIVE,
f"MAE de {metriques.mae:.2f} kWh au-delà de {plafond:.2f} kWh, "
"seuil dérivé de la fenêtre de référence.",
)
if (
seuils.seuil_biais > 0
and metriques.bias is not None
and abs(metriques.bias) > seuils.seuil_biais
):
return Verdict(
STATUT_DERIVE,
f"Biais de {metriques.bias:+.2f} kWh : le modèle se trompe toujours du même côté.",
)
return Verdict(STATUT_STABLE, None)
def _plafond(reference_mae: float | None, *, ratio: float, plancher: float) -> float | None:
if reference_mae is None:
return plancher or None
return max(plancher, reference_mae * ratio)
def _couverture(apparie: int, comptages: Sequence[ComptageStatut]) -> float | None:
"""Part des prévisions disponibles qui ont trouvé leur réalisé. Mesure l'ingestion et
l'ordonnancement, pas la qualité du modèle."""
disponibles = sum(c.nombre for c in comptages if c.status == STATUT_DISPONIBLE)
return apparie / disponibles if disponibles else None
def _part_insuffisante(comptages: Sequence[ComptageStatut]) -> float | None:
total = sum(c.nombre for c in comptages)
if not total:
return None
return sum(c.nombre for c in comptages if c.status == STATUT_INSUFFISANT) / total
+22 -288
View File
@@ -213,7 +213,7 @@
}, },
"security": [ "security": [
{ {
"Cookie de rafraîchissement": [] "CookieRafraichissement": []
} }
] ]
} }
@@ -252,7 +252,7 @@
}, },
"security": [ "security": [
{ {
"Cookie de rafraîchissement": [] "CookieRafraichissement": []
} }
] ]
} }
@@ -301,7 +301,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -347,7 +347,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -423,7 +423,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -673,7 +673,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
}, },
@@ -757,7 +757,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -771,7 +771,7 @@
"operationId": "update_user_api_v1_users__user_id__patch", "operationId": "update_user_api_v1_users__user_id__patch",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -889,7 +889,7 @@
"operationId": "reset_password_api_v1_users__user_id__password_reset_post", "operationId": "reset_password_api_v1_users__user_id__password_reset_post",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1023,7 +1023,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -1037,7 +1037,7 @@
"operationId": "get_site_api_v1_sites__site_id__get", "operationId": "get_site_api_v1_sites__site_id__get",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1124,7 +1124,7 @@
"operationId": "get_current_api_v1_sites__site_id__current_get", "operationId": "get_current_api_v1_sites__site_id__current_get",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1211,7 +1211,7 @@
"operationId": "list_alerts_api_v1_alerts_get", "operationId": "list_alerts_api_v1_alerts_get",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1361,7 +1361,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -1375,7 +1375,7 @@
"operationId": "get_recommendation_api_v1_recommendations__recommendation_id__get", "operationId": "get_recommendation_api_v1_recommendations__recommendation_id__get",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1462,7 +1462,7 @@
"operationId": "generate_recommendations_api_v1_recommendations_generate_post", "operationId": "generate_recommendations_api_v1_recommendations_generate_post",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1588,7 +1588,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -1602,7 +1602,7 @@
"operationId": "list_readings_api_v1_readings_get", "operationId": "list_readings_api_v1_readings_get",
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
], ],
"parameters": [ "parameters": [
@@ -1799,7 +1799,7 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
@@ -1855,98 +1855,10 @@
}, },
"security": [ "security": [
{ {
"Jeton d'accès": [] "JetonAcces": []
} }
] ]
} }
},
"/api/v1/monitoring/drift": {
"get": {
"tags": [
"monitoring"
],
"summary": "Dernier rapport de dérive par site, plus la ligne globale",
"operationId": "get_drift_api_v1_monitoring_drift_get",
"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": {
"type": "array",
"items": {
"$ref": "#/components/schemas/DriftReportResponse"
},
"title": "Response Get Drift Api V1 Monitoring Drift Get"
}
}
}
},
"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"
}
}
}
}
}
}
} }
}, },
"components": { "components": {
@@ -2065,180 +1977,6 @@
], ],
"title": "AlertType" "title": "AlertType"
}, },
"DriftReportResponse": {
"properties": {
"site_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"title": "Site Id"
},
"computed_at": {
"type": "string",
"format": "date-time",
"title": "Computed At"
},
"window_start": {
"type": "string",
"format": "date-time",
"title": "Window Start"
},
"window_end": {
"type": "string",
"format": "date-time",
"title": "Window End"
},
"reference_start": {
"anyOf": [
{
"type": "string",
"format": "date-time"
},
{
"type": "null"
}
],
"title": "Reference Start"
},
"reference_end": {
"anyOf": [
{
"type": "string",
"format": "date-time"
},
{
"type": "null"
}
],
"title": "Reference End"
},
"n_observations": {
"type": "integer",
"title": "N Observations"
},
"mae": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Mae"
},
"mape": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Mape"
},
"bias": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Bias"
},
"reference_mae": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Reference Mae"
},
"coverage_ratio": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Coverage Ratio"
},
"insufficient_data_ratio": {
"anyOf": [
{
"type": "number"
},
{
"type": "null"
}
],
"title": "Insufficient Data Ratio"
},
"model_references": {
"items": {
"type": "string"
},
"type": "array",
"title": "Model References"
},
"status": {
"$ref": "#/components/schemas/DriftStatus"
},
"reason": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"title": "Reason"
}
},
"type": "object",
"required": [
"site_id",
"computed_at",
"window_start",
"window_end",
"reference_start",
"reference_end",
"n_observations",
"mae",
"mape",
"bias",
"reference_mae",
"coverage_ratio",
"insufficient_data_ratio",
"model_references",
"status",
"reason"
],
"title": "DriftReportResponse"
},
"DriftStatus": {
"type": "string",
"enum": [
"stable",
"derive",
"indetermine"
],
"title": "DriftStatus"
},
"ErrorResponse": { "ErrorResponse": {
"properties": { "properties": {
"detail": { "detail": {
@@ -3487,13 +3225,13 @@
} }
}, },
"securitySchemes": { "securitySchemes": {
"Cookie de rafraîchissement": { "CookieRafraichissement": {
"type": "apiKey", "type": "apiKey",
"description": "Cookie `HttpOnly` posé par `/auth/login` et tourné par `/auth/refresh`. Il prend le préfixe `__Secure-` dès que l'API tourne derrière TLS, et n'est émis que vers `/api/v1/auth`.", "description": "Cookie `HttpOnly` posé par `/auth/login` et tourné par `/auth/refresh`. Il prend le préfixe `__Secure-` dès que l'API tourne derrière TLS, et n'est émis que vers `/api/v1/auth`.",
"in": "cookie", "in": "cookie",
"name": "ev_refresh" "name": "ev_refresh"
}, },
"Jeton d'accès": { "JetonAcces": {
"type": "http", "type": "http",
"scheme": "bearer" "scheme": "bearer"
} }
@@ -3539,10 +3277,6 @@
{ {
"name": "predictions", "name": "predictions",
"description": "Dernière prévision de consommation par site, calculée hors ligne par le pipeline de scoring (`ml/`) et simplement lue ici. Accessible à partir du rôle `lecteur`." "description": "Dernière prévision de consommation par site, calculée hors ligne par le pipeline de scoring (`ml/`) et simplement lue ici. Accessible à partir du rôle `lecteur`."
},
{
"name": "monitoring",
"description": "Surveillance de la dérive du modèle : écart entre les prévisions déjà écrites et les lectures réellement arrivées, par site et tous sites confondus. Réservé à partir du rôle `operateur`, qui agit sur un pipeline dégradé."
} }
] ]
} }
+2 -5
View File
@@ -87,11 +87,8 @@ disallow_untyped_defs = false
testpaths = ["tests"] testpaths = ["tests"]
asyncio_mode = "auto" asyncio_mode = "auto"
asyncio_default_fixture_loop_scope = "function" asyncio_default_fixture_loop_scope = "function"
addopts = "-q --strict-markers -m 'not integration and not chaine' --cov=app --cov-report=term-missing" addopts = "-q --strict-markers -m 'not integration' --cov=app --cov-report=term-missing"
markers = [ markers = ["integration: requiert une base PostgreSQL joignable, hors `make test`"]
"integration: requiert une base PostgreSQL joignable, hors `make test`",
"chaine: requiert en plus l'environnement uv de ml/, hors `make test` et hors `-m integration`",
]
[tool.coverage.run] [tool.coverage.run]
source = ["app"] source = ["app"]
-1
View File
@@ -56,7 +56,6 @@ ROLE_MINIMUM: Final[dict[Route, Role]] = {
("GET", "/api/v1/readings"): Role.LECTEUR, ("GET", "/api/v1/readings"): Role.LECTEUR,
("GET", "/api/v1/predictions"): Role.LECTEUR, ("GET", "/api/v1/predictions"): Role.LECTEUR,
("GET", "/api/v1/sensors/status"): Role.ADMIN, ("GET", "/api/v1/sensors/status"): Role.ADMIN,
("GET", "/api/v1/monitoring/drift"): Role.OPERATEUR,
("GET", "/api/v1/users"): Role.ADMIN, ("GET", "/api/v1/users"): Role.ADMIN,
("POST", "/api/v1/users"): Role.ADMIN, ("POST", "/api/v1/users"): Role.ADMIN,
("PATCH", "/api/v1/users/{user_id}"): Role.ADMIN, ("PATCH", "/api/v1/users/{user_id}"): Role.ADMIN,
-132
View File
@@ -1,132 +0,0 @@
"""Piège : ces fixtures valident leurs écritures, contrairement à celles de tests/repositories.
Un endpoint ouvre sa propre session par `get_session` : il ne verrait pas une ligne semée dans
une transaction en cours. Lui passer la session de la fixture par `dependency_overrides`
supprimerait justement ce que ces tests prouvent, et `RecommendationService.generate` valide de
toute façon lui-même. L'isolation vient donc de la marque portée par chaque `site_id`, et le
nettoyage est explicite, dans l'ordre imposé par les clés étrangères `RESTRICT`.
Contrainte : toutes ces fixtures sont à portée fonction. `engine_per_test` vide le cache du
moteur après chaque test ; une fixture de module verrait un moteur déjà fermé à son démontage,
et ses lignes resteraient en base.
"""
from collections.abc import AsyncIterator, Callable, Iterator
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from uuid import uuid4
import pytest
from fastapi import FastAPI
from sqlalchemy import delete, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.api.deps import get_current_principal
from app.core.principal import Principal
from app.core.roles import AccountKind, Role
from app.db.session import get_session_factory
from app.models.energy import Alert, DriftReport, Prediction, Reading, Recommendation, Site
from tests.repositories.test_alert import creer_alerte
from tests.repositories.test_prediction import creer_prediction
from tests.repositories.test_reading import creer_lecture
from tests.repositories.test_site import creer as creer_site
INSTANT = datetime(2026, 9, 16, 12, 0, tzinfo=UTC)
@dataclass(frozen=True)
class JeuMetier:
"""Identifiants seuls, jamais d'instance ORM : un attribut relu sur une session fermée
déclenche un `MissingGreenlet`."""
site_id: str
site_voisin: str
alert_id: int
prediction_id: int
instant: datetime
async def _supprime(session: AsyncSession, sites: list[str]) -> None:
# La suppression des recommandations est inconditionnelle : `POST /generate` en cree hors du
# controle de la fixture, et `alert` les retient par une cle etrangere `RESTRICT`.
alertes = select(Alert.alert_id).where(Alert.site_id.in_(sites))
await session.execute(delete(Recommendation).where(Recommendation.alert_id.in_(alertes)))
await session.execute(delete(Alert).where(Alert.site_id.in_(sites)))
await session.execute(delete(DriftReport).where(DriftReport.site_id.in_(sites)))
await session.execute(delete(Prediction).where(Prediction.site_id.in_(sites)))
await session.execute(delete(Reading).where(Reading.site_id.in_(sites)))
await session.execute(delete(Site).where(Site.site_id.in_(sites)))
await session.commit()
@pytest.fixture
def marque() -> str:
return uuid4().hex[:12]
@pytest.fixture
async def jeu_metier(marque: str) -> AsyncIterator[JeuMetier]:
"""Un site instrumenté, un site voisin, trois lectures horaires, une prédiction, une alerte.
Le voisin existe pour que les tests de filtre prouvent qu'ils écartent quelque chose.
"""
site_id = f"SITE-{marque}"
voisin = f"SITE-{marque}-VOISIN"
async with get_session_factory()() as session:
await creer_site(session, site_id=site_id, capacity_kw=100.0)
await creer_site(session, site_id=voisin, capacity_kw=100.0)
for decalage in range(3):
await creer_lecture(
session,
site_id=site_id,
timestamp=INSTANT - timedelta(hours=decalage),
consumption_kw=10.0 + decalage,
)
prediction = await creer_prediction(session, site_id=site_id, target_at=INSTANT)
alerte = await creer_alerte(session, site_id=site_id, timestamp=INSTANT)
jeu = JeuMetier(
site_id=site_id,
site_voisin=voisin,
alert_id=alerte.alert_id,
prediction_id=prediction.prediction_id,
instant=INSTANT,
)
await session.commit()
try:
yield jeu
finally:
async with get_session_factory()() as session:
await _supprime(session, [site_id, voisin])
@pytest.fixture
async def site_nu(marque: str) -> AsyncIterator[str]:
"""Un site sans lecture ni prédiction : le cas que seul un vrai `LEFT JOIN` distingue."""
site_id = f"SITE-{marque}-NU"
async with get_session_factory()() as session:
await creer_site(session, site_id=site_id, capacity_kw=100.0)
await session.commit()
try:
yield site_id
finally:
async with get_session_factory()() as session:
await _supprime(session, [site_id])
@pytest.fixture
def principal_injecte(app: FastAPI) -> Iterator[Callable[[Role], None]]:
def installe(role: Role = Role.LECTEUR) -> None:
app.dependency_overrides[get_current_principal] = lambda: Principal(
id=uuid4(),
email="parcours@enervision.fr",
role=role,
kind=AccountKind.HUMAIN,
must_change_password=False,
)
yield installe
app.dependency_overrides.pop(get_current_principal, None)
-105
View File
@@ -1,105 +0,0 @@
from collections.abc import Iterator, Sequence
from datetime import UTC, datetime, timedelta
from uuid import uuid4
import pytest
from fastapi import FastAPI
from httpx import AsyncClient
from app.api.deps import get_current_principal, get_drift_service
from app.core.principal import Principal
from app.core.roles import AccountKind, Role
from app.models.energy import DriftReport
INSTANT = datetime(2026, 9, 22, 12, tzinfo=UTC)
def operateur() -> Principal:
return Principal(
id=uuid4(),
email="operateur@enervision.fr",
role=Role.OPERATEUR,
kind=AccountKind.HUMAIN,
must_change_password=False,
)
def rapport(*, site_id: str | None) -> DriftReport:
return DriftReport(
drift_report_id=1,
computed_at=INSTANT,
site_id=site_id,
window_start=INSTANT - timedelta(hours=168),
window_end=INSTANT,
reference_start=None,
reference_end=None,
n_observations=48,
mae=1.5,
mape=12.0,
bias=0.3,
reference_mae=1.2,
coverage_ratio=0.95,
insufficient_data_ratio=0.0,
model_references=["lightgbm-aaa"],
status="stable",
reason=None,
)
class FauxService:
def __init__(self, rapports: Sequence[DriftReport]) -> None:
self.rapports = list(rapports)
self.site_demande: str | None = None
async def derniers(self, *, site_id: str | None = None) -> Sequence[DriftReport]:
self.site_demande = site_id
return self.rapports
@pytest.fixture
def servi(app: FastAPI) -> Iterator[list[DriftReport]]:
rapports = [rapport(site_id="SITE001"), rapport(site_id=None)]
service = FauxService(rapports)
app.dependency_overrides[get_current_principal] = operateur
app.dependency_overrides[get_drift_service] = lambda: service
yield rapports
app.dependency_overrides.pop(get_current_principal, None)
app.dependency_overrides.pop(get_drift_service, None)
async def test_drift_returns_the_latest_report_of_every_site(
servi: list[DriftReport], client: AsyncClient
) -> None:
reponse = await client.get("/api/v1/monitoring/drift")
assert reponse.status_code == 200
assert [ligne["site_id"] for ligne in reponse.json()] == ["SITE001", None]
async def test_drift_exposes_the_metrics_of_the_stored_report(
servi: list[DriftReport], client: AsyncClient
) -> None:
reponse = await client.get("/api/v1/monitoring/drift")
premier = reponse.json()[0]
assert premier["status"] == "stable"
assert premier["mae"] == 1.5
assert premier["model_references"] == ["lightgbm-aaa"]
@pytest.fixture
def sans_rapport(app: FastAPI) -> Iterator[None]:
app.dependency_overrides[get_current_principal] = operateur
app.dependency_overrides[get_drift_service] = lambda: FauxService([])
yield
app.dependency_overrides.pop(get_current_principal, None)
app.dependency_overrides.pop(get_drift_service, None)
async def test_drift_returns_an_empty_list_when_no_report_exists(
sans_rapport: None, client: AsyncClient
) -> None:
reponse = await client.get("/api/v1/monitoring/drift")
assert reponse.status_code == 200
assert reponse.json() == []
+2 -2
View File
@@ -104,8 +104,8 @@ def test_the_rate_limit_documents_the_delay_header(schema: dict[str, Any]) -> No
def test_the_refresh_cookie_appears_in_the_security_schemes(schema: dict[str, Any]) -> None: def test_the_refresh_cookie_appears_in_the_security_schemes(schema: dict[str, Any]) -> None:
schemes = schema["components"]["securitySchemes"] schemes = schema["components"]["securitySchemes"]
assert schemes["Cookie de rafraîchissement"]["in"] == "cookie" assert schemes["CookieRafraichissement"]["in"] == "cookie"
assert schemes["Cookie de rafraîchissement"]["name"] == "ev_refresh" assert schemes["CookieRafraichissement"]["name"] == "ev_refresh"
def test_each_tag_used_by_a_route_is_described(schema: dict[str, Any]) -> None: def test_each_tag_used_by_a_route_is_described(schema: dict[str, Any]) -> None:
@@ -1,92 +0,0 @@
from collections.abc import Callable
import pytest
from httpx import AsyncClient
from app.core.roles import Role
from tests.api.conftest import JeuMetier
pytestmark = pytest.mark.integration
async def genere(client: AsyncClient, site_id: str) -> dict[str, int]:
# Toujours borne a un site : sans `site_id`, le service examine toutes les alertes de la
# base, y compris celles d'un autre test, et le rapport cesse d'etre deterministe.
reponse = await client.post(f"/api/v1/recommendations/generate?site_id={site_id}")
assert reponse.status_code == 200
return dict(reponse.json())
async def test_generate_creates_a_recommendation_for_the_alert_of_the_requested_site(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.ADMIN)
rapport = await genere(client, jeu_metier.site_id)
assert rapport["alerts_examined"] == 1
assert rapport["recommendations_created"] >= 1
assert rapport["already_present"] == 0
async def test_generate_creates_nothing_more_when_it_runs_twice_on_the_same_alerts(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.ADMIN)
premier = await genere(client, jeu_metier.site_id)
second = await genere(client, jeu_metier.site_id)
assert second["recommendations_created"] == 0
assert second["already_present"] == premier["recommendations_created"]
async def test_generate_examines_no_alert_when_the_requested_site_has_none(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.ADMIN)
rapport = await genere(client, jeu_metier.site_voisin)
assert rapport["alerts_examined"] == 0
assert rapport["recommendations_created"] == 0
async def test_list_recommendations_returns_what_generate_persisted_in_another_session(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.ADMIN)
await genere(client, jeu_metier.site_id)
reponse = await client.get("/api/v1/recommendations")
assert reponse.status_code == 200
miennes = [r for r in reponse.json() if r["alert_id"] == jeu_metier.alert_id]
assert miennes != []
assert all(r["rule_reference"] for r in miennes)
async def test_get_recommendation_returns_the_row_created_by_generate(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.ADMIN)
await genere(client, jeu_metier.site_id)
liste = await client.get("/api/v1/recommendations")
creee = next(r for r in liste.json() if r["alert_id"] == jeu_metier.alert_id)
reponse = await client.get(f"/api/v1/recommendations/{creee['recommendation_id']}")
assert reponse.status_code == 200
assert reponse.json() == creee
async def test_get_recommendation_returns_404_when_the_identifier_is_unknown(
principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.LECTEUR)
reponse = await client.get("/api/v1/recommendations/9999999")
assert reponse.status_code == 404
assert reponse.json()["detail"] == "Recommandation introuvable"
@@ -1,80 +0,0 @@
from collections.abc import Callable
import pytest
from httpx import AsyncClient
from app.core.roles import Role
from app.db.session import get_session_factory
from tests.api.conftest import JeuMetier
from tests.repositories.test_reading import creer_lecture
pytestmark = pytest.mark.integration
async def test_list_sites_returns_the_seeded_site_with_its_stored_attributes(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.LECTEUR)
reponse = await client.get("/api/v1/sites")
assert reponse.status_code == 200
mien = next(site for site in reponse.json() if site["site_id"] == jeu_metier.site_id)
assert mien["capacity_kw"] == 100.0
assert mien["site_name"] == "Site de test"
async def test_get_site_returns_404_when_the_identifier_is_absent_from_the_database(
principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.LECTEUR)
reponse = await client.get("/api/v1/sites/SITE-JAMAIS-INSERE")
assert reponse.status_code == 404
assert reponse.json()["detail"] == "Site introuvable"
async def test_get_current_returns_the_most_recent_reading_when_several_hours_are_stored(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.LECTEUR)
reponse = await client.get(f"/api/v1/sites/{jeu_metier.site_id}/current")
assert reponse.status_code == 200
corps = reponse.json()
assert corps["consumption_kw"] == 10.0
assert corps["timestamp"].startswith("2026-09-16T12:00")
async def test_get_current_keeps_the_highest_reading_id_when_two_sources_share_the_timestamp(
jeu_metier: JeuMetier, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.LECTEUR)
async with get_session_factory()() as session:
await creer_lecture(
session,
site_id=jeu_metier.site_id,
timestamp=jeu_metier.instant,
source="api_history",
consumption_kw=999.0,
)
await session.commit()
reponse = await client.get(f"/api/v1/sites/{jeu_metier.site_id}/current")
assert reponse.json()["consumption_kw"] == 999.0
async def test_get_current_reports_a_critical_quality_when_the_site_has_no_reading(
site_nu: str, principal_injecte: Callable[[Role], None], client: AsyncClient
) -> None:
principal_injecte(Role.LECTEUR)
reponse = await client.get(f"/api/v1/sites/{site_nu}/current")
assert reponse.status_code == 200
corps = reponse.json()
assert corps["timestamp"] is None
assert corps["data_quality"] == "critical"
+2 -69
View File
@@ -1,5 +1,5 @@
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
from datetime import UTC, datetime, timedelta from datetime import UTC, datetime
from uuid import uuid4 from uuid import uuid4
import pytest import pytest
@@ -9,15 +9,7 @@ from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
from app.core.config import get_settings from app.core.config import get_settings
from app.models.energy import ( from app.models.energy import Alert, Dataset, Prediction, Reading, Recommendation, Site
Alert,
Dataset,
DriftReport,
Prediction,
Reading,
Recommendation,
Site,
)
pytestmark = pytest.mark.integration pytestmark = pytest.mark.integration
MOMENT = datetime(2024, 1, 1, tzinfo=UTC) MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
@@ -277,62 +269,3 @@ async def test_recommendation_is_unique_when_alert_and_rule_match(
with pytest.raises(IntegrityError): with pytest.raises(IntegrityError):
async with savepoint: async with savepoint:
await data_connection.execute(statement) await data_connection.execute(statement)
def _rapport(**remplacements: object) -> dict[str, object]:
defauts: dict[str, object] = {
"site_id": None,
"window_start": MOMENT,
"window_end": MOMENT,
"n_observations": 12,
"model_references": ["lightgbm-aaa"],
"status": "stable",
"reason": None,
}
return {**defauts, **remplacements}
async def test_drift_report_rejects_an_unknown_status(data_connection: AsyncConnection) -> None:
statement = insert(DriftReport).values(**_rapport(status="douteux", reason="x"))
savepoint = data_connection.begin_nested()
with pytest.raises(IntegrityError):
async with savepoint:
await data_connection.execute(statement)
async def test_drift_report_rejects_a_drift_without_a_reason(
data_connection: AsyncConnection,
) -> None:
statement = insert(DriftReport).values(**_rapport(status="derive"))
savepoint = data_connection.begin_nested()
with pytest.raises(IntegrityError):
async with savepoint:
await data_connection.execute(statement)
async def test_drift_report_accepts_one_global_row_without_a_site(
data_connection: AsyncConnection,
) -> None:
identifiant = (
await data_connection.execute(
insert(DriftReport).values(**_rapport()).returning(DriftReport.drift_report_id)
)
).scalar_one()
assert identifiant is not None
async def test_drift_report_is_unique_when_window_and_site_match(
data_connection: AsyncConnection,
) -> None:
fenetre = MOMENT + timedelta(days=1)
statement = insert(DriftReport).values(**_rapport(window_end=fenetre))
await data_connection.execute(statement)
savepoint = data_connection.begin_nested()
with pytest.raises(IntegrityError):
async with savepoint:
await data_connection.execute(statement)
@@ -1,187 +0,0 @@
from datetime import UTC, datetime, timedelta
import pytest
from sqlalchemy.dialects import postgresql
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.sql import ClauseElement
from app.repositories.drift import (
DriftRepository,
NouveauRapportDerive,
_lectures_retenues,
_predictions_retenues,
)
from tests.repositories.test_prediction import creer_prediction
from tests.repositories.test_reading import creer_lecture
from tests.repositories.test_site import creer as creer_site
DEBUT = datetime(2026, 9, 15, tzinfo=UTC)
FIN = datetime(2026, 9, 22, tzinfo=UTC)
CIBLE = datetime(2026, 9, 16, 12, tzinfo=UTC)
def sql(requete: ClauseElement) -> str:
return str(requete.compile(dialect=postgresql.dialect())) # type: ignore[no-untyped-call]
def rapport(**remplacements: object) -> NouveauRapportDerive:
defauts: dict[str, object] = {
"site_id": None,
"window_start": DEBUT,
"window_end": FIN,
"reference_start": None,
"reference_end": None,
"n_observations": 10,
"mae": 1.0,
"mape": 5.0,
"bias": 0.1,
"reference_mae": None,
"coverage_ratio": 1.0,
"insufficient_data_ratio": 0.0,
"model_references": ["lightgbm-aaa"],
"status": "stable",
"reason": None,
}
return NouveauRapportDerive(**{**defauts, **remplacements}) # type: ignore[arg-type]
def test_predictions_keep_one_row_per_site_and_target_in_sql() -> None:
requete = sql(_predictions_retenues(debut=DEBUT, fin=FIN, site_id=None).element)
assert "DISTINCT ON (prediction.site_id, prediction.target_at)" in requete
assert "prediction.prediction_id DESC" in requete
def test_readings_keep_one_row_per_site_and_instant_in_sql() -> None:
requete = sql(_lectures_retenues(debut=DEBUT, fin=FIN, site_id=None).element)
assert "DISTINCT ON (reading.site_id, reading.timestamp)" in requete
assert "reading.reading_id DESC" in requete
def test_predictions_restrict_themselves_to_the_requested_site_in_sql() -> None:
requete = sql(_predictions_retenues(debut=DEBUT, fin=FIN, site_id="SITE001").element)
assert requete.count("prediction.site_id = ") == 1
def test_readings_ignore_a_missing_consumption_in_sql() -> None:
requete = sql(_lectures_retenues(debut=DEBUT, fin=FIN, site_id=None).element)
assert "reading.consumption_kwh IS NOT NULL" in requete
@pytest.mark.integration
async def test_repository_pairs_a_prediction_with_the_reading_of_the_same_instant(
session: AsyncSession,
) -> None:
site = await creer_site(session)
await creer_prediction(session, site_id=site.site_id, target_at=CIBLE, predicted_value=12.0)
await creer_lecture(session, site_id=site.site_id, timestamp=CIBLE, consumption_kwh=10.0)
paires = await DriftRepository(session).paires(debut=DEBUT, fin=FIN, site_id=site.site_id)
await session.rollback()
assert [(p.predicted_value, p.actual_value) for p in paires] == [(12.0, 10.0)]
@pytest.mark.integration
async def test_repository_keeps_the_latest_run_when_several_predictions_share_a_target(
session: AsyncSession,
) -> None:
site = await creer_site(session)
await creer_prediction(session, site_id=site.site_id, target_at=CIBLE, predicted_value=12.0)
await creer_prediction(session, site_id=site.site_id, target_at=CIBLE, predicted_value=99.0)
await creer_lecture(session, site_id=site.site_id, timestamp=CIBLE, consumption_kwh=10.0)
paires = await DriftRepository(session).paires(debut=DEBUT, fin=FIN, site_id=site.site_id)
await session.rollback()
assert [p.predicted_value for p in paires] == [99.0]
@pytest.mark.integration
async def test_repository_keeps_one_reading_per_instant_when_two_sources_wrote_the_same_hour(
session: AsyncSession,
) -> None:
site = await creer_site(session)
await creer_prediction(session, site_id=site.site_id, target_at=CIBLE, predicted_value=12.0)
await creer_lecture(
session, site_id=site.site_id, timestamp=CIBLE, source="api_current", consumption_kwh=10.0
)
await creer_lecture(
session, site_id=site.site_id, timestamp=CIBLE, source="api_history", consumption_kwh=20.0
)
paires = await DriftRepository(session).paires(debut=DEBUT, fin=FIN, site_id=site.site_id)
await session.rollback()
assert [p.actual_value for p in paires] == [20.0]
@pytest.mark.integration
async def test_repository_excludes_an_insufficient_data_prediction_from_the_pairs(
session: AsyncSession,
) -> None:
site = await creer_site(session)
await creer_prediction(
session,
site_id=site.site_id,
target_at=CIBLE,
predicted_value=None,
status="insufficient_data",
failure_reason="historique trop court",
)
await creer_lecture(session, site_id=site.site_id, timestamp=CIBLE, consumption_kwh=10.0)
depot = DriftRepository(session)
paires = await depot.paires(debut=DEBUT, fin=FIN, site_id=site.site_id)
comptages = await depot.comptages(debut=DEBUT, fin=FIN, site_id=site.site_id)
await session.rollback()
assert paires == []
assert [(c.status, c.nombre) for c in comptages] == [("insufficient_data", 1)]
@pytest.mark.integration
async def test_repository_excludes_a_target_outside_the_window(session: AsyncSession) -> None:
site = await creer_site(session)
hors_fenetre = FIN + timedelta(hours=1)
await creer_prediction(
session, site_id=site.site_id, target_at=hors_fenetre, predicted_value=12.0
)
await creer_lecture(session, site_id=site.site_id, timestamp=hors_fenetre, consumption_kwh=10.0)
paires = await DriftRepository(session).paires(debut=DEBUT, fin=FIN, site_id=site.site_id)
await session.rollback()
assert paires == []
@pytest.mark.integration
async def test_repository_reads_back_the_global_report_it_wrote(session: AsyncSession) -> None:
depot = DriftRepository(session)
fenetre = datetime(2035, 3, 1, tzinfo=UTC)
ecrites = await depot.enregistre([rapport(window_end=fenetre)])
derniers = await depot.derniers()
globaux = [r for r in derniers if r.site_id is None and r.window_end == fenetre]
await session.rollback()
assert ecrites == 1
assert len(globaux) == 1
@pytest.mark.integration
async def test_repository_ignores_a_second_report_for_the_same_window_and_site(
session: AsyncSession,
) -> None:
depot = DriftRepository(session)
fenetre = datetime(2035, 4, 1, tzinfo=UTC)
premiere = await depot.enregistre([rapport(window_end=fenetre)])
seconde = await depot.enregistre([rapport(window_end=fenetre, status="derive", reason="x")])
await session.rollback()
assert premiere == 1
assert seconde == 0
-260
View File
@@ -1,260 +0,0 @@
from collections.abc import Sequence
from datetime import UTC, datetime, timedelta
import pytest
from app.repositories.drift import ComptageStatut, PaireDerive
from app.services.drift import (
STATUT_DERIVE,
STATUT_INDETERMINE,
STATUT_STABLE,
DriftService,
Seuils,
mesure,
)
INSTANT = datetime(2026, 9, 22, 12, 0, tzinfo=UTC)
def paire(
*, site_id: str = "SITE001", prevu: float, reel: float, reference: str = "lightgbm-aaa"
) -> PaireDerive:
return PaireDerive(
site_id=site_id,
target_at=INSTANT,
predicted_value=prevu,
actual_value=reel,
model_reference=reference,
)
def paires(
*, site_id: str = "SITE001", nombre: int, prevu: float, reel: float
) -> list[PaireDerive]:
return [paire(site_id=site_id, prevu=prevu, reel=reel) for _ in range(nombre)]
class FauxDepot:
def __init__(
self,
*,
recentes: Sequence[PaireDerive] = (),
anciennes: Sequence[PaireDerive] = (),
comptages: Sequence[ComptageStatut] = (),
) -> None:
self.recentes = list(recentes)
self.anciennes = list(anciennes)
self._comptages = list(comptages)
self.fenetres: list[tuple[datetime, datetime]] = []
async def paires(
self, *, debut: datetime, fin: datetime, site_id: str | None = None
) -> Sequence[PaireDerive]:
self.fenetres.append((debut, fin))
return self.recentes if len(self.fenetres) == 1 else self.anciennes
async def comptages(
self, *, debut: datetime, fin: datetime, site_id: str | None = None
) -> Sequence[ComptageStatut]:
return self._comptages
def service(depot: FauxDepot, **surcharges: object) -> DriftService:
return DriftService(depot, seuils=Seuils(**surcharges)) # type: ignore[arg-type]
def test_drift_averages_the_absolute_gap_between_forecast_and_actual() -> None:
metriques = mesure([paire(prevu=12.0, reel=10.0), paire(prevu=8.0, reel=10.0)])
assert metriques.mae == 2.0
assert metriques.n_observations == 2
def test_drift_computes_a_signed_bias_when_the_model_overforecasts() -> None:
metriques = mesure([paire(prevu=12.0, reel=10.0), paire(prevu=14.0, reel=10.0)])
assert metriques.bias == 3.0
def test_drift_computes_a_negative_bias_when_the_model_underforecasts() -> None:
metriques = mesure([paire(prevu=8.0, reel=10.0), paire(prevu=6.0, reel=10.0)])
assert metriques.bias == -3.0
def test_drift_excludes_a_zero_actual_from_the_mape_only() -> None:
metriques = mesure([paire(prevu=11.0, reel=10.0), paire(prevu=5.0, reel=0.0)])
assert metriques.mape == 10.0
assert metriques.n_observations == 2
assert metriques.mae == 3.0
def test_drift_reports_no_mape_when_every_actual_is_zero() -> None:
metriques = mesure([paire(prevu=1.0, reel=0.0)])
assert metriques.mape is None
def test_drift_lists_every_model_reference_seen_in_the_window() -> None:
metriques = mesure(
[paire(prevu=10.0, reel=10.0, reference="lightgbm-bbb"), paire(prevu=10.0, reel=10.0)]
)
assert metriques.model_references == ["lightgbm-aaa", "lightgbm-bbb"]
async def test_drift_reports_indetermine_when_the_window_holds_too_few_observations() -> None:
depot = FauxDepot(recentes=paires(nombre=3, prevu=10.0, reel=10.0))
rapports = await service(depot, min_observations=24).evaluate(now=INSTANT)
assert {rapport.status for rapport in rapports} == {STATUT_INDETERMINE}
assert all(rapport.reason for rapport in rapports)
async def test_drift_reports_derive_when_the_recent_mae_exceeds_the_reference_ratio() -> None:
depot = FauxDepot(
recentes=paires(nombre=30, prevu=14.0, reel=10.0),
anciennes=paires(nombre=30, prevu=11.0, reel=10.0),
comptages=[ComptageStatut(site_id="SITE001", status="available", nombre=30)],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
global_ = next(rapport for rapport in rapports if rapport.site_id is None)
assert global_.status == STATUT_DERIVE
assert global_.mae == 4.0
assert global_.reference_mae == 1.0
async def test_drift_reports_stable_when_the_recent_mae_stays_close_to_the_reference() -> None:
depot = FauxDepot(
recentes=paires(nombre=30, prevu=11.0, reel=10.0),
anciennes=paires(nombre=30, prevu=11.0, reel=10.0),
comptages=[ComptageStatut(site_id="SITE001", status="available", nombre=30)],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
global_ = next(rapport for rapport in rapports if rapport.site_id is None)
assert global_.status == STATUT_STABLE
assert global_.reason is None
async def test_drift_reports_indetermine_when_the_reference_window_is_empty() -> None:
depot = FauxDepot(
recentes=paires(nombre=30, prevu=14.0, reel=10.0),
comptages=[ComptageStatut(site_id="SITE001", status="available", nombre=30)],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
global_ = next(rapport for rapport in rapports if rapport.site_id is None)
assert global_.status == STATUT_INDETERMINE
assert global_.reference_mae is None
async def test_drift_reports_derive_when_the_coverage_ratio_falls_under_the_threshold() -> None:
depot = FauxDepot(
recentes=paires(nombre=30, prevu=10.0, reel=10.0),
anciennes=paires(nombre=30, prevu=10.0, reel=10.0),
comptages=[ComptageStatut(site_id="SITE001", status="available", nombre=100)],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
global_ = next(rapport for rapport in rapports if rapport.site_id is None)
assert global_.status == STATUT_DERIVE
assert global_.coverage_ratio == 0.3
async def test_drift_reports_one_line_per_site_and_one_global_line() -> None:
depot = FauxDepot(
recentes=[
*paires(site_id="SITE001", nombre=12, prevu=10.0, reel=10.0),
*paires(site_id="SITE002", nombre=12, prevu=10.0, reel=10.0),
],
comptages=[
ComptageStatut(site_id="SITE001", status="available", nombre=12),
ComptageStatut(site_id="SITE002", status="available", nombre=12),
],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
assert [rapport.site_id for rapport in rapports] == ["SITE001", "SITE002", None]
assert next(r for r in rapports if r.site_id is None).n_observations == 24
async def test_drift_still_reports_a_site_that_stopped_being_scored() -> None:
depot = FauxDepot(
recentes=paires(site_id="SITE001", nombre=30, prevu=10.0, reel=10.0),
anciennes=[
*paires(site_id="SITE001", nombre=30, prevu=10.0, reel=10.0),
*paires(site_id="SITE002", nombre=30, prevu=10.0, reel=10.0),
],
comptages=[ComptageStatut(site_id="SITE001", status="available", nombre=30)],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
disparu = next(rapport for rapport in rapports if rapport.site_id == "SITE002")
assert disparu.status == STATUT_INDETERMINE
assert disparu.n_observations == 0
async def test_drift_measures_the_share_of_sites_left_without_enough_history() -> None:
depot = FauxDepot(
recentes=paires(nombre=30, prevu=10.0, reel=10.0),
comptages=[
ComptageStatut(site_id="SITE001", status="available", nombre=30),
ComptageStatut(site_id="SITE001", status="insufficient_data", nombre=10),
],
)
rapports = await service(depot, min_observations=10).evaluate(now=INSTANT)
assert next(r for r in rapports if r.site_id is None).insufficient_data_ratio == 0.25
async def test_drift_closes_the_window_before_the_grace_delay() -> None:
depot = FauxDepot()
await service(depot, grace=timedelta(hours=2), fenetre=timedelta(hours=168)).evaluate(
now=INSTANT
)
recente, reference = depot.fenetres
assert recente[1] == INSTANT - timedelta(hours=2)
assert recente[0] == INSTANT - timedelta(hours=170)
assert reference[1] == recente[0]
async def test_drift_truncates_the_reference_instant_to_the_hour() -> None:
premier = FauxDepot()
second = FauxDepot()
await service(premier).evaluate(now=INSTANT + timedelta(minutes=17, microseconds=3))
await service(second).evaluate(now=INSTANT + timedelta(minutes=48))
assert premier.fenetres[0] == second.fenetres[0]
@pytest.mark.parametrize(
("prevu", "attendu"),
[(10.0, STATUT_STABLE), (30.0, STATUT_DERIVE)],
ids=["mae_stable", "mae_triplee"],
)
async def test_drift_compares_the_recent_window_to_the_reference_one(
prevu: float, attendu: str
) -> None:
depot = FauxDepot(
recentes=paires(nombre=30, prevu=prevu, reel=10.0),
anciennes=paires(nombre=30, prevu=10.0, reel=10.0),
comptages=[ComptageStatut(site_id="SITE001", status="available", nombre=30)],
)
rapports = await service(depot, min_observations=10, mae_plancher=1.0).evaluate(now=INSTANT)
assert next(r for r in rapports if r.site_id is None).status == attendu
-223
View File
@@ -1,223 +0,0 @@
"""Piege : ce fichier porte le marqueur `chaine`, pas `integration` - test_the_ml_binaries...()
Il lance les vrais binaires `enervision_ml.train` et `enervision_ml.score` dans l'environnement
uv de `ml/`, que le job `integration` de `backend.yml` n'installe pas. Un marqueur distinct evite
que ce job, et `make test`, ne le selectionnent et n'echouent faute de `ml/.venv`.
"""
import math
import os
import subprocess
from collections.abc import AsyncIterator, Iterator
from dataclasses import dataclass, field
from datetime import UTC, datetime, timedelta
from functools import partial
from pathlib import Path
from typing import Any
from uuid import uuid4
import anyio
import pytest
from fastapi import FastAPI
from httpx import AsyncClient
from sqlalchemy import delete, insert, make_url
from app.api.deps import get_current_principal
from app.core.config import get_settings
from app.core.principal import Principal
from app.core.roles import AccountKind, Role
from app.db.session import get_session_factory
from app.models.energy import Prediction, Reading, Site
pytestmark = pytest.mark.chaine
RACINE = Path(__file__).resolve().parents[3]
ML = RACINE / "ml"
PYTHON_ML = Path(os.environ.get("ML_PYTHON", ML / ".venv" / "bin" / "python"))
HEURES_COMPLETES = 400
HEURES_INSUFFISANTES = 100
def lecteur() -> Principal:
return Principal(
id=uuid4(),
email="lecteur@enervision.fr",
role=Role.LECTEUR,
kind=AccountKind.HUMAIN,
must_change_password=False,
)
def url_ml() -> str:
"""Derive la chaine du pipeline de celle du backend plutot que de la recopier : les deux
cotes visent ainsi la meme base, dans leur dialecte respectif."""
return (
make_url(get_settings().database_url)
.set(drivername="postgresql+psycopg")
.render_as_string(hide_password=False)
)
def lance_ml(module: str, *arguments: str, journal: Path) -> subprocess.CompletedProcess[str]:
if not PYTHON_ML.exists():
pytest.fail(
f"Environnement ml/ absent ({PYTHON_ML}). Lancer `cd ml && uv sync --all-groups`."
)
return subprocess.run( # noqa: S603 -- argv en liste, sans shell, binaire resolu dans le depot
[str(PYTHON_ML), "-m", module, *arguments],
cwd=ML,
text=True,
capture_output=True,
timeout=600,
check=False,
env={
**os.environ,
"ML_DATABASE_URL": url_ml(),
"MLFLOW_TRACKING_URI": f"sqlite:///{journal}/mlflow.db",
},
)
async def executer(module: str, *arguments: str, journal: Path) -> subprocess.CompletedProcess[str]:
resultat = await anyio.to_thread.run_sync(
partial(lance_ml, module, *arguments, journal=journal)
)
assert resultat.returncode == 0, resultat.stderr
return resultat
@dataclass
class Parc:
sites: list[str] = field(default_factory=list)
def lignes_horaires(site_id: str, *, heures: int, fin: datetime) -> list[dict[str, Any]]:
return [
{
"site_id": site_id,
"timestamp": fin - timedelta(hours=decalage),
"source": "api_history",
"consumption_kwh": 50.0 + math.sin(decalage / 12.0) * 10.0,
"temperature_celsius": 15.0,
"humidity_percent": 50.0,
"solar_irradiance_wm2": 0.0,
"is_working_hours": True,
"raw_data": {},
}
for decalage in reversed(range(heures))
]
@pytest.fixture
async def parc() -> AsyncIterator[Parc]:
"""Deux sites dotes d'un historique complet, un troisieme qui n'atteint pas le lag de 168 h.
Les ecritures sont validees : les binaires ML ouvrent leur propre connexion et ne verraient
pas une transaction en cours.
"""
fin = datetime.now(UTC).replace(minute=0, second=0, microsecond=0) - timedelta(hours=1)
marque = uuid4().hex[:12]
complets = [f"TEST-{marque}-A", f"TEST-{marque}-B"]
partiel = f"TEST-{marque}-C"
parc = Parc(sites=[*complets, partiel])
async with get_session_factory()() as session:
await session.execute(
insert(Site),
[
{
"site_id": site_id,
"site_name": f"Site {site_id}",
"site_type": "office",
"capacity_kw": 100.0,
}
for site_id in parc.sites
],
)
for site_id in complets:
await session.execute(
insert(Reading), lignes_horaires(site_id, heures=HEURES_COMPLETES, fin=fin)
)
await session.execute(
insert(Reading), lignes_horaires(partiel, heures=HEURES_INSUFFISANTES, fin=fin)
)
await session.commit()
try:
yield parc
finally:
async with get_session_factory()() as session:
await session.execute(delete(Prediction).where(Prediction.site_id.in_(parc.sites)))
await session.execute(delete(Reading).where(Reading.site_id.in_(parc.sites)))
await session.execute(delete(Site).where(Site.site_id.in_(parc.sites)))
await session.commit()
@pytest.fixture
def principal_lecteur(app: FastAPI) -> Iterator[None]:
app.dependency_overrides[get_current_principal] = lecteur
yield
app.dependency_overrides.pop(get_current_principal, None)
async def resume_du_site(client: AsyncClient, site_id: str) -> dict[str, Any]:
reponse = await client.get("/api/v1/predictions")
assert reponse.status_code == 200
sites = reponse.json()["sites"]
return next(site for site in sites if site["site_id"] == site_id)
async def entraine_et_score(parc: Parc, tmp_path: Path, *arguments: str) -> Path:
modele = tmp_path / "lightgbm-consumption.txt"
await executer(
"enervision_ml.train",
"--model-output",
str(modele),
"--mlflow-tracking-uri",
f"sqlite:///{tmp_path}/mlflow.db",
journal=tmp_path,
)
await executer("enervision_ml.score", "--model", str(modele), *arguments, journal=tmp_path)
return modele
async def test_the_ml_binaries_produce_a_prediction_that_the_api_serves(
parc: Parc, tmp_path: Path, client: AsyncClient, principal_lecteur: None
) -> None:
await entraine_et_score(parc, tmp_path)
servi = await resume_du_site(client, parc.sites[0])
assert servi["prediction"]["status"] == "available"
assert servi["prediction"]["predicted_value"] is not None
assert servi["prediction"]["target_metric"] == "consumption_kwh"
async def test_the_api_exposes_the_failure_reason_of_a_site_without_enough_history(
parc: Parc, tmp_path: Path, client: AsyncClient, principal_lecteur: None
) -> None:
await entraine_et_score(parc, tmp_path)
servi = await resume_du_site(client, parc.sites[-1])
assert servi["prediction"]["status"] == "insufficient_data"
assert servi["prediction"]["predicted_value"] is None
assert servi["prediction"]["failure_reason"] is not None
async def test_the_api_serves_the_latest_run_when_the_score_cli_runs_twice(
parc: Parc, tmp_path: Path, client: AsyncClient, principal_lecteur: None
) -> None:
modele = await entraine_et_score(parc, tmp_path)
premier = await resume_du_site(client, parc.sites[0])
await executer("enervision_ml.score", "--model", str(modele), journal=tmp_path)
second = await resume_du_site(client, parc.sites[0])
assert second["prediction"]["created_at"] > premier["prediction"]["created_at"]
assert second["prediction"]["model_reference"] == premier["prediction"]["model_reference"]
-108
View File
@@ -1,108 +0,0 @@
from datetime import UTC, datetime, timedelta
import pytest
from app.monitoring import drift as cli
from app.repositories.drift import NouveauRapportDerive
from app.services.drift import STATUT_DERIVE, STATUT_STABLE, Seuils
INSTANT = datetime(2026, 9, 22, 12, tzinfo=UTC)
def rapport(*, site_id: str | None, status: str, reason: str | None = None) -> NouveauRapportDerive:
return NouveauRapportDerive(
site_id=site_id,
window_start=INSTANT - timedelta(hours=168),
window_end=INSTANT,
reference_start=None,
reference_end=None,
n_observations=48,
mae=1.5,
mape=12.0,
bias=0.3,
reference_mae=1.2,
coverage_ratio=1.0,
insufficient_data_ratio=0.0,
model_references=["lightgbm-aaa"],
status=status,
reason=reason,
)
def installe(monkeypatch: pytest.MonkeyPatch, rapports: list[NouveauRapportDerive]) -> None:
async def fausse_execution(
*, now: datetime | None, site_id: str | None, seuils: Seuils | None
) -> list[NouveauRapportDerive]:
return rapports
monkeypatch.setattr(cli, "run_drift", fausse_execution)
def test_parse_args_defaults_to_the_standard_window() -> None:
arguments = cli.parse_args([])
assert arguments.window_hours == 168
assert arguments.grace_hours == 2
assert arguments.fail_on_drift is False
def test_parse_args_reads_the_site_id() -> None:
assert cli.parse_args(["--site-id", "SITE001"]).site_id == "SITE001"
def test_parse_args_parses_the_instant_option() -> None:
arguments = cli.parse_args(["--now", "2026-09-22T12:00:00+00:00"])
assert arguments.now == INSTANT
def test_parse_instant_treats_a_naive_datetime_as_utc() -> None:
assert cli._parse_instant("2026-09-22T12:00:00") == INSTANT
def test_seuils_depuis_translates_the_hour_options_into_durations() -> None:
seuils = cli.seuils_depuis(cli.parse_args(["--window-hours", "24", "--grace-hours", "1"]))
assert seuils.fenetre == timedelta(hours=24)
assert seuils.grace == timedelta(hours=1)
def test_main_prints_the_verdict_of_every_line(
monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str]
) -> None:
installe(
monkeypatch,
[
rapport(site_id="SITE001", status=STATUT_STABLE),
rapport(site_id=None, status=STATUT_STABLE),
],
)
code = cli.main([])
sortie = capsys.readouterr().out
assert code == 0
assert "SITE001" in sortie
assert "TOUS SITES" in sortie
def test_main_exits_non_zero_when_drift_is_detected_and_the_flag_is_set(
monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str]
) -> None:
installe(monkeypatch, [rapport(site_id=None, status=STATUT_DERIVE, reason="MAE doublée")])
code = cli.main(["--fail-on-drift"])
assert code == 1
assert "MAE doublée" in capsys.readouterr().out
def test_main_exits_zero_when_drift_is_detected_without_the_flag(
monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str]
) -> None:
installe(monkeypatch, [rapport(site_id=None, status=STATUT_DERIVE, reason="MAE doublée")])
code = cli.main([])
assert code == 0
assert capsys.readouterr().out != ""
+6 -26
View File
@@ -90,25 +90,14 @@ consommation prévue de **l'heure suivant sa dernière lecture connue**, et écr
### Ce que le run écrit, et ce qu'il n'écrase pas ### Ce que le run écrit, et ce qu'il n'écrase pas
La table `prediction` **n'a pas de contrainte d'unicité sur `(site_id, target_at)`** : chaque run La table `prediction` **n'a pas de contrainte d'unicité sur `(site_id, target_at)`** : chaque run
insère une ligne de plus au lieu d'écraser la précédente. C'est délibéré, et c'est ce qui rend insère une ligne de plus au lieu d'écraser la précédente. C'est délibéré, et c'est ce qui rendra
possible la comparaison prévision contre réalisé. La surveillance de dérive s'en sert : elle possible la comparaison prévision contre réalisé, donc la surveillance de dérive (#44, #45), qui
retient, pour chaque `(site_id, target_at)`, la ligne du run le plus récent, celle-là même que n'existe pas encore.
sert `GET /api/v1/predictions`. Voir l'[ADR 0011](adr/0011-surveillance-de-derive-dans-le-backend.md).
Trois contraintes de cohérence sont portées par la base et non par le code applicatif : Trois contraintes de cohérence sont portées par la base et non par le code applicatif :
`status = 'available'` exige une `predicted_value` et interdit un `failure_reason` ; `status = 'available'` exige une `predicted_value` et interdit un `failure_reason` ;
`insufficient_data` et `error` exigent l'inverse ; `target_metric` est bornée à `insufficient_data` et `error` exigent l'inverse ; `target_metric` est bornée à
`consumption_kwh` ou `consumption_kw`, et la forme énergie impose une `period_minutes`. Elles `consumption_kwh` ou `consumption_kw`, et la forme énergie impose une `period_minutes`.
sont vérifiées depuis le code qui écrit par `ml/tests/test_score_integration.py`, sur une vraie
base : un double ne prouverait rien d'une contrainte SQL.
**`--now` borne la fenêtre des deux côtés.** `load_recent_from_database` exige un `until` autant
qu'un `since`, et le scoring lui passe l'instant de référence. Sans cette borne haute,
`build_scoring_frame` repartait de la dernière lecture de toute la table quelle que soit la valeur
demandée : `target_at` valait toujours « fin du jeu + 1 h », et l'âge de la dernière lecture
devenait négatif sans franchir le seuil de péremption. Rejouer le scoring sur des instants passés
produit désormais des prévisions dont le réalisé existe déjà, ce dont la surveillance de dérive a
besoin pour se démontrer sur un jeu figé.
### `model_reference` est un hachage, pas un nom de fichier ### `model_reference` est un hachage, pas un nom de fichier
@@ -153,10 +142,6 @@ flowchart LR
train -- "models/*.txt + run MLflow" --> score train -- "models/*.txt + run MLflow" --> score
score -- "INSERT" --> prediction score -- "INSERT" --> prediction
prediction -- "lecture seule" --> route prediction -- "lecture seule" --> route
prediction -- "prévu" --> derive["app.monitoring.drift<br/>écart prévu / réalisé"]
reading -- "réalisé" --> derive
derive -- "INSERT" --> rapport[("drift_report")]
rapport -- "lecture seule" --> monitoring["GET /api/v1/monitoring/drift"]
``` ```
**La règle, en une phrase : FastAPI ne fait jamais tourner LightGBM.** **La règle, en une phrase : FastAPI ne fait jamais tourner LightGBM.**
@@ -178,12 +163,8 @@ flowchart LR
Le corollaire est qu'il n'y a **aucune prévision à la demande** : la fraîcheur d'une prévision est Le corollaire est qu'il n'y a **aucune prévision à la demande** : la fraîcheur d'une prévision est
celle du dernier run de scoring. Ce run est ordonnancé par Airflow, DAG `ml_score` en `@hourly` celle du dernier run de scoring. Ce run est ordonnancé par Airflow, DAG `ml_score` en `@hourly`
(issue #115) ; seuls le mode `--csv` et un lancement local restent manuels, tout comme (issue #115) ; seuls le mode `--csv` et un lancement local restent manuels, tout comme
l'entraînement, dont le DAG `ml_train` n'a pas de planification. l'entraînement, dont le DAG `ml_train` n'a pas de planification. La dette qui subsiste est la
surveillance de dérive, portée par les issues #44 et #45.
La surveillance de dérive traverse cette frontière **dans le sens de la table vers le backend**,
sans la percer : elle relit `prediction` et `reading` en SQL, ne charge aucun modèle, et n'appelle
pas MLflow. Son calcul, son seuil et son refus de comparer à la métrique d'entraînement sont dans
l'[ADR 0011](adr/0011-surveillance-de-derive-dans-le-backend.md).
--- ---
@@ -194,4 +175,3 @@ l'[ADR 0011](adr/0011-surveillance-de-derive-dans-le-backend.md).
- [ADR 0006](adr/0006-moteur-de-regles-dans-le-backend.md) : ce qui consomme les prédictions - [ADR 0006](adr/0006-moteur-de-regles-dans-le-backend.md) : ce qui consomme les prédictions
- [`architecture/20-backend.md`](architecture/20-backend.md) : le contrat de `GET /predictions` - [`architecture/20-backend.md`](architecture/20-backend.md) : le contrat de `GET /predictions`
- [`architecture/40-data.md`](architecture/40-data.md) : le modèle de données - [`architecture/40-data.md`](architecture/40-data.md) : le modèle de données
- [ADR 0011](adr/0011-surveillance-de-derive-dans-le-backend.md) : la surveillance de dérive
-1
View File
@@ -17,4 +17,3 @@
| [0008](adr/0008-airflow-execute-le-code-du-backend.md) | Airflow exécute le code du backend en sous-processus, dans son propre environnement | | [0008](adr/0008-airflow-execute-le-code-du-backend.md) | Airflow exécute le code du backend en sous-processus, dans son propre environnement |
| [0009](adr/0009-deux-environnements-compose-sur-la-vm-eni.md) | Deux environnements sur la VM ENI, un projet Compose chacun, déployés par un runner auto-hébergé | | [0009](adr/0009-deux-environnements-compose-sur-la-vm-eni.md) | Deux environnements sur la VM ENI, un projet Compose chacun, déployés par un runner auto-hébergé |
| [0010](adr/0010-terraform-provisionne-github-actions-deploie.md) | Terraform provisionne la machine, GitHub Actions déploie l'application | | [0010](adr/0010-terraform-provisionne-github-actions-deploie.md) | Terraform provisionne la machine, GitHub Actions déploie l'application |
| [0011](adr/0011-surveillance-de-derive-dans-le-backend.md) | La surveillance de dérive vit dans le backend et écrit sa propre table |
@@ -1,120 +0,0 @@
# 0011 - La surveillance de dérive vit dans le backend et écrit sa propre table
- Statut : accepté
- Date : 2026-09-22
## Contexte
L'issue #45 demande des tests d'intégration API ↔ DB ↔ ML. Trois documents du dépôt annoncent
par ailleurs, depuis le jalon J3, une surveillance de dérive qui n'existe nulle part :
`docs/architecture/00-vue-ensemble.md` (« Surveillance de dérive (EC06, #44/#45) pas encore
construite »), `docs/ML-START.md` (« la dette qui subsiste est la surveillance de dérive »), et
le docstring de `write_predictions()` dans `ml/enervision_ml/score.py`, qui justifie l'absence
d'unicité sur `(site_id, target_at)` par la comparaison future entre prévu et réalisé.
La matière première est en base : `prediction` porte ce que le modèle a annoncé, `reading` ce
qui est réellement arrivé. Restaient trois questions : où vit le calcul, à quoi on compare, et
où atterrit le résultat.
## Décision
**Le calcul vit dans `apps/backend`** : `repositories/drift.py` pour le SQL, `services/drift.py`
pour la logique, `monitoring/drift.py` pour la CLI, `api/v1/endpoints/monitoring.py` pour la
lecture. Le dossier `ml/` ne gagne pas une ligne.
**Le résultat est persisté** dans une table `drift_report`, une ligne par site plus une ligne
globale que `site_id` à NULL désigne.
**La comparaison oppose deux fenêtres vives de 168 h**, la récente et celle qui la précède, et
le verdict a trois valeurs : `stable`, `derive`, `indetermine`.
### Pourquoi le backend, alors que le sujet est le modèle
- **`prediction` n'est pas dans le périmètre de `ML_DATABASE_URL`.** `enervision_ml/config.py`,
`docs/ML-START.md` et l'[ADR 0003](0003-autorisation-rbac-a-trois-roles.md) désignent pour
cette variable un rôle PostgreSQL restreint **en lecture sur `reading` et `site`**. Mettre la
dérive dans `ml/` obligerait à élargir ce rôle à `prediction`, et à l'écriture : ce serait
contredire par le code la dette de moindre privilège que ces trois documents ont posée par
écrit.
- **L'alignement prévu contre réalisé existe déjà ici, une fois.** `AlertService._detect_anomaly`
croise `reading` et `prediction` sur le même instant, et `PredictionRepository.list_since`
porte déjà le piège des runs empilés. Le réécrire en SQL brut dans `ml/` créerait une seconde
source de vérité sur « quelle prédiction correspond à quelle lecture », ce que
l'[ADR 0006](0006-moteur-de-regles-dans-le-backend.md) a déjà refusé pour les règles.
- **La frontière de `docs/ML-START.md` tient.** FastAPI ne fait toujours pas tourner LightGBM :
la dérive lit deux tables et compare des nombres, elle n'évalue aucun modèle.
**Conséquence assumée** : `enervision_ml.metrics.regression_metrics` n'est pas réutilisable, le
backend n'important pas `enervision_ml`. MAE, MAPE et biais sont donc réécrits, une quinzaine de
lignes. Cette duplication n'est pas celle que `build_features` interdit : une divergence de
features est silencieuse et ruine les prévisions sans erreur, une divergence sur une moyenne
d'écarts absolus est attrapée par le premier test à valeurs connues.
### Ce qu'on mesure, et les deux dédoublonnages obligatoires
La paire est `prediction ⋈ reading` sur `(site_id, target_at = timestamp)`, restreinte aux
prédictions `available`. Elle exige un `DISTINCT ON` **des deux côtés** :
- `prediction` n'a pas d'unicité sur `(site_id, target_at)`, chaque run de scoring empile une
ligne. On retient la plus récente, celle que sert `GET /api/v1/predictions`, départagée par
`prediction_id` : `created_at` vaut l'heure de début de transaction et ne distingue pas deux
lignes du même run.
- `uq_reading_source` autorise deux lectures au même instant quand la `source` diffère. Sans
dédoublonnage, la jointure compterait cette heure deux fois et pondérerait doublement le site.
La fenêtre est **fermée à droite par un délai de grâce de 2 h** : le réalisé de la dernière
heure n'est pas encore ingéré, et l'inclure ferait chuter le taux de couverture à chaque
exécution, pour une raison qui n'a rien à voir avec le modèle.
Métriques retenues : `mae` (la métrique même qu'optimise LightGBM), **`bias` signé** (une MAE qui
monte dit « moins bon », un biais qui s'éloigne de zéro dit « le modèle se trompe toujours du
même côté », signature d'un décalage de distribution), `mape`, `n_observations`,
`coverage_ratio` et `insufficient_data_ratio` (qui mesurent le pipeline, pas le modèle), et la
liste des `model_references` vus dans la fenêtre : une MAE qui saute à l'instant exact où le
modèle change n'est pas une dérive, c'est une régression de réentraînement.
## Alternatives écartées
| Écartée | Raison |
|---|---|
| Comparer à la métrique MLflow de l'entraînement | Ce ne sont pas les mêmes grandeurs : `train.py` mesure un backtest où la météo de l'heure cible est connue, le scoring prévoit une heure future dont la météo est `NaN` et dont `is_working_hours` est recopié. Le verdict serait « dérive » dès le premier jour. Et le backend devrait importer `mlflow`, ce que la frontière de ML-START interdit. |
| Écrire le résultat dans `alert` | `ck_alert_source` et `ck_alert_type` bornent les valeurs autorisées, `alert.site_id` est `NOT NULL` et n'accueillerait donc pas la ligne globale, et toute alerte est ensuite relue par le moteur de recommandations, qui devrait apprendre une règle qui ne le concerne pas (ADR 0006). |
| Une jauge Prometheus | `monitoring/` ne contient que des `.gitkeep` et aucun collecteur ne lit `/metrics` : une jauge que personne ne scrute n'est pas une preuve. Le calcul est de surcroît un traitement par lot, pas le processus qui sert l'API : la jauge disparaîtrait avec lui. |
| Ne rien persister, journaliser seulement | La question posée à un jury est « comment savez-vous que le modèle se dégrade ? ». La réponse est une série dans le temps, pas une ligne de journal perdue avec le conteneur. Sans ligne écrite, l'endpoint n'a rien à lire et le test d'intégration rien à vérifier. |
| Une tâche de plus dans le DAG `alertes` | La fenêtre fait 168 h : la recalculer chaque heure écrirait vingt-quatre lignes par jour pour un verdict qui ne bouge pas à cette cadence. Surtout, un échec de dérive ferait rougir `alertes` et laisserait croire que la détection a échoué. |
## Conséquences
- Une migration ajoute `drift_report`. Son idempotence passe par un **index unique à
`coalesce(site_id, '')`** et non par une `UniqueConstraint` : deux lignes globales ont toutes
deux `site_id` à NULL, et NULL n'est égal à rien, pas même à lui-même. Même forme que
`uq_reading_source`. Cet index n'a de sens que parce que `evaluate()` **tronque son instant de
référence à l'heure** : avec les microsecondes de `now()`, deux exécutions ne porteraient jamais
la même clé et l'index ne dédoublonnerait rien.
- **Sans fenêtre de référence, le verdict est `indetermine`, pas `stable`.** Au premier
lancement, et après tout trou d'ingestion de plus de 168 h, il n'y a rien à quoi comparer :
annoncer `stable` serait affirmer ce que la donnée ne dit pas.
- **Un site présent dans la fenêtre de référence et absent de la récente reçoit sa ligne**, à
zéro observation. Un site qui cesse d'être scoré est exactement la panne que cette surveillance
existe pour dire : le taire en ne produisant aucune ligne serait l'inverse du besoin.
- `GET /api/v1/monitoring/drift` est réservé à partir du rôle `operateur` : c'est l'opérateur
qui agit sur un pipeline dégradé, pas l'administrateur de comptes. La route est classée dans
`tests/api/acces.py`, donc couverte gratuitement par la matrice de rôles rejouée avec de vrais
jetons.
- Un DAG `derive` quotidien l'ordonnance, sans reprise : rejouer une dérive la redéclarerait à
l'identique.
- La CLI sort en code non nul sous `--fail-on-drift` seulement. Par défaut, constater une dérive
n'est pas un échec d'exécution.
## Effet de bord assumé sur le pipeline
La dérive n'a de matière que si des paires prévu/réalisé existent. Or `enervision_ml.score --now`
ne rejouait pas l'historique : `load_recent_from_database` n'avait pas de borne haute et
`build_scoring_frame` repartait de la dernière lecture connue, si bien que `target_at` valait
toujours « fin du jeu + 1 h » et que l'âge de la dernière lecture devenait négatif sans franchir
le seuil de péremption. Sur le jeu historique, figé au 31/12/2024, aucune boucle de rattrapage
n'aurait donc rien produit de vérifiable.
`until` est devenu obligatoire sur ce chargeur, et le scoring lui passe son instant de référence.
Le comportement en exploitation ne change pas, aucune lecture n'étant postérieure à l'heure
courante ; seul le rattrapage sur données passées devient possible.
+4 -4
View File
@@ -70,7 +70,7 @@ Le lien `front -.-> api` reste en pointillé : le frontend appelle bien une API,
intercepteur répond à sa place tant que les endpoints n'existent pas. Voir intercepteur répond à sa place tant que les endpoints n'existent pas. Voir
[30-frontend.md](30-frontend.md). [30-frontend.md](30-frontend.md).
Le lien `airflow --> db` est maintenant en trait plein : cinq DAGs tournent, deux pour Le lien `airflow --> db` est maintenant en trait plein : quatre DAGs tournent, deux pour
l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la
génération des recommandations (issue #116), et `historical_import` pour l'ingestion du dataset génération des recommandations (issue #116), et `historical_import` pour l'ingestion du dataset
historique (issue #119). L'orchestration de l'import API Mock et la réconciliation globale des historique (issue #119). L'orchestration de l'import API Mock et la réconciliation globale des
@@ -86,16 +86,16 @@ collecteur ne vient le lire.
| Backend | FastAPI, Python 3.14 | `apps/backend` | `En cours` | Factory, configuration, journalisation, 2 sondes de santé, `/metrics`, contrat OpenAPI versionné, routes `sites`, `alerts`, `recommendations`, `stats/summary`, `readings`, `sensors/status` et `predictions` en lecture (endpoints → services → repositories → models) | | Backend | FastAPI, Python 3.14 | `apps/backend` | `En cours` | Factory, configuration, journalisation, 2 sondes de santé, `/metrics`, contrat OpenAPI versionné, routes `sites`, `alerts`, `recommendations`, `stats/summary`, `readings`, `sensors/status` et `predictions` en lecture (endpoints → services → repositories → models) |
| Frontend | Angular 22, Node 24 | `apps/frontend` | `En cours` | Tableau de bord sur route `/dashboard`, authentification complète (garde de route, intercepteur de jeton), cinq services HTTP, graphiques Chart.js. `stats`/`alerts` sur fixtures, `predictions` branché sur l'API réelle | | Frontend | Angular 22, Node 24 | `apps/frontend` | `En cours` | Tableau de bord sur route `/dashboard`, authentification complète (garde de route, intercepteur de jeton), cinq services HTTP, graphiques Chart.js. `stats`/`alerts` sur fixtures, `predictions` branché sur l'API réelle |
| Base | PostgreSQL 17 + TimescaleDB | `db` | `Fait` | Bootstrap de l'extension, base de test, chaîne Alembic. Schéma applicatif créé (`site`, `dataset`, `reading` en hypertable, `prediction`, `alert`, `recommendation`) | | Base | PostgreSQL 17 + TimescaleDB | `db` | `Fait` | Bootstrap de l'extension, base de test, chaîne Alembic. Schéma applicatif créé (`site`, `dataset`, `reading` en hypertable, `prediction`, `alert`, `recommendation`) |
| ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`, orchestré par Airflow (`ml_train`/`ml_score`). Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../ML-START.md). Surveillance de dérive livrée côté backend (`app.monitoring.drift`, table `drift_report`, `GET /monitoring/drift`, DAG `derive`), voir [ADR 0011](../adr/0011-surveillance-de-derive-dans-le-backend.md) | | ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`, orchestré par Airflow (`ml_train`/`ml_score`). Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../ML-START.md). Surveillance de dérive (EC06, #44/#45) pas encore construite |
| Infra | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Provisionnement de la VM par Terraform, qui installe Docker, prépare les deux environnements et enregistre le runner, jamais appliqué ([ADR 0010](../adr/0010-terraform-provisionne-github-actions-deploie.md)). Module d'installation k3s jamais appliqué, aucune ressource Kubernetes déclarée | | Infra | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Provisionnement de la VM par Terraform, qui installe Docker, prépare les deux environnements et enregistre le runner, jamais appliqué ([ADR 0010](../adr/0010-terraform-provisionne-github-actions-deploie.md)). Module d'installation k3s jamais appliqué, aucune ressource Kubernetes déclarée |
| Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API | | Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API |
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Cinq DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes`, `historical_import` et `derive` (quotidien, surveillance de dérive). Le DAG historique orchestre `app.etl.historical_import` et charge `dataset`, `site` et `reading`. L'orchestration API Mock reste à compléter dans #15 | | ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Quatre DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes` et `historical_import`. Le DAG historique orchestre `app.etl.historical_import` et charge `dataset`, `site` et `reading`. L'orchestration API Mock reste à compléter dans #15 |
| CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 7 workflows, 19 jobs : lint, typage, tests avec seuil de couverture bloquant, tests d'intégration sur TimescaleDB réel, audit de dépendances, SAST Bandit, quality gate SonarCloud, intégrité des DAGs Airflow, formatage et validation du Terraform. Déploiement continu vers la VM ENI écrit par `deploy.yml`, `dev` en recette et `main` en production après approbation ([ADR 0009](../adr/0009-deux-environnements-compose-sur-la-vm-eni.md)), mais jamais exécuté : la machine n'est pas provisionnée et le runner n'y est pas enregistré. Détail dans [50-cicd.md](50-cicd.md) | | CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 7 workflows, 19 jobs : lint, typage, tests avec seuil de couverture bloquant, tests d'intégration sur TimescaleDB réel, audit de dépendances, SAST Bandit, quality gate SonarCloud, intégrité des DAGs Airflow, formatage et validation du Terraform. Déploiement continu vers la VM ENI écrit par `deploy.yml`, `dev` en recette et `main` en production après approbation ([ADR 0009](../adr/0009-deux-environnements-compose-sur-la-vm-eni.md)), mais jamais exécuté : la machine n'est pas provisionnée et le runner n'y est pas enregistré. Détail dans [50-cicd.md](50-cicd.md) |
## Flux bout en bout ## Flux bout en bout
Statut : `En cours`. **Le chemin de lecture tourne** : base, API et frontend. **Le chemin Statut : `En cours`. **Le chemin de lecture tourne** : base, API et frontend. **Le chemin
d'ingestion dessiné ci-dessous n'existe pas** : les DAGs livrés (`ml_train`, `ml_score`, d'ingestion dessiné ci-dessous n'existe pas** : les trois DAGs livrés (`ml_train`, `ml_score`,
issue #115 ; `alertes`, issue #116) orchestrent le pipeline ML et la détection d'alertes, pas issue #115 ; `alertes`, issue #116) orchestrent le pipeline ML et la détection d'alertes, pas
l'ingestion, qui reste lancée à la main par les scripts d'import (issues #15 et #16). l'ingestion, qui reste lancée à la main par les scripts d'import (issues #15 et #16).
-1
View File
@@ -86,7 +86,6 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md).
| `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` | | `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` |
| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` | | `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` |
| `historical_import` | manuelle | `app.etl.historical_import`, dans `/opt/backend/.venv` ; les fichiers de `data/raw` sont montés en lecture seule dans `/opt/data/raw` | | `historical_import` | manuelle | `app.etl.historical_import`, dans `/opt/backend/.venv` ; les fichiers de `data/raw` sont montés en lecture seule dans `/opt/data/raw` |
| `derive` | `30 5 * * *` | `app.monitoring.drift`, dans `/opt/backend/.venv` ; quotidien parce que sa fenêtre couvre 168 h, et sans reprise parce qu'une dérive n'est pas une panne passagère |
Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique. Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique.
Il reste manuel, car le dataset sert à initialiser l'environnement. Le montage Il reste manuel, car le dataset sert à initialiser l'environnement. Le montage
+8 -35
View File
@@ -12,11 +12,11 @@ Les quatre couches existent désormais, portées par l'authentification.
```mermaid ```mermaid
flowchart TB flowchart TB
ep["endpoints<br/>health, auth, users, sites, alerts,<br/>recommendations, stats, readings, sensors,<br/>predictions, monitoring"] ep["endpoints<br/>health, auth, users, sites, alerts,<br/>recommendations, stats, readings, sensors, predictions"]
sc["schemas<br/>Pydantic"] sc["schemas<br/>Pydantic"]
sv["services<br/>AuthService, UserService,<br/>SiteService, AlertService, RecommendationService,<br/>StatsService, ReadingService, SensorService,<br/>PredictionService, DriftService"] sv["services<br/>AuthService, UserService,<br/>SiteService, AlertService, RecommendationService,<br/>StatsService, ReadingService, SensorService, PredictionService"]
rp["repositories<br/>user, refresh_token,<br/>login_attempt, audit_log,<br/>site, alert, recommendation, reading,<br/>prediction, drift"] rp["repositories<br/>user, refresh_token,<br/>login_attempt, audit_log,<br/>site, alert, recommendation, reading, prediction"]
md["models<br/>13 tables"] md["models<br/>10 tables"]
db[("PostgreSQL")] db[("PostgreSQL")]
ep --> sc ep --> sc
@@ -151,7 +151,6 @@ Deux fichiers d'environnement, deux usages : `.env` à la racine alimente `docke
| 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 |
| GET | `/api/v1/predictions` | Dernière prévision de consommation par site, calculée hors ligne par le pipeline de scoring (`ml/`). `lecteur` | 401, 403, 500 | | GET | `/api/v1/predictions` | Dernière prévision de consommation par site, calculée hors ligne par le pipeline de scoring (`ml/`). `lecteur` | 401, 403, 500 |
| GET | `/api/v1/monitoring/drift` | Dernier rapport de dérive par site, plus la ligne globale. `operateur` | 401, 403, 422, 500 |
| GET | `/metrics` | Format Prometheus, hors du schéma. Jeton requis si `APP_METRICS_TOKEN` est posé | | | GET | `/metrics` | Format Prometheus, hors du schéma. Jeton requis si `APP_METRICS_TOKEN` est posé | |
| GET | `/docs`, `/redoc`, `/openapi.json` | Hors du schéma. Fermés en `staging` et en `prod` | | | GET | `/docs`, `/redoc`, `/openapi.json` | Hors du schéma. Fermés en `staging` et en `prod` | |
@@ -224,34 +223,6 @@ par exemple `limit` hors bornes). Un datetime sans fuseau dans `start`/`end` est
l'UTC plutôt que rejeté : le comparer tel quel à `reading.timestamp` (`timestamptz`) échouerait l'UTC plutôt que rejeté : le comparer tel quel à `reading.timestamp` (`timestamptz`) échouerait
côté pilote, en `500` plutôt qu'un refus propre. côté pilote, en `500` plutôt qu'un refus propre.
### Surveillance de dérive
`DriftService.evaluate()` joint `prediction` et `reading` sur `(site_id, target_at = timestamp)`
et compare deux fenêtres vives de 168 h, la récente et celle qui la précède. Il rend une ligne par
site plus une ligne globale, que `DriftRepository.enregistre()` écrit dans `drift_report` avec
`ON CONFLICT DO NOTHING` sur `uq_drift_report_window`. L'instant de référence est tronqué à
l'heure, ce qui est la condition pour que cet index serve : rejouer la commande dans la même
heure n'ajoute rien.
| Métrique | Ce qu'elle dit |
|---|---|
| `mae` | Erreur moyenne en kWh, la métrique même qu'optimise LightGBM |
| `bias` | Erreur moyenne **signée** : c'est elle qui distingue un modèle plus bruyant d'un modèle qui se trompe systématiquement du même côté |
| `mape` | Comparable entre sites de tailles différentes, hors réalisés nuls |
| `coverage_ratio` | Part des prévisions disponibles qui ont trouvé leur réalisé : mesure le pipeline, pas le modèle |
| `insufficient_data_ratio` | Part des sites privés d'historique suffisant |
| `model_references` | Les modèles vus dans la fenêtre : une MAE qui saute à l'instant où le modèle change est une régression de réentraînement, pas une dérive |
Le verdict a trois valeurs, `stable`, `derive` et `indetermine` : sous un nombre minimal
d'observations, ou faute de fenêtre de référence à laquelle comparer, le service dit qu'il ne
sait pas plutôt que de rendre un chiffre trompeur. Un site qui figure dans la fenêtre de
référence mais plus dans la récente reçoit sa ligne à zéro observation : cesser d'être scoré est
la panne que cette surveillance existe pour dire. La
fenêtre est fermée à droite par un délai de grâce de 2 h, le temps que l'ingestion livre le
réalisé de la dernière heure. `python -m app.monitoring.drift` l'exécute, le DAG `derive`
l'ordonnance, et `GET /api/v1/monitoring/drift` sert le dernier rapport de chaque site. Les
arbitrages sont dans l'[ADR 0011](../adr/0011-surveillance-de-derive-dans-le-backend.md).
### Détection d'alertes internes ### Détection d'alertes internes
`AlertService` n'est plus lecture seule : `AlertService.detect()` compare les `reading` (et, pour `AlertService` n'est plus lecture seule : `AlertService.detect()` compare les `reading` (et, pour
@@ -360,8 +331,10 @@ pas prise :
| `license_info` | Aucune licence n'est choisie | | `license_info` | Aucune licence n'est choisie |
| `contact` | Aucun canal de support n'existe | | `contact` | Aucun canal de support n'existe |
Deux schémas de sécurité sont déclarés : `Jeton d'accès` pour le porteur JWT, et Deux schémas de sécurité sont déclarés : `JetonAcces` pour le porteur JWT, et
`Cookie de rafraîchissement` pour `/auth/refresh` et `/auth/logout`. **Le second est purement `CookieRafraichissement` pour `/auth/refresh` et `/auth/logout`, des noms ASCII délibérés (issue
#41 : un outillage tiers comme ZAP peut mal analyser un nom de schéma accentué dans le contrat).
**Le second est purement
documentaire** : son `auto_error=False` garantit qu'il ne décide d'aucun refus. Le passer à vrai documentaire** : son `auto_error=False` garantit qu'il ne décide d'aucun refus. Le passer à vrai
ferait répondre 403 avant d'atteindre `lit_le_cookie()`, et `/auth/refresh` cesserait de rendre le ferait répondre 403 avant d'atteindre `lit_le_cookie()`, et `/auth/refresh` cesserait de rendre le
401 sur lequel le frontend déclenche sa déconnexion. 401 sur lequel le frontend déclenche sa déconnexion.
+3 -12
View File
@@ -48,8 +48,7 @@ Statut : `Fait`.
et refuse de s'appliquer si l'extension TimescaleDB manque. et refuse de s'appliquer si l'extension TimescaleDB manque.
- Les révisions suivantes créent les tables liées à l'authentification : - Les révisions suivantes créent les tables liées à l'authentification :
`app_user`, `login_attempt`, `audit_log` et `refresh_token`. `app_user`, `login_attempt`, `audit_log` et `refresh_token`.
- La révision `e6d2026091501` crée six des sept tables Data et déclare l'hypertable `reading`. - La révision `e6d2026091501` crée les six tables Data et déclare l'hypertable `reading`.
- La révision `d3f1a2b7c904` ajoute `drift_report`, la septième.
- La révision `c0adab96238c` ajoute les tables `password_reset_attempt` - La révision `c0adab96238c` ajoute les tables `password_reset_attempt`
et `password_reset_token`. et `password_reset_token`.
@@ -262,8 +261,8 @@ Cette modélisation prend en compte :
- leurs métadonnées JSON ; - leurs métadonnées JSON ;
- les données de l'API Mock. - les données de l'API Mock.
Elle comprend sept tables Data, depuis le stockage des mesures jusqu'aux recommandations Elle comprend six tables Data, depuis le stockage des mesures jusqu'aux recommandations proposées
proposées à l'utilisateur, et jusqu'au suivi de la dérive du modèle. à l'utilisateur.
### Schéma de données ### Schéma de données
@@ -287,18 +286,11 @@ Chaque table remplit un rôle précis dans le traitement et l'exploitation des d
| `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 |
| `drift_report` | Suivre l'écart entre prévisions et réalisé, par site et tous sites confondus | Surveillance de dérive d'EnerVision |
Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`. Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`.
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles. Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
Les lignes de `drift_report` sont écrites par `app.monitoring.drift`, ordonnancé par le DAG
`derive`. Une ligne dont le `site_id` est `NULL` porte le résultat global, tous sites confondus :
c'est pourquoi l'unicité passe par un index sur `coalesce(site_id, '')` et non par une contrainte,
qui ne dédoublonnerait jamais deux lignes globales. Le calcul, ses seuils et ce qu'il refuse de
comparer sont dans l'[ADR 0011](../adr/0011-surveillance-de-derive-dans-le-backend.md).
Les lignes de `recommendation` sont écrites par le moteur de règles du backend Les lignes de `recommendation` sont écrites par le moteur de règles du backend
(`app/services/recommendation_rules.py`), déclenché par `POST /api/v1/recommendations/generate`, (`app/services/recommendation_rules.py`), déclenché par `POST /api/v1/recommendations/generate`,
par `make recommendations`, ou par la seconde tâche du DAG `alertes`, à partir des alertes déjà en par `make recommendations`, ou par la seconde tâche du DAG `alertes`, à partir des alertes déjà en
@@ -312,7 +304,6 @@ n'ajoute aucune ligne.
- Les mesures API ne sont pas rattachées à un dataset historique. - Les mesures API ne sont pas rattachées à un dataset historique.
- 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.
- Un site possède plusieurs rapports de dérive ; un rapport global n'est rattaché à aucun site.
## Ingestion des données historiques ## Ingestion des données historiques
+117 -38
View File
@@ -68,23 +68,36 @@ flowchart TB
push --> mv & ms push --> mv & ms
push --> av & ab push --> av & ab
push --> it push --> it
push --> sb1 & sb2 --> sscan push --> sb1 & sb2 & sb3 --> sscan
subgraph cd["Déploiement · deploy.yml"] subgraph cd["Déploiement · deploy.yml"]
dep["deploy<br/>runner eni-g3, environnement rec ou prod"] dep["deploy<br/>runner eni-g3, environnement rec ou prod"]
end end
push -->|"push sur dev ou main"| dep push -->|"push sur dev ou main"| dep
planifie["chaque lundi 3h UTC,<br/>ou à la main"]
subgraph dastw["DAST · dast.yml"]
zscan["zap<br/>seed + scan actif OWASP ZAP"]
end
planifie --> zscan
push -->|"PR sur dast.yml<br/>ou dast-token.sh"| zscan
``` ```
## Déclenchement ## Déclenchement
Les six workflows hébergés par GitHub se déclenchent sur `push` **et** sur `pull_request`, Les six workflows hébergés par GitHub qui vérifient le code se déclenchent sur `push` **et** sur
filtrés par **chemin** : `backend.yml` sur `apps/backend/**`, `frontend.yml` sur `pull_request`, filtrés par **chemin** : `backend.yml` sur `apps/backend/**`, `frontend.yml` sur
`apps/frontend/**`, `ml.yml` sur `ml/**`, `infra.yml` sur `infra/terraform/**`, `airflow.yml` sur `apps/frontend/**`, `ml.yml` sur `ml/**`, `infra.yml` sur `infra/terraform/**`, `airflow.yml` sur
`etl/airflow/**` **plus des chemins de `ml/` et de `apps/backend/`**, chacun incluant son propre `etl/airflow/**` **plus des chemins de `ml/` et de `apps/backend/`**, chacun incluant son propre
fichier de workflow dans le filtre pour qu'une modification du pipeline déclenche le pipeline. fichier de workflow dans le filtre pour qu'une modification du pipeline déclenche le pipeline.
`dast.yml` s'en écarte volontairement (détail dans sa propre section plus bas) : aucun
déclenchement sur `push`, seulement `workflow_dispatch`, une planification hebdomadaire, et
`pull_request` restreint à ses deux seuls fichiers. Un scan actif est trop long pour tourner à
chaque commit.
Le filtre d'`airflow.yml` mérite un mot : il inclut `ml/pyproject.toml`, `ml/uv.lock`, Le filtre d'`airflow.yml` mérite un mot : il inclut `ml/pyproject.toml`, `ml/uv.lock`,
`ml/enervision_ml/**`, `apps/backend/pyproject.toml`, `apps/backend/uv.lock` et `ml/enervision_ml/**`, `apps/backend/pyproject.toml`, `apps/backend/uv.lock` et
`apps/backend/app/**` parce que l'image Airflow copie le code et les dépendances des deux `apps/backend/app/**` parce que l'image Airflow copie le code et les dépendances des deux
@@ -108,7 +121,8 @@ environnement.
## Déploiement ## Déploiement
`deploy.yml` est le septième workflow, et le seul qui ne tourne pas chez GitHub : il s'exécute sur `deploy.yml` est le huitième workflow (`backend`, `frontend`, `ml`, `infra`, `airflow`,
`sonarqube`, `dast`, plus lui-même), et le seul qui ne tourne pas chez GitHub : il s'exécute sur
un runner auto-hébergé installé sur la VM ENI, label `eni-g3`, parce que les runners hébergés ne un runner auto-hébergé installé sur la VM ENI, label `eni-g3`, parce que les runners hébergés ne
joignent pas une adresse privée d'école. Le runner se connecte en sortie vers GitHub, aucun port joignent pas une adresse privée d'école. Le runner se connecte en sortie vers GitHub, aucun port
entrant n'est ouvert. entrant n'est ouvert.
@@ -155,8 +169,6 @@ dans [10-infra.md](10-infra.md).
| Typage `mypy` | backend (`app`), ml (strict) | zéro erreur | Bloque | | Typage `mypy` | backend (`app`), ml (strict) | zéro erreur | Bloque |
| Tests unitaires `pytest` | backend, ml | **`--cov-fail-under=85`** côté backend | Bloque | | Tests unitaires `pytest` | backend, ml | **`--cov-fail-under=85`** côté backend | Bloque |
| Tests d'intégration | backend | marqueur `integration`, base réelle | Bloque | | Tests d'intégration | backend | marqueur `integration`, base réelle | Bloque |
| Tests d'intégration ML ↔ DB | ml | marqueur `integration`, base réelle migrée par Alembic | Bloque |
| Chaîne ML → DB → API | ml | marqueur `chaine`, vrais binaires en sous-processus | Bloque |
| Audit de dépendances `pip-audit` | backend | sur le **verrou figé** | Bloque | | Audit de dépendances `pip-audit` | backend | sur le **verrou figé** | Bloque |
| Audit de dépendances `npm audit` | frontend | `--audit-level=high` | Bloque | | Audit de dépendances `npm audit` | frontend | `--audit-level=high` | Bloque |
| **SAST `bandit`** | backend (`app`), ml (`enervision_ml`) | **MEDIUM et au-dessus** | Bloque | | **SAST `bandit`** | backend (`app`), ml (`enervision_ml`) | **MEDIUM et au-dessus** | Bloque |
@@ -256,36 +268,112 @@ Ils ne transitent ni par git ni par GitHub, et le runner, qui travaille dans ce
à recevoir. Le revers : ils ne sont sauvegardés nulle part ailleurs. Un `.env` perdu se à recevoir. Le revers : ils ne sont sauvegardés nulle part ailleurs. Un `.env` perdu se
régénère, ce qui invalide les sessions et les connexions chiffrées par Airflow. régénère, ce qui invalide les sessions et les connexions chiffrées par Airflow.
### Pourquoi le job d'intégration ML installe aussi le backend ## Scan DAST (OWASP ZAP)
Le schéma de la base n'a qu'une source, les sept révisions Alembic de `apps/backend/alembic` : le Statut : `En cours`. Le workflow `dast.yml` attaque l'API **en fonctionnement**, ce que ni Bandit,
backend est propriétaire du schéma, `ml/` n'en est que consommateur. Reconstruire ce schéma à la ni `pip-audit`, ni Sonar ne font. Il se lance à la main (`workflow_dispatch`), chaque lundi à 3h
main dans le job ML donnerait un job vert sur une base qui n'est pas la nôtre, exactement l'erreur UTC, et sur une PR qui modifie le scan lui-même. Pas à chaque PR : un scan actif dure plusieurs
qu'évite déjà le choix de l'image `timescaledb-ha` plutôt qu'un `postgres` nu. Le job installe minutes.
donc les deux environnements uv, applique `alembic upgrade head`, puis joue `-m integration` côté
`ml/` et `-m chaine` côté backend.
Conséquence sur le déclenchement : les `paths` de `ml.yml` incluent `apps/backend/alembic/**` et Le job démarre sur le runner la base (même image TimescaleDB que `docker-compose.yml`, base
`apps/backend/app/**`. Sans eux, une migration qui renomme une colonne de `reading` ne jetable), applique les migrations, y sème un site et deux relevés (`db/seeds/` est vide, pas
déclencherait pas ce job, le SQL brut du pipeline dériverait du schéma, et **rien ne casserait encore d'outillage de jeu de données pour la CI ; sans données, `GET /sites` rend `[]`, chaque
avant la production**. `app/**` en entier, et non les seuls modèles : ce job est le seul à jouer `/{site_id}` rend 404, et le scan actif ne frappe que des gestionnaires d'erreur), démarre le
`-m chaine`, or la chaîne traverse les endpoints, les services et les schémas jusqu'à backend, puis `scripts/dast-token.sh` crée un compte **`lecteur`** et rend son jeton.
`GET /predictions`. Un filtre plus étroit laisserait le test muet sur la PR même qui le casse. Le
prix est qu'une PR backend lance aussi le lint et le typage de `ml/` : environ deux minutes de
runner, en parallèle. Même arbitrage que le filtre d'`airflow.yml`, qui écoute déjà `ml/**` et
`apps/backend/app/**` parce que son image réunit les deux.
Le marqueur `chaine` est distinct d'`integration` pour une raison mécanique : le job `integration` ZAP charge le contrat `/openapi.json` depuis un fichier (`zap-api-scan.py -f openapi -t
de `backend.yml` n'installe pas `ml/.venv`, et sélectionnerait sinon un test qui lance les /zap/wrk/openapi.json`) et en importe les 26 opérations **quel que soit le jeton** : c'est le
binaires du pipeline. Il est aussi exclu d'`addopts`, sans quoi `make test` échouerait sur tout contrat qui décide de ce qui est exploré, pas l'authentification. Le jeton ne change que les
poste où `ml/` n'est pas installé. réponses obtenues sur les routes gardées : sans lui, elles répondraient toutes `401` plutôt que
de dérouler leur logique. Huit routes n'exigent aucun jeton porteur (les deux sondes, `login`,
`refresh`, `logout`, `forgot-password`, `reset-password` et `reset-password/validate`) et
répondent donc pareil avec ou sans lui.
Décisions à savoir défendre :
- **Le compte du scan est `lecteur`, jamais `admin`.** Un scan actif avec un jeton admin frapperait
`POST /users` et la réinitialisation de mots de passe pour de bon. Le script passe par un admin
jetable pour créer le lecteur (l'API n'a pas d'inscription publique) puis ne s'en sert plus.
- **Un compte neuf est en `must_change_password`**, et toute route gardée le refuse tant que le
mot de passe n'est pas changé. Le script fait ce changement et vérifie `GET /sites` = 200 avant
de rendre le jeton ; sans cela, tout le scan authentifié ne testerait que des `403`.
`POST /auth/password` rend déjà un nouveau jeton valide (l'`iat` tronqué documenté dans
`app/api/deps.py` ne le rejette pas comme antérieur à la session) : le script s'en sert
directement plutôt que de se reconnecter, deux hachages Argon2id (19456 Kio chacun) et deux
allers-retours de refresh-token de moins sur le chemin critique de la CI.
- **`APP_ACCESS_TOKEN_TTL_SECONDS=3600`** (plafond de la configuration) : le jeton par défaut
dure 15 minutes. `scanner.maxScanDurationInMins=15` (ci-dessous) borne le scan actif très en
dessous, marge comprise pour les étapes qui l'entourent.
- **Le jeton ne transite ni par `${{ }}` dans le script de l'étape, ni par l'argv de `docker
run`.** Le premier finirait en clair dans le fichier de commande que GitHub écrit sur le disque
du runner pour toute la durée de l'étape ; le second serait visible par `ps aux` et par
`docker inspect zap` tant que le conteneur existe. Il est écrit dans un fichier de
configuration ZAP séparé (`-configfile`), monté en lecture seule hors de `/zap/wrk` pour ne
jamais atterrir dans l'artefact publié. ZAP journalise malgré tout la valeur de chaque
`-config`/`-configfile` chargé à un niveau visible sans `-d` : les copies de `zap.log` et
`zap-stdout.log` publiées en artefact sont donc caviardées avant publication.
**Deux pièges d'autorisation** sur ce fichier de configuration (`zap-auth.conf`), tous les deux
propres au montage bind Docker : le conteneur y lit avec son propre uid (1000), distinct de celui
du runner qui l'a écrit, sans remappage automatique.
- Un `chmod 600` seul rend le fichier illisible pour le conteneur (« File not readable :
/zap/auth.conf »). ZAP échoue dès le lancement, mais `zap-api-scan.py` attend les `-T` minutes
complètes avant d'abandonner : dix minutes qui ressemblent à un scan actif, pour un daemon mort
depuis le début. Corrigé par `sudo chown 1000:1000` du fichier avant de le passer à `644`.
- Ce `chown` déplace la propriété du fichier hors de l'utilisateur du runner : un `chmod` qui
suit sans `sudo` échoue alors (« Operation not permitted »), et le `-e` implicite des étapes
bash de GitHub Actions arrête toute l'étape avant même `docker run` — un scan « réussi » en une
fraction de seconde, sans le moindre journal ni rapport produit. Les deux commandes doivent
passer par `sudo`.
Les routes d'authentification qui changent l'état du compte (`login`, `password`, `logout-all`,
`forgot-password`, `reset-password`) sont exclues du scan actif : elles y déclencheraient la
limitation de débit et fermeraient les sessions sans rien apprendre de plus.
**Un scan vert n'est pas un scan qui a testé quelque chose.** Deux garde-fous, eux, **bloquent** :
- **Moins de 80% des opérations du contrat importées.** Constaté une première fois : 2 URL sur 26
opérations importées, ZAP n'avait envoyé que des requêtes vouées au 404 (l'analyseur de ZAP
refusait alors le nom accentué d'un des deux schémas de sécurité du contrat, corrigé depuis en
ASCII côté backend). Le seuil est dérivé du contrat (`zap-out/openapi.json`, présent à cette
étape) plutôt que d'un nombre fixe : un contrat qui grossit ne doit pas rendre la garde plus
permissive qu'elle ne l'était.
- **Aucune réponse 2xx.** Constaté une deuxième fois, cause différente : la clé de configuration
du nom d'en-tête pour la règle Replacer est `matchstr`, pas `matchstring` (celui-ci n'existe que
pour le job d'automatisation ZAP, pas pour `-config`) ; ZAP acceptait la mauvaise clé sans
erreur et laissait le nom d'en-tête vide, qu'uvicorn refusait par un `400` sur **toute** requête,
y compris les routes publiques. Piège de conception rencontré en corrigeant cette garde : borner
le *pourcentage* de 4xx ne marche pas, un scan actif fuzze délibérément un grand nombre
d'entrées invalides, si bien qu'un scan sain contre l'API seedée reste à 98% de 4xx avec
seulement 1% de 2xx. C'est la forme normale d'un scan actif. Le signal qui distingue vraiment un
scan cassé (2xx nul, absent du rapport dans les deux incidents) d'un scan sain (2xx non nul,
aussi faible soit-il) est l'absence de succès, pas la part d'échecs. Les deux gardes lisent
`zap-out/zap-report.json` (champs structurés `insights[]`), pas le texte libre du rapport
Markdown.
Le journal interne de ZAP (`zap.log`) et sa sortie complète (`zap-stdout.log`) sont publiés dans
l'artefact `zap-report` (dossier `zap-logs/`, propriété du runner : `zap-out/` bascule sous l'uid
1000 du conteneur ZAP dès que le contrat y est copié, le runner n'y écrit plus ensuite) pour
diagnostiquer un futur import raté.
**Non bloquant pour l'instant** (`continue-on-error`, sur la seule étape du scan) pour ce qui est
des alertes elles-mêmes. Le volume d'un premier passage trié est inconnu ; le rapport
HTML/JSON/Markdown est publié en artefact `zap-report`, et sa synthèse (jusqu'aux tableaux
d'alertes, sans le détail par alerte) dans le résumé du job. Fixer un seuil viendra une fois les
alertes triées.
**Limite à ne pas oublier :** le scan tape la configuration par défaut du backend (`APP_ENV=local`,
pas de TLS, pas de reverse proxy). Il remontera des alertes qui n'existent pas derrière le proxy
(HSTS absent...) et ne dit **rien** des en-têtes ni du TLS que le proxy pose en production. Un
second passage sur la stack complète reste à faire.
## Ce qui manque, et pourquoi ## Ce qui manque, et pourquoi
| Manque | Issue | Conséquence assumée | | Manque | Issue | Conséquence assumée |
|---|---|---| |---|---|---|
| Images publiées et promues par digest (GHCR) | aucune | Chaque environnement reconstruit ses images : la production n'exécute pas l'artefact validé en recette, mais un second build du même commit | | Images publiées et promues par digest (GHCR) | aucune | Chaque environnement reconstruit ses images : la production n'exécute pas l'artefact validé en recette, mais un second build du même commit |
| DAST (OWASP ZAP) | #41 | Aucune vérification sur l'application en fonctionnement, seulement sur le code et les dépendances | | DAST bloquant | #41 | Le scan ZAP existe mais ne bloque rien : aucun seuil n'est fixé tant que les alertes du premier passage ne sont pas triées |
| Tests end to end | #46 | Les parcours utilisateur ne sont pas vérifiés en CI | | Tests end to end | #46 | Les parcours utilisateur ne sont pas vérifiés en CI |
| Tests de charge | #47 | Aucun garde-fou de performance | | Tests de charge | #47 | Aucun garde-fou de performance |
| Scan d'image de conteneur | aucune | Les `Dockerfile` sont construits en local, pas analysés | | Scan d'image de conteneur | aucune | Les `Dockerfile` sont construits en local, pas analysés |
@@ -293,17 +381,8 @@ poste où `ml/` n'est pas installé.
## Reproduire la CI en local ## Reproduire la CI en local
`make check` enchaîne formatage, analyse statique, typage et tests du backend, c'est à dire le job `make check` enchaîne formatage, analyse statique, typage et tests du backend, c'est à dire le job
`verification`. `make ml-check` fait la même chose pour le module ML. `verification`. `make ml-check` fait la même chose pour le module ML. Les tests d'intégration
demandent une base : `make db-up` puis `uv run pytest -m integration`.
Les tests d'intégration demandent une base **migrée**, et `db/init` ne crée `enervision_test` que
vide :
```bash
make db-up migrate-test # la base de test reçoit les sept révisions Alembic
make test-integration # backend, marqueur `integration`
make ml-test-integration # pipeline ML, marqueur `integration`
make test-chaine # vrais binaires ML puis relecture par l'API, marqueur `chaine`
```
Le SAST se rejoue à l'identique : `uvx bandit==1.9.4 --recursive app --severity-level medium Le SAST se rejoue à l'identique : `uvx bandit==1.9.4 --recursive app --severity-level medium
--confidence-level medium` depuis `apps/backend`, et la même commande sur `enervision_ml` depuis --confidence-level medium` depuis `apps/backend`, et la même commande sur `enervision_ml` depuis
+1
View File
@@ -38,6 +38,7 @@ lecture seule ; plusieurs lignes resteront à compléter une fois les endpoints
| Caviardage des jetons, empreintes, mots de passe et cookies dans les journaux | `app/core/logging.py` | A09, A02 | | Caviardage des jetons, empreintes, mots de passe et cookies dans les journaux | `app/core/logging.py` | A09, A02 |
| Cinq gardes de configuration qui refusent le démarrage plutôt que de dégrader silencieusement | `app/core/config.py` | A05 | | Cinq gardes de configuration qui refusent le démarrage plutôt que de dégrader silencieusement | `app/core/config.py` | A05 |
| Documentation interactive fermée hors développement, `/metrics` derrière un jeton, sonde qui ne publie plus de version | `app/main.py`, `app/api/security.py` | A05 | | Documentation interactive fermée hors développement, `/metrics` derrière un jeton, sonde qui ne publie plus de version | `app/main.py`, `app/api/security.py` | A05 |
| Scan dynamique OWASP ZAP de l'API authentifiée (compte `lecteur` jetable), non bloquant, configuration par défaut du backend uniquement (ni TLS ni en-têtes du reverse proxy) | `.github/workflows/dast.yml`, `scripts/dast-token.sh` | A05, API8 Security Misconfiguration |
| En-têtes `nosniff`, `DENY`, `no-referrer`, et `no-store` sur les routes d'authentification | `app/api/middleware.py` | A05 | | En-têtes `nosniff`, `DENY`, `no-referrer`, et `no-store` sur les routes d'authentification | `app/api/middleware.py` | A05 |
| Refus de rétrograder ou désactiver le dernier administrateur actif | `app/services/user.py` | A04 Insecure Design | | Refus de rétrograder ou désactiver le dernier administrateur actif | `app/services/user.py` | A04 Insecure Design |
| Amorçage du premier administrateur hors dépôt, mot de passe jamais dans `argv` ni dans Git | `app/cli.py` | A02, A05 | | Amorçage du premier administrateur hors dépôt, mot de passe jamais dans `argv` ni dans Git | `app/cli.py` | A02, A05 |
+1 -1
View File
@@ -673,6 +673,6 @@ Le DAG `historical_import` est déclenché manuellement. Il exécute
`/opt/data/raw`. L'orchestration de l'import API Mock et la réconciliation globale des deux `/opt/data/raw`. L'orchestration de l'import API Mock et la réconciliation globale des deux
sources restent couvertes par l'issue #15. sources restent couvertes par l'issue #15.
Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py`, `alertes.py`, `historical_import.py` et `derive.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui). Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py` et `alertes.py` et `historical_import.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui).
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning. Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
-49
View File
@@ -1,49 +0,0 @@
"""DAG de surveillance de la dérive du modèle de prévision (issue #45).
Quotidien, pas horaire : la fenêtre mesurée couvre 168 h, la recalculer chaque heure écrirait
vingt-quatre lignes par jour pour un verdict qui ne bouge pas à cette cadence. Planifié après
les scorings de la nuit, et décalé de `ml_score` (à l'heure pile) comme de `alertes` (à la
quinzième minute).
Sans `--now` : le service tronque son instant de référence à l'heure, donc deux exécutions de
la même heure portent la même clé `uq_drift_report_window` et la seconde n'écrit rien.
Tâche distincte du DAG `alertes` plutôt qu'ajoutée à lui : un échec de dérive y ferait croire
que la détection d'alertes a échoué, et ce DAG porte un budget temporel déjà argumenté face à
son pas horaire.
"""
from __future__ import annotations
from datetime import datetime, timedelta
from airflow.providers.standard.operators.bash import BashOperator
from airflow.sdk import DAG
# Le backend a son propre environnement uv dans l'image (ADR 0008). `--no-sync` et
# `env -u VIRTUAL_ENV` : cf. `ml_train.py`, même raisonnement.
COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m"
# Piège : aucune reprise. Une dérive n'est pas un échec transitoire, la rejouer la redéclarerait
# à l'identique ; et la cadence quotidienne pardonne une connexion perdue.
TENTATIVES = 0
PLAFOND = timedelta(minutes=10)
with DAG(
dag_id="derive",
description=(
"Compare les prévisions déjà écrites aux lectures réellement arrivées "
"(app.monitoring.drift)."
),
schedule="30 5 * * *",
start_date=datetime(2026, 1, 1),
catchup=False,
max_active_runs=1,
tags=["ml", "monitoring"],
) as dag:
BashOperator(
task_id="derive",
bash_command=f"{COMMANDE_BACKEND} app.monitoring.drift",
retries=TENTATIVES,
execution_timeout=PLAFOND,
)
+1 -15
View File
@@ -10,14 +10,13 @@ from airflow.sdk import BaseOperator
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags" DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import", "derive"] DAG_IDS = ["ml_train", "ml_score", "alertes", "historical_import"]
TACHES = [ TACHES = [
("ml_train", "train"), ("ml_train", "train"),
("ml_score", "score"), ("ml_score", "score"),
("alertes", "detection"), ("alertes", "detection"),
("alertes", "recommandations"), ("alertes", "recommandations"),
("historical_import", "import_historical"), ("historical_import", "import_historical"),
("derive", "derive"),
] ]
@@ -160,19 +159,6 @@ def test_historical_import_retries_after_a_transient_failure(dagbag: DagBag) ->
assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1 assert dagbag.dags["historical_import"].get_task("import_historical").retries >= 1
def test_derive_runs_once_a_day(dagbag: DagBag) -> None:
assert dagbag.dags["derive"].timetable.expression == "30 5 * * *"
def test_derive_calls_the_backend_drift_module(dagbag: DagBag) -> None:
assert "app.monitoring.drift" in dagbag.dags["derive"].get_task("derive").bash_command
def test_derive_never_retries_a_detected_drift(dagbag: DagBag) -> None:
# Une derive n'est pas une panne passagere : la rejouer la redeclarerait a l'identique.
assert dagbag.dags["derive"].get_task("derive").retries == 0
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES) @pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
def test_tasks_never_resync_the_baked_environment( def test_tasks_never_resync_the_baked_environment(
dagbag: DagBag, dag_id: str, task_id: str dagbag: DagBag, dag_id: str, task_id: str
+6 -17
View File
@@ -111,23 +111,12 @@ Depuis la racine du monorepo, via le `Makefile` : `make install-ml`, `make ml-li
## Ou ecrire les tests ## Ou ecrire les tests
Deux regimes, separes par le marqueur `integration` que `pytest` ecarte par defaut. Aucun test ne touche PostgreSQL ni un serveur MLflow distant : `enervision_ml.data.load_from_csv`
et le chargement CSV de test suffisent a exercer `build_features` sur des donnees reelles ou
**Sans base** : `enervision_ml.data.load_from_csv` et le chargement CSV de test suffisent a synthetiques, et `enervision_ml.train.train()` accepte un `tracking_uri` SQLite isole (`tmp_path`
exercer `build_features` sur des donnees reelles ou synthetiques, et `enervision_ml.train.train()` pytest) pour un test de bout en bout sans effet de bord. `enervision_ml.data.load_from_database`
accepte un `tracking_uri` SQLite isole (`tmp_path` pytest) pour un test de bout en bout sans effet n'est pas encore couvert : il n'existe aucune base PostgreSQL a interroger en CI ni dans cet
de bord. environnement de developpement pour le moment.
**Avec base**, sous `integration` : `test_data_integration.py` confronte les neuf colonnes du
contrat au schema Alembic reel, et `test_score_integration.py` verifie les contraintes de
`prediction` depuis le code qui ecrit. Les fixtures sont dans `tests/conftest.py`, qui refuse de
demarrer si `ML_DATABASE_URL` ne vise pas `enervision_test`.
make db-up migrate-test ml-test-integration
Regle a tenir : **toute requete SQL nouvelle porte un test `integration`**. Le schema vit dans
`apps/backend/alembic`, pas ici : sans ce garde-fou, une migration qui renomme une colonne casse
le pipeline en production sans qu'aucun test ne rougisse.
## Piege a connaitre ## Piege a connaitre
+7 -18
View File
@@ -43,10 +43,6 @@ NUMERIC_COLUMNS = [
"capacity_kw", "capacity_kw",
] ]
# Piege : `reading.is_working_hours` est nullable et entre dans les features. Une seule lecture a
# NULL rend la colonne `object`, que LightGBM refuse ("pandas dtypes must be int, float or bool").
FLAG_COLUMNS = ["is_working_hours"]
_READING_QUERY = text( _READING_QUERY = text(
""" """
SELECT SELECT
@@ -80,7 +76,7 @@ _RECENT_READING_QUERY = text(
s.capacity_kw s.capacity_kw
FROM reading r FROM reading r
JOIN site s ON s.site_id = r.site_id JOIN site s ON s.site_id = r.site_id
WHERE r.timestamp >= :since AND r.timestamp <= :until WHERE r.timestamp >= :since
ORDER BY r.site_id, r.timestamp ORDER BY r.site_id, r.timestamp
""" """
) )
@@ -94,21 +90,14 @@ def load_from_database(connection: Connectable) -> pd.DataFrame:
return _typer(frame[OUTPUT_COLUMNS]) return _typer(frame[OUTPUT_COLUMNS])
def load_recent_from_database( def load_recent_from_database(connection: Connectable, *, since: datetime) -> pd.DataFrame:
connection: Connectable, *, since: datetime, until: datetime """Lit `reading` + `site` depuis `since` seulement, pour le scoring.
) -> pd.DataFrame:
"""Lit `reading` + `site` sur la fenetre `[since, until]`, pour le scoring.
Piege evite cote bas : un `SELECT` sans borne sur l'hypertable complete juste pour scorer le Piege evite : un `SELECT` sans borne sur l'hypertable complete juste pour scorer le prochain
prochain pas horaire serait la meme erreur que celle corrigee sur `GET /readings` (fenetre non pas horaire serait la meme erreur que celle corrigee sur `GET /readings` (fenetre non
plafonnee sur une table pouvant porter des annees d'historique). plafonnee sur une table pouvant porter des annees d'historique).
Piege evite cote haut : `until` est obligatoire, et c'est ce qui donne son sens a `--now`.
Sans lui, `build_scoring_frame` repartait de la derniere lecture de toute la table quel que
soit l'instant demande, donc `target_at` valait toujours "fin du jeu + 1h" et l'age de la
derniere lecture devenait negatif sans que rien ne le signale.
""" """
frame = pd.read_sql(_RECENT_READING_QUERY, connection, params={"since": since, "until": until}) frame = pd.read_sql(_RECENT_READING_QUERY, connection, params={"since": since})
return _typer(frame[OUTPUT_COLUMNS]) return _typer(frame[OUTPUT_COLUMNS])
@@ -137,6 +126,6 @@ def _typer(frame: pd.DataFrame) -> pd.DataFrame:
scoring -- ce n'est pas un effet de bord limite aux colonnes mesurees. scoring -- ce n'est pas un effet de bord limite aux colonnes mesurees.
""" """
typee = frame.copy() typee = frame.copy()
for colonne in (*NUMERIC_COLUMNS, *FLAG_COLUMNS): for colonne in NUMERIC_COLUMNS:
typee[colonne] = pd.to_numeric(typee[colonne], errors="coerce") typee[colonne] = pd.to_numeric(typee[colonne], errors="coerce")
return typee return typee
+2 -3
View File
@@ -204,8 +204,7 @@ def _load_recent_from_csv(csv_path: Path, *, now: datetime | None) -> tuple[pd.D
instant = now or ( instant = now or (
brute["timestamp"].max().to_pydatetime() if not brute.empty else datetime.now(UTC) brute["timestamp"].max().to_pydatetime() if not brute.empty else datetime.now(UTC)
) )
fenetre = (brute["timestamp"] >= instant - LOOKBACK) & (brute["timestamp"] <= instant) return brute[brute["timestamp"] >= instant - LOOKBACK], instant
return brute[fenetre], instant
def _score_frame( def _score_frame(
@@ -241,7 +240,7 @@ def run_scoring(
engine = create_engine(config.database_url()) engine = create_engine(config.database_url())
try: try:
instant = now or datetime.now(UTC) instant = now or datetime.now(UTC)
recent = load_recent_from_database(engine, since=instant - LOOKBACK, until=instant) recent = load_recent_from_database(engine, since=instant - LOOKBACK)
resultats = _score_frame(recent, model_path=model_path, site_id=site_id, instant=instant) resultats = _score_frame(recent, model_path=model_path, site_id=site_id, instant=instant)
reference = model_reference(model_path) reference = model_reference(model_path)
View File
-332
View File
@@ -1,332 +0,0 @@
"""Piege : deux fixtures d'acces a la base, jamais interchangeables - `connexion_ml` et `parc`.
`connexion_ml` ouvre une transaction annulee a la fin du test : rien ne subsiste, et rien n'est
visible hors de cette connexion. Elle sert aux fonctions qui recoivent leur connexion en
argument (`load_from_database`, `load_recent_from_database`, `write_predictions`).
`run_scoring` fabrique en revanche son propre engine depuis `ML_DATABASE_URL` : il ne verrait
pas des lignes semees dans une transaction non validee, et ses propres ecritures survivraient a
l'annulation. Les tests qui l'appellent passent donc par `parc`, qui valide ce qu'il ecrit et
nettoie lui-meme, dans l'ordre impose par les cles etrangeres `RESTRICT`.
"""
import math
import os
from collections.abc import Iterator
from dataclasses import dataclass, field
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Any
from uuid import uuid4
import lightgbm as lgb
import pandas as pd
import pytest
from sqlalchemy import Connection, Engine, Row, bindparam, create_engine, text
from sqlalchemy.engine import URL, make_url
from enervision_ml.features import TARGET_COLUMN, build_features, feature_columns
BASE_ATTENDUE = "enervision_test"
# Piege : `load_from_database` lit toute la table, et `enervision_test` est partagee entre un run
# local et la CI. Les tests ancrent donc leurs lectures au-dela de tout jeu de donnees reel
# (l'historique s'arrete au 31/12/2024) pour que leur borne `since` ne ramene qu'eux.
ANCRAGE = datetime(2035, 1, 1, tzinfo=UTC)
SITE_TYPE = "office"
CAPACITY_KW = 100.0
_INSERT_SITE = text(
"""
INSERT INTO site (site_id, site_name, site_type, capacity_kw)
VALUES (:site_id, :site_name, :site_type, :capacity_kw)
"""
)
# `source = 'api_history'` impose `dataset_id IS NULL` (ck_reading_dataset_source), ce qui evite
# de creer une ligne `dataset`. `raw_data` est NOT NULL, d'ou le litteral jsonb.
_INSERT_READING = text(
"""
INSERT INTO reading (
site_id, timestamp, source, consumption_kwh, temperature_celsius,
humidity_percent, solar_irradiance_wm2, is_working_hours, raw_data
) VALUES (
:site_id, :timestamp, :source, :consumption_kwh, :temperature_celsius,
:humidity_percent, :solar_irradiance_wm2, :is_working_hours, '{}'::jsonb
)
"""
)
_SELECT_PREDICTIONS = text(
"""
SELECT target_at, predicted_value, status, failure_reason, model_reference
FROM prediction
WHERE site_id = :site_id
ORDER BY prediction_id
"""
)
_INSERT_PREDICTION = text(
"""
INSERT INTO prediction (
site_id, target_at, target_metric, period_minutes,
predicted_value, model_reference, status, failure_reason
) VALUES (
:site_id, :target_at, 'consumption_kwh', 60,
:predicted_value, :model_reference, :status, :failure_reason
)
"""
)
# Ordre impose par les cles etrangeres `RESTRICT` : une lecture avant son site, une prediction
# avant sa lecture.
_SUPPRESSIONS = tuple(
text(requete).bindparams(bindparam("sites", expanding=True))
for requete in (
"DELETE FROM prediction WHERE site_id IN :sites",
"DELETE FROM reading WHERE site_id IN :sites",
"DELETE FROM site WHERE site_id IN :sites",
)
)
def insere_site(
connexion: Connection,
*,
site_type: str = SITE_TYPE,
capacity_kw: float | None = CAPACITY_KW,
) -> str:
site_id = f"TEST-{uuid4().hex[:12]}"
connexion.execute(
_INSERT_SITE,
{
"site_id": site_id,
"site_name": "Site de test",
"site_type": site_type,
"capacity_kw": capacity_kw,
},
)
return site_id
def insere_lectures(
connexion: Connection,
site_id: str,
*,
heures: int,
fin: datetime,
valeur: float = 50.0,
source: str = "api_history",
is_working_hours: bool | None = True,
) -> list[datetime]:
"""Grille horaire contigue finissant a `fin`, incluse.
Contigue parce que les lags de `build_features` sont des `shift()` positionnels : un trou
dans la grille decalerait le lag de 168 h sans qu'aucune erreur ne se declenche.
"""
instants = [fin - timedelta(hours=decalage) for decalage in reversed(range(heures))]
connexion.execute(
_INSERT_READING,
[
{
"site_id": site_id,
"timestamp": instant,
"source": source,
"consumption_kwh": valeur + math.sin(rang / 12.0) * 10.0,
"temperature_celsius": 15.0,
"humidity_percent": 50.0,
"solar_irradiance_wm2": 0.0,
"is_working_hours": is_working_hours,
}
for rang, instant in enumerate(instants)
],
)
return instants
def insere_lecture(
connexion: Connection,
site_id: str,
*,
instant: datetime,
consumption_kwh: float | None = 50.0,
source: str = "api_history",
is_working_hours: bool | None = True,
) -> None:
"""Une lecture isolee, quand le test pilote sa valeur plutot que sa forme."""
connexion.execute(
_INSERT_READING,
{
"site_id": site_id,
"timestamp": instant,
"source": source,
"consumption_kwh": consumption_kwh,
"temperature_celsius": 15.0,
"humidity_percent": 50.0,
"solar_irradiance_wm2": 0.0,
"is_working_hours": is_working_hours,
},
)
def insere_prediction(
connexion: Connection,
site_id: str,
*,
target_at: datetime,
predicted_value: float | None = 42.0,
model_reference: str = "lightgbm-test000000",
status: str = "available",
failure_reason: str | None = None,
) -> None:
connexion.execute(
_INSERT_PREDICTION,
{
"site_id": site_id,
"target_at": target_at,
"predicted_value": predicted_value,
"model_reference": model_reference,
"status": status,
"failure_reason": failure_reason,
},
)
@pytest.fixture(scope="session")
def url_ml() -> URL:
valeur = os.environ.get("ML_DATABASE_URL")
if not valeur:
pytest.fail("ML_DATABASE_URL absente. Voir `make ml-test-integration`.")
url = make_url(valeur)
if url.database != BASE_ATTENDUE:
pytest.fail(
f"Ces tests ecrivent et suppriment : ML_DATABASE_URL doit viser {BASE_ATTENDUE}, "
f"pas {url.database}."
)
return url
@pytest.fixture(scope="session")
def moteur_ml(url_ml: URL) -> Iterator[Engine]:
moteur = create_engine(url_ml)
try:
yield moteur
finally:
moteur.dispose()
@pytest.fixture
def connexion_ml(moteur_ml: Engine) -> Iterator[Connection]:
with moteur_ml.connect() as connexion:
transaction = connexion.begin()
try:
yield connexion
finally:
transaction.rollback()
@dataclass
class Parc:
"""Semis valide en base, et son nettoyage, pour les tests qui appellent `run_scoring`.
Chaque `site_id` porte une marque unique : la base de test est partagee entre un run local
et la CI.
"""
moteur: Engine
sites: list[str] = field(default_factory=list)
def site(self, *, site_type: str = SITE_TYPE, capacity_kw: float | None = CAPACITY_KW) -> str:
with self.moteur.begin() as connexion:
site_id = insere_site(connexion, site_type=site_type, capacity_kw=capacity_kw)
self.sites.append(site_id)
return site_id
def lectures(self, site_id: str, **arguments: Any) -> list[datetime]:
with self.moteur.begin() as connexion:
return insere_lectures(connexion, site_id, **arguments)
def lecture(self, site_id: str, **arguments: Any) -> None:
with self.moteur.begin() as connexion:
insere_lecture(connexion, site_id, **arguments)
def prediction(self, site_id: str, **arguments: Any) -> None:
with self.moteur.begin() as connexion:
insere_prediction(connexion, site_id, **arguments)
def predictions_ecrites(self, site_id: str) -> list[Row[Any]]:
with self.moteur.connect() as connexion:
return list(connexion.execute(_SELECT_PREDICTIONS, {"site_id": site_id}))
def nettoie(self) -> None:
if not self.sites:
return
with self.moteur.begin() as connexion:
for suppression in _SUPPRESSIONS:
connexion.execute(suppression, {"sites": self.sites})
@pytest.fixture
def parc(moteur_ml: Engine) -> Iterator[Parc]:
semis = Parc(moteur=moteur_ml)
try:
yield semis
finally:
semis.nettoie()
def trame_synthetique(*, sites: int = 2, heures: int = 400) -> pd.DataFrame:
"""Lectures horaires deterministes, assez longues pour que le lag de 168 h existe."""
depart = datetime(2024, 1, 1, tzinfo=UTC)
morceaux = [
pd.DataFrame(
{
"site_id": f"SITE{numero:03d}",
"timestamp": [depart + timedelta(hours=rang) for rang in range(heures)],
TARGET_COLUMN: [
50.0 + 10.0 * math.sin(rang / 12.0) + numero * 5.0 for rang in range(heures)
],
"temperature_celsius": 15.0,
"humidity_percent": 50.0,
"solar_irradiance_wm2": 0.0,
"is_working_hours": True,
"site_type": SITE_TYPE,
"capacity_kw": CAPACITY_KW,
}
)
for numero in range(sites)
]
return pd.concat(morceaux, ignore_index=True)
@pytest.fixture(scope="session")
def modele_jetable(tmp_path_factory: pytest.TempPathFactory) -> Path:
"""Booster reel entraine sur une trame synthetique, ecrit dans un repertoire temporaire.
Ni `ml/models/` (ignore par git, et le polluer serait un effet de bord), ni
`enervision_ml.train.train()` (qui journalise dans MLflow sans garde). Le typage `category`
de `site_type` reproduit celui de l'entrainement : c'est le `pandas_categorical` enregistre
dans le modele que `score()` devra retrouver.
"""
features = build_features(trame_synthetique()).dropna(subset=feature_columns())
typee = features.copy()
typee["site_type"] = typee["site_type"].astype("category")
donnees = lgb.Dataset(
typee[feature_columns()],
label=typee[TARGET_COLUMN],
categorical_feature=["site_type"],
)
booster = lgb.train(
{"objective": "regression", "num_leaves": 7, "min_data_in_leaf": 5, "verbosity": -1},
donnees,
num_boost_round=5,
)
chemin = tmp_path_factory.mktemp("modele") / "lightgbm-consumption.txt"
booster.save_model(str(chemin))
return chemin
-189
View File
@@ -1,189 +0,0 @@
from datetime import timedelta
from pathlib import Path
import pandas as pd
import pytest
from sqlalchemy import Connection
from enervision_ml.data import (
OUTPUT_COLUMNS,
load_from_csv,
load_from_database,
load_recent_from_database,
)
from tests.conftest import ANCRAGE, insere_lecture, insere_lectures, insere_site
pytestmark = pytest.mark.integration
def du_site(frame: pd.DataFrame, site_id: str) -> pd.DataFrame:
"""Piege : les chargeurs ne filtrent pas par site, et `enervision_test` est partagee avec
les tests qui valident leurs ecritures. Juger le contenu de toute la fenetre les couplerait."""
return frame[frame["site_id"] == site_id].reset_index(drop=True)
def test_load_from_database_returns_the_nine_contract_columns(connexion_ml: Connection) -> None:
site_id = insere_site(connexion_ml)
insere_lectures(connexion_ml, site_id, heures=3, fin=ANCRAGE)
frame = load_from_database(connexion_ml)
assert list(frame.columns) == OUTPUT_COLUMNS
def test_load_from_database_joins_the_site_attributes_to_every_reading(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml, site_type="factory", capacity_kw=250.0)
insere_lectures(connexion_ml, site_id, heures=3, fin=ANCRAGE)
frame = load_from_database(connexion_ml)
mien = frame[frame["site_id"] == site_id]
assert len(mien) == 3
assert set(mien["site_type"]) == {"factory"}
assert set(mien["capacity_kw"]) == {250.0}
def test_load_recent_from_database_excludes_readings_before_the_since_bound(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
insere_lectures(connexion_ml, site_id, heures=5, fin=ANCRAGE)
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE - timedelta(hours=2), until=ANCRAGE
)
assert list(du_site(frame, site_id)["timestamp"]) == [
ANCRAGE - timedelta(hours=2),
ANCRAGE - timedelta(hours=1),
ANCRAGE,
]
def test_load_recent_from_database_includes_a_reading_exactly_at_the_since_bound(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
insere_lecture(connexion_ml, site_id, instant=ANCRAGE)
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3)
)
assert len(du_site(frame, site_id)) == 1
def test_load_recent_from_database_keeps_timestamps_timezone_aware(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
insere_lecture(connexion_ml, site_id, instant=ANCRAGE)
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3)
)
assert frame["timestamp"].dt.tz is not None
def test_load_recent_from_database_orders_readings_by_site_then_timestamp(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
for decalage in (2, 0, 1):
insere_lecture(connexion_ml, site_id, instant=ANCRAGE + timedelta(hours=decalage))
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3)
)
assert list(du_site(frame, site_id)["timestamp"]) == [
ANCRAGE,
ANCRAGE + timedelta(hours=1),
ANCRAGE + timedelta(hours=2),
]
def test_load_recent_from_database_returns_the_contract_columns_even_without_any_row(
connexion_ml: Connection,
) -> None:
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE + timedelta(days=365), until=ANCRAGE + timedelta(days=400)
)
assert frame.empty
assert list(frame.columns) == OUTPUT_COLUMNS
def test_load_recent_from_database_types_a_fully_null_capacity_kw_as_float64(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml, capacity_kw=None)
insere_lectures(connexion_ml, site_id, heures=3, fin=ANCRAGE)
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE - timedelta(hours=2), until=ANCRAGE
)
assert frame["capacity_kw"].dtype == "float64"
assert du_site(frame, site_id)["capacity_kw"].isna().all()
def test_load_recent_from_database_types_a_null_is_working_hours_as_float64(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
insere_lecture(connexion_ml, site_id, instant=ANCRAGE, is_working_hours=None)
insere_lecture(
connexion_ml, site_id, instant=ANCRAGE + timedelta(hours=1), is_working_hours=True
)
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3)
)
assert frame["is_working_hours"].dtype == "float64"
assert list(du_site(frame, site_id)["is_working_hours"].isna()) == [True, False]
def test_both_loaders_produce_the_same_columns_in_the_same_order(
connexion_ml: Connection, tmp_path: Path
) -> None:
site_id = insere_site(connexion_ml)
insere_lectures(connexion_ml, site_id, heures=2, fin=ANCRAGE)
csv_path = tmp_path / "lectures.csv"
pd.DataFrame(
{
"site_id": [site_id],
"timestamp": [ANCRAGE],
"consumption_kwh": [50.0],
"temperature_celsius": [15.0],
"humidity_percent": [50.0],
"solar_irradiance_wm2": [0.0],
"is_working_hours": [True],
"site_type": ["office"],
}
).to_csv(csv_path, index=False)
depuis_la_base = load_recent_from_database(
connexion_ml, since=ANCRAGE - timedelta(hours=1), until=ANCRAGE
)
depuis_le_csv = load_from_csv(csv_path)
assert list(depuis_la_base.columns) == list(depuis_le_csv.columns)
assert depuis_la_base.dtypes.to_dict() == depuis_le_csv.dtypes.to_dict()
def test_load_recent_from_database_excludes_readings_after_the_until_bound(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
insere_lectures(connexion_ml, site_id, heures=5, fin=ANCRAGE + timedelta(hours=4))
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE - timedelta(days=1), until=ANCRAGE
)
assert list(du_site(frame, site_id)["timestamp"]) == [ANCRAGE]
-21
View File
@@ -275,24 +275,3 @@ def test_run_scoring_in_csv_mode_scores_without_touching_a_database(tmp_path: Pa
assert {r.site_id for r in resultats} == {"site-a", "site-b"} assert {r.site_id for r in resultats} == {"site-a", "site-b"}
assert all(r.status == "available" for r in resultats) assert all(r.status == "available" for r in resultats)
assert all(r.predicted_value == 7.0 for r in resultats) assert all(r.predicted_value == 7.0 for r in resultats)
def test_run_scoring_in_csv_mode_targets_the_hour_after_the_reference_instant(
tmp_path: Path,
) -> None:
depart = datetime(2026, 1, 1, tzinfo=UTC)
frame = make_recent("site-a", heures=400, depart=depart)
csv_path = tmp_path / "recent.csv"
frame.to_csv(csv_path, index=False)
model_path = tmp_path / "model.txt"
model_path.write_bytes(b"peu importe le contenu pour ce test")
rattrapage = depart + timedelta(hours=300)
with pytest.MonkeyPatch.context() as monkeypatch:
monkeypatch.setattr(
"enervision_ml.score.lgb.Booster", lambda model_file: FakeBooster(valeur=7.0)
)
resultats = run_scoring(model_path=model_path, csv_path=csv_path, now=rattrapage)
assert [r.target_at for r in resultats] == [rattrapage + timedelta(hours=1)]
-238
View File
@@ -1,238 +0,0 @@
from datetime import datetime, timedelta
from pathlib import Path
from typing import Any
import pytest
from sqlalchemy import Connection, Row, text
from sqlalchemy.exc import IntegrityError
from enervision_ml.score import (
INSUFFICIENT_DATA_REASON,
LOOKBACK,
MAX_STALENESS,
ScoredSite,
model_reference,
run_scoring,
write_predictions,
)
from tests.conftest import ANCRAGE, Parc, insere_site
pytestmark = pytest.mark.integration
REFERENCE = "lightgbm-000000000000"
_SELECT = text(
"""
SELECT target_at, target_metric, period_minutes, predicted_value,
model_reference, status, failure_reason
FROM prediction
WHERE site_id = :site_id
ORDER BY prediction_id
"""
)
def lignes(connexion: Connection, site_id: str) -> list[Row[Any]]:
return list(connexion.execute(_SELECT, {"site_id": site_id}))
def disponible(
site_id: str,
*,
target_at: datetime = ANCRAGE,
predicted_value: float | None = 12.5,
) -> ScoredSite:
return ScoredSite(
site_id=site_id,
target_at=target_at,
status="available",
predicted_value=predicted_value,
failure_reason=None,
)
def test_write_predictions_inserts_one_row_per_scored_site(connexion_ml: Connection) -> None:
premier = insere_site(connexion_ml)
second = insere_site(connexion_ml)
write_predictions(connexion_ml, [disponible(premier), disponible(second)], reference=REFERENCE)
assert len(lignes(connexion_ml, premier)) == 1
assert len(lignes(connexion_ml, second)) == 1
def test_write_predictions_stores_the_model_reference_and_the_hourly_period(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
write_predictions(connexion_ml, [disponible(site_id)], reference=REFERENCE)
ligne = lignes(connexion_ml, site_id)[0]
assert ligne.model_reference == REFERENCE
assert ligne.target_metric == "consumption_kwh"
assert ligne.period_minutes == 60
def test_write_predictions_stacks_a_second_run_instead_of_overwriting_the_first(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
write_predictions(
connexion_ml, [disponible(site_id, predicted_value=10.0)], reference=REFERENCE
)
write_predictions(
connexion_ml, [disponible(site_id, predicted_value=20.0)], reference=REFERENCE
)
assert [ligne.predicted_value for ligne in lignes(connexion_ml, site_id)] == [10.0, 20.0]
def test_write_predictions_writes_nothing_when_no_site_was_scored(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
write_predictions(connexion_ml, [], reference=REFERENCE)
assert lignes(connexion_ml, site_id) == []
def test_write_predictions_rejects_an_available_row_without_a_predicted_value(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
with pytest.raises(IntegrityError, match="ck_prediction_status"):
write_predictions(
connexion_ml, [disponible(site_id, predicted_value=None)], reference=REFERENCE
)
def test_write_predictions_rejects_an_insufficient_data_row_carrying_a_value(
connexion_ml: Connection,
) -> None:
site_id = insere_site(connexion_ml)
incoherent = ScoredSite(
site_id=site_id,
target_at=ANCRAGE,
status="insufficient_data",
predicted_value=12.5,
failure_reason=INSUFFICIENT_DATA_REASON,
)
with pytest.raises(IntegrityError, match="ck_prediction_status"):
write_predictions(connexion_ml, [incoherent], reference=REFERENCE)
def test_write_predictions_rejects_a_prediction_for_an_unknown_site(
connexion_ml: Connection,
) -> None:
with pytest.raises(IntegrityError, match="fk_prediction_site"):
write_predictions(connexion_ml, [disponible("SITE-INCONNU")], reference=REFERENCE)
def test_run_scoring_writes_an_available_prediction_for_a_site_with_a_full_week(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=200, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE)
ligne = parc.predictions_ecrites(site_id)[0]
assert ligne.status == "available"
assert ligne.predicted_value is not None
assert ligne.target_at == ANCRAGE + timedelta(hours=1)
def test_run_scoring_writes_insufficient_data_when_the_weekly_lag_is_missing(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=100, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE)
ligne = parc.predictions_ecrites(site_id)[0]
assert ligne.status == "insufficient_data"
assert ligne.predicted_value is None
assert ligne.failure_reason == INSUFFICIENT_DATA_REASON
def test_run_scoring_writes_a_staleness_reason_when_the_last_reading_is_too_old(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=200, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE + MAX_STALENESS + timedelta(hours=1))
ligne = parc.predictions_ecrites(site_id)[0]
assert ligne.status == "insufficient_data"
assert ligne.failure_reason != INSUFFICIENT_DATA_REASON
def test_run_scoring_writes_nothing_when_every_reading_is_older_than_the_window(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=200, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE + LOOKBACK + timedelta(days=1))
assert parc.predictions_ecrites(site_id) == []
def test_run_scoring_only_writes_the_site_that_was_requested(
parc: Parc, modele_jetable: Path
) -> None:
demande = parc.site()
ignore = parc.site()
parc.lectures(demande, heures=200, fin=ANCRAGE)
parc.lectures(ignore, heures=200, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, site_id=demande, now=ANCRAGE)
assert len(parc.predictions_ecrites(demande)) == 1
assert parc.predictions_ecrites(ignore) == []
def test_run_scoring_uses_the_model_file_hash_as_model_reference(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=200, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE)
ligne = parc.predictions_ecrites(site_id)[0]
assert ligne.model_reference == model_reference(modele_jetable)
def test_run_scoring_appends_a_second_row_when_it_runs_twice(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=200, fin=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE)
run_scoring(model_path=modele_jetable, now=ANCRAGE)
ecrites = parc.predictions_ecrites(site_id)
assert len(ecrites) == 2
assert ecrites[0].target_at == ecrites[1].target_at
def test_run_scoring_targets_the_hour_after_the_reference_instant(
parc: Parc, modele_jetable: Path
) -> None:
site_id = parc.site()
parc.lectures(site_id, heures=200, fin=ANCRAGE + timedelta(hours=48))
rattrapage = ANCRAGE
run_scoring(model_path=modele_jetable, now=rattrapage)
ligne = parc.predictions_ecrites(site_id)[0]
assert ligne.target_at == rattrapage + timedelta(hours=1)
+8
View File
@@ -1,3 +1,11 @@
# Scripts # Scripts
Outillage local du monorepo. Les taches courantes passent par le `Makefile` racine. Outillage local du monorepo. Les taches courantes passent par le `Makefile` racine.
## dast-token.sh
Prépare le scan DAST (`.github/workflows/dast.yml`) : sur une API déjà démarrée, crée un compte
`lecteur` jetable, lui fait passer le changement de mot de passe obligatoire et écrit son jeton
d'accès sur la sortie standard. À lancer depuis `apps/backend`, contre une base **jetable** (il y
crée deux comptes) : `BASE_URL=http://localhost:8000 ../../scripts/dast-token.sh`. Nécessite `curl`,
`jq` et `openssl`.
+78
View File
@@ -0,0 +1,78 @@
#!/usr/bin/env bash
# Prépare le scan DAST : crée un compte `lecteur` sur une API déjà démarrée, lui fait passer le
# changement de mot de passe obligatoire, et écrit son jeton d'accès sur la sortie standard.
#
# Piège : un compte neuf est en `must_change_password`, et toute route gardée le refuse tant que
# le mot de passe n'a pas été changé. Sans cette étape, ZAP ne verrait que 403 sur les routes
# gardées et le scan ne testerait rien de l'API authentifiée.
#
# Contrainte : le compte du scan est `lecteur`, jamais `admin`. Un scan actif avec un jeton admin
# frapperait POST /users ou la réinitialisation de mots de passe pour de bon.
#
# L'administrateur n'existe que pour créer ce compte (l'API n'a pas d'inscription publique).
# À lancer depuis apps/backend, dans un environnement où DATABASE_URL et APP_SECRET_KEY visent
# une base JETABLE : le script y crée deux comptes.
set -euo pipefail
BASE_URL="${BASE_URL:-http://localhost:8000}"
API="$BASE_URL/api/v1"
SUFFIXE="$(openssl rand -hex 4)"
EMAIL_ADMIN="dast-admin-$SUFFIXE@enervision.fr"
EMAIL_LECTEUR="dast-lecteur-$SUFFIXE@enervision.fr"
# Classes exigées par le validateur : majuscule, minuscule, chiffre, caractère spécial.
nouveau_mot_de_passe() { echo "Dast-$(openssl rand -hex 12)-Aa1!"; }
# Tout ce qui n'est pas la sortie finale part sur stderr : la sortie standard ne porte que le jeton.
journal() { echo "dast-token: $*" >&2; }
connexion() {
local email="$1" mot_de_passe="$2"
curl -fsS -X POST "$API/auth/login" -H 'Content-Type: application/json' \
-d "$(jq -n --arg e "$email" --arg p "$mot_de_passe" '{email:$e, password:$p}')" \
| jq -r '.access_token'
}
# Rend le nouveau jeton d'accès : `/auth/password` en émet un (avec l'`iat` de la session en
# cours, cf. le piège documenté dans `app/api/deps.py`), pas seulement une confirmation. S'y fier
# évite une reconnexion, donc un second hachage Argon2id (19456 Kio) et un aller-retour de
# refresh-token superflus sur le chemin critique de la CI.
changer_mot_de_passe() {
local jeton="$1" ancien="$2" nouveau="$3"
curl -fsS -X POST "$API/auth/password" \
-H "Authorization: Bearer $jeton" -H 'Content-Type: application/json' \
-d "$(jq -n --arg a "$ancien" --arg n "$nouveau" '{current_password:$a, new_password:$n}')" \
| jq -r '.access_token'
}
journal "création de l'administrateur $EMAIL_ADMIN"
if ! SORTIE="$(uv run --frozen --no-sync --no-build python -m app.cli create-admin --email "$EMAIL_ADMIN" --generate)"; then
journal "la création de l'administrateur a échoué :"
journal "$SORTIE"
exit 1
fi
MDP_ADMIN="$(sed -n 's/^Mot de passe généré, il ne sera plus affiché : //p' <<<"$SORTIE")"
[[ -n "$MDP_ADMIN" ]] || { journal "mot de passe administrateur introuvable dans la sortie :"; journal "$SORTIE"; exit 1; }
JETON="$(connexion "$EMAIL_ADMIN" "$MDP_ADMIN")"
NOUVEAU_ADMIN="$(nouveau_mot_de_passe)"
JETON="$(changer_mot_de_passe "$JETON" "$MDP_ADMIN" "$NOUVEAU_ADMIN")"
journal "création du lecteur $EMAIL_LECTEUR"
REPONSE="$(curl -fsS -X POST "$API/users" -H "Authorization: Bearer $JETON" \
-H 'Content-Type: application/json' \
-d "$(jq -n --arg e "$EMAIL_LECTEUR" '{email:$e, role:"lecteur"}')")"
MDP_TEMPORAIRE="$(jq -r '.temporary_password // empty' <<<"$REPONSE")"
[[ -n "$MDP_TEMPORAIRE" ]] || { journal "mot de passe temporaire introuvable dans la réponse de POST /users :"; journal "$REPONSE"; exit 1; }
JETON="$(connexion "$EMAIL_LECTEUR" "$MDP_TEMPORAIRE")"
NOUVEAU_LECTEUR="$(nouveau_mot_de_passe)"
JETON="$(changer_mot_de_passe "$JETON" "$MDP_TEMPORAIRE" "$NOUVEAU_LECTEUR")"
# Vérifie que le jeton ouvre bien une route gardée avant de le rendre.
CODE="$(curl -sS -o /dev/null -w '%{http_code}' "$API/sites" -H "Authorization: Bearer $JETON")"
[[ "$CODE" == "200" ]] || { journal "GET /sites répond $CODE avec le jeton du lecteur, attendu 200"; exit 1; }
journal "jeton du lecteur prêt"
echo "$JETON"