From b00c39277bf5ca18da53300a058b29cb989da651 Mon Sep 17 00:00:00 2001 From: Dorian Date: Wed, 23 Sep 2026 14:45:38 +0200 Subject: [PATCH 1/6] fix(backend,airflow,ml): cloture la reconciliation entre les deux sources de lectures (#15) --- apps/backend/app/etl/mock_api_import.py | 32 ++++++++ .../backend/tests/etl/test_mock_api_import.py | 77 +++++++++++++++++++ docs/architecture/00-vue-ensemble.md | 19 +++-- docs/architecture/40-data.md | 37 ++++++++- etl/airflow/dags/mock_api_import.py | 12 ++- etl/airflow/tests/test_dags.py | 2 +- ml/enervision_ml/data.py | 14 +++- ml/tests/conftest.py | 43 +++++++++-- ml/tests/test_data_integration.py | 54 ++++++++++++- 9 files changed, 269 insertions(+), 21 deletions(-) diff --git a/apps/backend/app/etl/mock_api_import.py b/apps/backend/app/etl/mock_api_import.py index 65fcea5..c2c105c 100644 --- a/apps/backend/app/etl/mock_api_import.py +++ b/apps/backend/app/etl/mock_api_import.py @@ -19,6 +19,7 @@ from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine from app.core.config import get_settings +from app.etl.historical_import import SOURCE_NAME as SOURCE_CSV SOURCE_HISTORY = "api_history" @@ -244,6 +245,35 @@ def build_reading_row( } +# `uq_reading_source` autorise deux lignes au même (site_id, timestamp) dès que `source` diffère : +# sans ce garde-fou, importer une fenêtre déjà couverte par le dataset historique (source='csv') +# dupliquerait silencieusement chaque point plutôt que de lever une erreur, et rien côté lecture +# (pipeline ML, GET /readings) ne saurait laquelle des deux lectures retenir. +OVERLAP_CHECK = text( + "SELECT count(*) FROM reading WHERE source = :source_csv " + "AND timestamp >= :start_time AND timestamp < :end_time" +) + + +async def refuse_if_overlaps_historical_dataset( + connection: AsyncConnection, + start_time: datetime, + end_time: datetime, +) -> None: + resultat = await connection.execute( + OVERLAP_CHECK, + {"source_csv": SOURCE_CSV, "start_time": start_time, "end_time": end_time}, + ) + nombre = resultat.scalar_one() + + if nombre > 0: + raise ValueError( + f"La fenêtre [{start_time.isoformat()}, {end_time.isoformat()}) recouvre " + f"{nombre} lecture(s) déjà importée(s) du dataset historique (source='{SOURCE_CSV}') : " + "import refusé pour éviter un doublon inter-source." + ) + + # Le conflit vise l'index unique uq_reading_source plutôt que la table entière : sans cible # nommée, DO NOTHING avalerait aussi une violation de clé primaire. READING_INSERT = text( @@ -345,6 +375,8 @@ async def import_mock_api_history( try: async with engine.begin() as connection: + await refuse_if_overlaps_historical_dataset(connection, start_time, end_time) + await upsert_sites( connection, sites, diff --git a/apps/backend/tests/etl/test_mock_api_import.py b/apps/backend/tests/etl/test_mock_api_import.py index fa654bc..562699f 100644 --- a/apps/backend/tests/etl/test_mock_api_import.py +++ b/apps/backend/tests/etl/test_mock_api_import.py @@ -357,6 +357,31 @@ async def test_upsert_sites_with_empty_list_does_nothing() -> None: connection.execute.assert_not_awaited() +async def test_refuse_if_overlaps_historical_dataset_lets_a_clear_window_through() -> None: + connection = AsyncMock() + connection.execute.return_value.scalar_one = MagicMock(return_value=0) + + await mock_api_import.refuse_if_overlaps_historical_dataset( + connection, + datetime.fromisoformat("2026-01-01T00:00:00"), + datetime.fromisoformat("2026-01-01T01:00:00"), + ) + + connection.execute.assert_awaited_once() + + +async def test_refuse_if_overlaps_historical_dataset_rejects_a_window_already_in_the_csv() -> None: + connection = AsyncMock() + connection.execute.return_value.scalar_one = MagicMock(return_value=5) + + with pytest.raises(ValueError, match="doublon inter-source"): + await mock_api_import.refuse_if_overlaps_historical_dataset( + connection, + datetime.fromisoformat("2023-06-15T12:00:00"), + datetime.fromisoformat("2023-06-15T13:00:00"), + ) + + async def test_import_mock_api_history_dry_run_does_not_write( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -454,6 +479,7 @@ async def test_import_mock_api_history_loads_data( ) connection = AsyncMock() + connection.execute.return_value.scalar_one = MagicMock(return_value=0) transaction_context = MagicMock() transaction_context.__aenter__ = AsyncMock( @@ -502,7 +528,58 @@ async def test_import_mock_api_history_loads_data( [make_site()], ) + assert connection.execute.await_count == 2 + dernier_appel = connection.execute.await_args_list[-1] + assert dernier_appel.args[0] is READING_INSERT + engine.dispose.assert_awaited_once() + + +async def test_import_mock_api_history_refuses_when_it_overlaps_the_historical_dataset( + monkeypatch: pytest.MonkeyPatch, +) -> None: + def handler(request: Request) -> Response: + if request.url.path == "/api/v1/sites": + return Response(status_code=200, json=[make_site()]) + + if request.url.path == "/api/v1/readings": + return Response(status_code=200, json=[make_reading()]) + + return Response(status_code=404) + + client = AsyncClient(transport=MockTransport(handler), base_url="https://mock.test") + + monkeypatch.setattr(mock_api_import, "create_mock_api_client", lambda: client) + monkeypatch.setattr( + mock_api_import, + "get_settings", + lambda: SimpleNamespace(database_url="postgresql+asyncpg://test:test@localhost/test"), + ) + + connection = AsyncMock() + connection.execute.return_value.scalar_one = MagicMock(return_value=3) + + transaction_context = MagicMock() + transaction_context.__aenter__ = AsyncMock(return_value=connection) + transaction_context.__aexit__ = AsyncMock(return_value=None) + + engine = MagicMock() + engine.begin.return_value = transaction_context + engine.dispose = AsyncMock() + + monkeypatch.setattr(mock_api_import, "create_async_engine", MagicMock(return_value=engine)) + upsert_sites_mock = AsyncMock() + monkeypatch.setattr(mock_api_import, "upsert_sites", upsert_sites_mock) + + with pytest.raises(ValueError, match="doublon inter-source"): + await mock_api_import.import_mock_api_history( + start_time=datetime.fromisoformat("2023-06-15T12:00:00"), + end_time=datetime.fromisoformat("2023-06-15T13:00:00"), + limit=60, + dry_run=False, + ) + connection.execute.assert_awaited_once() + upsert_sites_mock.assert_not_awaited() engine.dispose.assert_awaited_once() diff --git a/docs/architecture/00-vue-ensemble.md b/docs/architecture/00-vue-ensemble.md index 5d4ddb6..1aece1f 100644 --- a/docs/architecture/00-vue-ensemble.md +++ b/docs/architecture/00-vue-ensemble.md @@ -70,11 +70,17 @@ 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 [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 : six 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 -génération des recommandations (issue #116), `historical_import` pour le dataset historique -(issue #119) et `mock_api_import` pour l'ingestion horaire de l'API Mock (issue #15). -La réconciliation globale des données provenant des deux sources reste à compléter dans l'issue #15. +génération des recommandations (issue #116), un pour la surveillance de dérive (issue #45), +`historical_import` pour le dataset historique (issue #119) et `mock_api_import` pour l'ingestion +horaire de l'API Mock (issue #15). La réconciliation entre les deux sources de lectures (issue +#15) est tranchée : le trou entre la fin de l'historique (31/12/2024) et le début de l'ingestion +API Mock est accepté comme définitivement perdu, aucune mesure réelle n'existant pour cette +période. `mock_api_import` refuse toute fenêtre qui recouvrirait des lectures déjà importées du +CSV plutôt que de laisser les deux sources dupliquer silencieusement un même instant, et le +pipeline ML déduplique par construction (`DISTINCT ON`, source `csv` préférée) au cas où un +recouvrement se produirait malgré tout — voir [40-data.md](40-data.md). Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun collecteur ne vient le lire. @@ -89,7 +95,7 @@ collecteur ne vient le lire. | 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 0013](../adr/0013-surveillance-de-derive-dans-le-backend.md) | | 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 | -| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Six DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes`, `historical_import`, `mock_api_import` et `derive` (quotidien, surveillance de dérive). L'import historique reste manuel et l'import API Mock s'exécute chaque heure. La réconciliation globale des deux sources reste à compléter dans l'issue #15. | +| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Six DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes`, `historical_import`, `mock_api_import` et `derive` (quotidien, surveillance de dérive). L'import historique reste manuel et l'import API Mock s'exécute chaque heure. Réconciliation entre les deux sources (issue #15) : trou temporel accepté, recouvrement refusé à l'ingestion et dédupliqué en défense côté ML, voir [40-data.md](40-data.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 @@ -99,7 +105,8 @@ Statut : `En cours`. **Le chemin de lecture tourne** entre la base, l'API et le dataset CSV/JSON sur déclenchement manuel et `mock_api_import` collecte chaque heure les mesures de l'API Mock. Les DAGs `ml_train` et `ml_score` (issue #115), `alertes` (issue #116) et `derive` (issue #45) portent le pipeline ML, la détection d'alertes et la surveillance de dérive. La -réconciliation globale des données provenant des deux sources reste à compléter dans l'issue #15. +réconciliation entre les deux sources de lectures (issue #15) est close : voir +[40-data.md](40-data.md) pour le détail du garde-fou d'ingestion et de la déduplication ML. ```mermaid sequenceDiagram diff --git a/docs/architecture/40-data.md b/docs/architecture/40-data.md index 335bafd..9ca6706 100644 --- a/docs/architecture/40-data.md +++ b/docs/architecture/40-data.md @@ -441,6 +441,16 @@ Les paramètres de ligne de commande disponibles pour l'import sont : --dry-run ``` +**Piège sur `--limit`** : l'API ne renvoie pas un flux à un rythme naturel, elle répartit +exactement `limit` lectures, espacées uniformément, sur toute la fenêtre `[start_time, end_time)` +demandée. Une fenêtre d'une heure avec `limit=1000` renvoie donc 1000 lectures espacées de 3,6 +secondes à l'intérieur de cette heure, pas une lecture horaire — vérifié empiriquement en +interrogeant directement l'API. Le seul réglage qui produise une lecture par heure, alignée sur +l'heure et cohérente avec le grain horaire du reste du schéma (`period_minutes=60`, historique +CSV à une ligne par heure), est `limit` = nombre d'heures de la fenêtre. Le DAG `mock_api_import` +interroge toujours une fenêtre d'1h (`interval=timedelta(hours=1)`, voir +[10-infra.md](10-infra.md)), donc `limit=1`. + ### Flux d'ingestion API Mock ```text @@ -492,7 +502,7 @@ réponse est donc traitée comme une entrée hostile, conformément à API10 dan [la traçabilité OWASP](owasp-traceabilite.md). Le risque premier n'est pas la fausse alerte, c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction. -Quatre garde-fous, tous dans `mock_api_import.py` : +Cinq garde-fous, tous dans `mock_api_import.py` : | Garde-fou | Mise en œuvre | |---|---| @@ -500,6 +510,7 @@ Quatre garde-fous, tous dans `mock_api_import.py` : | Taille de tableau plafonnée | `MAX_SITES` sites, et au plus `--limit` mesures par site | | Bornes physiques | `PHYSICAL_BOUNDS`, une plage par grandeur | | Frontière d'anti-corruption | `build_site_row()` et `build_reading_row()`, qui ne recopient que les champs attendus | +| Refus de recouvrir l'historique | `refuse_if_overlaps_historical_dataset()`, voir ci-dessous | Une valeur hors bornes, d'un type inattendu, `NaN` ou infinie devient `NULL`. Elle laisse sa trace dans `null_reasons` sous la forme `out_of_physical_bounds:`, et `data_quality` @@ -510,6 +521,30 @@ d'origine intacte : rien n'est perdu, seule son exploitation est bornée. Le plafond de taille s'applique après désérialisation de la réponse. Borner le corps HTTP lui-même demanderait une lecture en flux, et reste à faire. +### Réconciliation entre les deux sources (issue #15) + +`historical_import` (source `csv`) et `mock_api_import` (source `api_history`) écrivent toutes +deux dans `reading`. Trois décisions ferment cette réconciliation : + +- **Le trou temporel est accepté.** Le dataset historique s'arrête au 31/12/2024, et + `mock_api_import` n'importe que l'heure précédant chaque déclenchement : rien ne comble + automatiquement la période intermédiaire, et rien ne le pourra jamais — aucune mesure réelle + n'existe pour ces instants. +- **Le recouvrement est refusé à l'ingestion.** `uq_reading_source` autorise deux lignes au même + `(site_id, timestamp)` dès que `source` diffère : rien dans le schéma n'empêche donc un import + Mock API manuel avec une fenêtre passée (le script accepte `--start-time`/`--end-time` + arbitraires) de dupliquer un point déjà couvert par le CSV. `import_mock_api_history()` appelle + `refuse_if_overlaps_historical_dataset()` avant toute écriture : si la fenêtre demandée recouvre + au moins une lecture `source='csv'`, l'import est refusé (`ValueError`) plutôt que d'écrire un + doublon inter-source silencieux. +- **Le pipeline ML déduplique en défense.** Le garde-fou ci-dessus protège l'ingestion, pas + la lecture : si un recouvrement se produisait malgré tout (import direct en base, contournement + du script), `ml/enervision_ml/data.py` ne doit pas casser silencieusement l'hypothèse de + `build_features` (« une ligne par `(site_id, timestamp)` »). `load_from_database()` et + `load_recent_from_database()` utilisent donc `SELECT DISTINCT ON (site_id, timestamp)`, `source + = 'csv'` gagnant sur `'api_history'` en cas d'égalité — l'historique étant une source vérifiée, + l'API Mock une entrée hostile (cf. ci-dessus). + ### Qualité des données de l'API Mock Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion. diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index f0043d0..e649be5 100644 --- a/etl/airflow/dags/mock_api_import.py +++ b/etl/airflow/dags/mock_api_import.py @@ -18,9 +18,15 @@ from airflow.timetables.trigger import CronTriggerTimetable # Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" -# Le pipeline backend et l'API acceptent au maximum 1 000 lectures par site. -# Cette marge évite de perdre silencieusement une lecture si une heure en contient plus de 60. -LIMITE_LECTURES = 1000 +# L'API Mock ne renvoie pas un flux au rythme naturel : elle répartit exactement `limit` +# lectures, espacées uniformément, sur toute la fenêtre demandée (vérifié empiriquement : +# une fenêtre d'1h avec `limit=1000` renvoie 1000 lectures espacées de 3,6s à l'intérieur de +# cette heure, pas une lecture horaire). Comme la fenêtre de ce DAG est toujours 1h +# (`interval=timedelta(hours=1)` ci-dessous), `limit=1` est ce qui produit une lecture par +# heure, alignée sur l'heure, cohérente avec le grain horaire du reste du schéma +# (`period_minutes=60`, historique CSV à 1 ligne/heure). Un `limit` plus grand ici fabriquerait +# des lectures infra-horaires, incompatibles avec les lags positionnels de `build_features`. +LIMITE_LECTURES = 1 # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # inférieure au pas horaire du DAG. diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index d03dd97..4e4e077 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -118,7 +118,7 @@ def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande - assert "--limit 1000" in commande + assert "--limit 1" in commande @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) diff --git a/ml/enervision_ml/data.py b/ml/enervision_ml/data.py index e74da91..3fa64d2 100644 --- a/ml/enervision_ml/data.py +++ b/ml/enervision_ml/data.py @@ -47,9 +47,15 @@ NUMERIC_COLUMNS = [ # `bool` : `astype(bool)` ferait un `True` d'une absence, et les deux chargeurs divergeraient. FLAG_COLUMNS = ["is_working_hours"] +# `uq_reading_source` autorise deux lignes au meme (site_id, timestamp) des que `source` differe +# (cf. `app/etl/mock_api_import.py`, qui refuse desormais d'importer une fenetre deja couverte par +# le CSV, mais ne protege pas le sens inverse). `build_features` suppose une ligne par +# (site_id, timestamp) sans doublon : le `DISTINCT ON` l'impose plutot que de la supposer. +# 'csv' gagne sur 'api_history' en cas de recouvrement, l'historique etant une source verifiee +# alors que l'API Mock est traitee comme une entree hostile (cf. OWASP API10). _READING_QUERY = text( """ - SELECT + SELECT DISTINCT ON (r.site_id, r.timestamp) r.site_id, r.timestamp, r.consumption_kwh, @@ -61,14 +67,14 @@ _READING_QUERY = text( s.capacity_kw FROM reading r JOIN site s ON s.site_id = r.site_id - ORDER BY r.site_id, r.timestamp + ORDER BY r.site_id, r.timestamp, (r.source = 'csv') DESC, r.reading_id DESC """ ) _RECENT_READING_QUERY = text( """ - SELECT + SELECT DISTINCT ON (r.site_id, r.timestamp) r.site_id, r.timestamp, r.consumption_kwh, @@ -81,7 +87,7 @@ _RECENT_READING_QUERY = text( FROM reading r JOIN site s ON s.site_id = r.site_id WHERE r.timestamp >= :since AND r.timestamp <= :until - ORDER BY r.site_id, r.timestamp + ORDER BY r.site_id, r.timestamp, (r.source = 'csv') DESC, r.reading_id DESC """ ) diff --git a/ml/tests/conftest.py b/ml/tests/conftest.py index 33f65b9..b2b7a49 100644 --- a/ml/tests/conftest.py +++ b/ml/tests/conftest.py @@ -16,7 +16,7 @@ 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 typing import Any, cast from uuid import uuid4 import lightgbm as lgb @@ -44,20 +44,30 @@ _INSERT_SITE = text( """ ) -# `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. +# `source = 'api_history'` impose `dataset_id IS NULL` (ck_reading_dataset_source) : le defaut +# `dataset_id=None` evite de creer une ligne `dataset` pour la plupart des tests. `source='csv'` +# impose l'inverse, d'ou `insere_dataset()` quand un test a besoin de cette source precise. +# `raw_data` est NOT NULL, d'ou le litteral jsonb. _INSERT_READING = text( """ INSERT INTO reading ( - site_id, timestamp, source, consumption_kwh, temperature_celsius, + site_id, timestamp, source, dataset_id, consumption_kwh, temperature_celsius, humidity_percent, solar_irradiance_wm2, is_working_hours, raw_data ) VALUES ( - :site_id, :timestamp, :source, :consumption_kwh, :temperature_celsius, + :site_id, :timestamp, :source, :dataset_id, :consumption_kwh, :temperature_celsius, :humidity_percent, :solar_irradiance_wm2, :is_working_hours, '{}'::jsonb ) """ ) +_INSERT_DATASET = text( + """ + INSERT INTO dataset (dataset_name, archive_sha256, storage_uri, source_timezone, metadata) + VALUES (:dataset_name, :archive_sha256, :storage_uri, 'UTC', '{}'::jsonb) + RETURNING dataset_id + """ +) + _SELECT_PREDICTIONS = text( """ SELECT target_at, predicted_value, status, failure_reason, model_reference @@ -111,6 +121,25 @@ def insere_site( return site_id +def insere_dataset(connexion: Connection) -> int: + """Ligne `dataset` minimale, requise pour inserer une lecture `source='csv'` + + (`ck_reading_dataset_source` impose `dataset_id IS NOT NULL` pour cette seule source). + """ + marque = uuid4().hex + return cast( + int, + connexion.execute( + _INSERT_DATASET, + { + "dataset_name": f"jeu de test {marque}", + "archive_sha256": marque.rjust(64, "0"), + "storage_uri": f"file:///test/{marque}.csv", + }, + ).scalar_one(), + ) + + def insere_lectures( connexion: Connection, site_id: str, @@ -119,6 +148,7 @@ def insere_lectures( fin: datetime, valeur: float = 50.0, source: str = "api_history", + dataset_id: int | None = None, is_working_hours: bool | None = True, ) -> list[datetime]: """Grille horaire contigue finissant a `fin`, incluse. @@ -134,6 +164,7 @@ def insere_lectures( "site_id": site_id, "timestamp": instant, "source": source, + "dataset_id": dataset_id, "consumption_kwh": valeur + math.sin(rang / 12.0) * 10.0, "temperature_celsius": 15.0, "humidity_percent": 50.0, @@ -153,6 +184,7 @@ def insere_lecture( instant: datetime, consumption_kwh: float | None = 50.0, source: str = "api_history", + dataset_id: int | None = None, is_working_hours: bool | None = True, ) -> None: """Une lecture isolee, quand le test pilote sa valeur plutot que sa forme.""" @@ -162,6 +194,7 @@ def insere_lecture( "site_id": site_id, "timestamp": instant, "source": source, + "dataset_id": dataset_id, "consumption_kwh": consumption_kwh, "temperature_celsius": 15.0, "humidity_percent": 50.0, diff --git a/ml/tests/test_data_integration.py b/ml/tests/test_data_integration.py index 3b5ea1a..9555d02 100644 --- a/ml/tests/test_data_integration.py +++ b/ml/tests/test_data_integration.py @@ -11,7 +11,7 @@ from enervision_ml.data import ( load_from_database, load_recent_from_database, ) -from tests.conftest import ANCRAGE, insere_lecture, insere_lectures, insere_site +from tests.conftest import ANCRAGE, insere_dataset, insere_lecture, insere_lectures, insere_site pytestmark = pytest.mark.integration @@ -39,6 +39,58 @@ def test_load_from_database_joins_the_site_attributes_to_every_reading( assert set(mien["capacity_kw"]) == {250.0} +def test_load_from_database_deduplicates_two_sources_at_the_same_instant( + connexion_ml: Connection, +) -> None: + # `uq_reading_source` autorise deux lignes au meme (site_id, timestamp) des que `source` + # differe : le garde-fou vit dans `mock_api_import.py`, pas dans le schema. Le chargeur ML + # doit donc imposer lui-meme "une ligne par (site_id, timestamp)", pas la supposer. + site_id = insere_site(connexion_ml) + dataset_id = insere_dataset(connexion_ml) + insere_lecture( + connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history" + ) + insere_lecture( + connexion_ml, + site_id, + instant=ANCRAGE, + consumption_kwh=99.0, + source="csv", + dataset_id=dataset_id, + ) + + frame = load_from_database(connexion_ml) + + mien = frame[frame["site_id"] == site_id] + assert len(mien) == 1 + assert mien["consumption_kwh"].iloc[0] == 99.0 + + +def test_load_recent_from_database_prefers_csv_when_two_sources_share_an_instant( + connexion_ml: Connection, +) -> None: + site_id = insere_site(connexion_ml) + dataset_id = insere_dataset(connexion_ml) + insere_lecture( + connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history" + ) + insere_lecture( + connexion_ml, + site_id, + instant=ANCRAGE, + consumption_kwh=99.0, + source="csv", + dataset_id=dataset_id, + ) + + frame = load_recent_from_database( + connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3) + ) + + assert len(frame) == 1 + assert frame["consumption_kwh"].iloc[0] == 99.0 + + def test_load_recent_from_database_excludes_readings_before_the_since_bound( connexion_ml: Connection, ) -> None: From b78322bd6178ce57cf215b63e5ff6352741e7096 Mon Sep 17 00:00:00 2001 From: Dorian Date: Wed, 23 Sep 2026 15:19:39 +0200 Subject: [PATCH 2/6] fix(backend,airflow,ml): applique les corrections de revue sur la PR #162 --- apps/backend/app/etl/mock_api_import.py | 159 +++++++---- .../backend/tests/etl/test_mock_api_import.py | 252 +++++++++++------- docs/architecture/00-vue-ensemble.md | 2 +- docs/architecture/40-data.md | 59 ++-- etl/airflow/dags/mock_api_import.py | 23 +- etl/airflow/tests/test_dags.py | 3 +- ml/tests/test_data_integration.py | 19 +- 7 files changed, 328 insertions(+), 189 deletions(-) diff --git a/apps/backend/app/etl/mock_api_import.py b/apps/backend/app/etl/mock_api_import.py index c2c105c..ca24d2b 100644 --- a/apps/backend/app/etl/mock_api_import.py +++ b/apps/backend/app/etl/mock_api_import.py @@ -1,17 +1,19 @@ # Contrainte : la réponse de l'API Mock est une entrée hostile, pas une source de confiance. # Voir OWASP API10 dans docs/architecture/owasp-traceabilite.md. Rien de ce qu'elle renvoie # n'atteint la base sans passer par build_site_row() ou build_reading_row() : seuls les champs -# attendus sont recopiés, les grandeurs physiques sont bornées par PHYSICAL_BOUNDS et la taille -# des tableaux est plafonnée par MAX_SITES et par --limit. Une valeur hors bornes devient NULL -# et laisse sa trace dans null_reasons plutôt que de lever : le mock émet des anomalies par -# construction, et raw_data conserve de toute façon la réponse d'origine intacte. +# attendus sont recopiés, les grandeurs physiques sont bornées par PHYSICAL_BOUNDS, la taille des +# tableaux est plafonnée par MAX_SITES et par limit_for_window() (dérivé de la fenêtre, jamais +# fourni par l'appelant), et les lectures dont le timestamp déborde de la fenêtre demandée sont +# écartées (fetch_readings). Une valeur hors bornes devient NULL et laisse sa trace dans +# null_reasons plutôt que de lever : le mock émet des anomalies par construction, et raw_data +# conserve de toute façon la réponse d'origine intacte. from __future__ import annotations import argparse import asyncio import json -from datetime import datetime +from datetime import UTC, datetime from typing import Any import httpx @@ -182,6 +184,21 @@ async def upsert_sites( ) +def _timestamp_in_window( + reading: dict[str, Any], start_time: datetime, end_time: datetime +) -> bool: + valeur = reading.get("timestamp") + if not isinstance(valeur, str): + return False + + try: + instant = parse_datetime(valeur) + except ValueError: + return False + + return start_time <= instant < end_time + + async def fetch_readings( client: httpx.AsyncClient, site_id: str, @@ -209,7 +226,22 @@ async def fetch_readings( if len(payload) > limit: raise ValueError(f"La réponse /api/v1/readings dépasse la limite demandée de {limit}.") - return payload + # Le garde-fou `refuse_if_overlaps_historical_dataset` ne vérifie que la fenêtre demandée : + # une réponse (bug du mock, ou hostile) dont les `timestamp` débordent de + # `[start_time, end_time)` contournerait ce contrôle et écrirait exactement le doublon + # inter-source qu'il doit empêcher. Écarter ces lectures ici rend le contrôle par fenêtre + # suffisant. + dans_la_fenetre = [ + lecture + for lecture in payload + if isinstance(lecture, dict) and _timestamp_in_window(lecture, start_time, end_time) + ] + + if len(dans_la_fenetre) != len(payload): + ecartees = len(payload) - len(dans_la_fenetre) + print(f"{site_id}: {ecartees} lecture(s) hors fenêtre écartée(s).") + + return dans_la_fenetre def build_reading_row( @@ -247,8 +279,9 @@ def build_reading_row( # `uq_reading_source` autorise deux lignes au même (site_id, timestamp) dès que `source` diffère : # sans ce garde-fou, importer une fenêtre déjà couverte par le dataset historique (source='csv') -# dupliquerait silencieusement chaque point plutôt que de lever une erreur, et rien côté lecture -# (pipeline ML, GET /readings) ne saurait laquelle des deux lectures retenir. +# dupliquerait silencieusement chaque point plutôt que de lever une erreur. Ce garde-fou protège +# l'ingestion ; il ne dit rien de la lecture (`GET /readings` renvoie les deux lignes en cas de +# doublon malgré tout, cf. la section réconciliation de 40-data.md). OVERLAP_CHECK = text( "SELECT count(*) FROM reading WHERE source = :source_csv " "AND timestamp >= :start_time AND timestamp < :end_time" @@ -332,41 +365,41 @@ def build_reading_batch( return [build_reading_row(reading) for reading in readings] +def limit_for_window(start_time: datetime, end_time: datetime) -> int: + """Nombre de lectures à demander pour que l'API Mock en rende une par heure, alignée. + + L'API ne renvoie pas un flux à un rythme naturel : elle répartit exactement `limit` lectures, + espacées uniformément, sur toute la fenêtre `[start_time, end_time)` demandée (vérifié + empiriquement). `limit` = nombre d'heures de la fenêtre est donc le seul réglage cohérent avec + le grain horaire du reste du schéma (`period_minutes=60`, historique CSV à une ligne/heure) ; + un `limit` plus grand fabriquerait des lectures infra-horaires, incompatibles avec les lags + positionnels de `build_features`. La fenêtre doit donc couvrir un nombre entier d'heures. + """ + duree = end_time - start_time + heures, reste = divmod(duree.total_seconds(), 3600) + + if reste != 0: + raise ValueError( + "La fenêtre doit couvrir un nombre entier d'heures pour obtenir une lecture par " + f"heure alignée : [{start_time.isoformat()}, {end_time.isoformat()}) n'en couvre pas." + ) + + if heures > MAX_LIMIT: + raise ValueError( + f"La fenêtre demandée couvre {int(heures)}h, au-delà du plafond de {MAX_LIMIT} " + "lectures accepté par l'API Mock." + ) + + return int(heures) + + async def import_mock_api_history( start_time: datetime, end_time: datetime, - limit: int, dry_run: bool, ) -> None: settings = get_settings() - - async with create_mock_api_client() as client: - sites = await fetch_sites(client) - - print(f"Sites récupérés : {len(sites)}") - - all_readings: list[dict[str, Any]] = [] - - for site in sites: - site_id = read_text(site, "site_id") - - readings = await fetch_readings( - client=client, - site_id=site_id, - start_time=start_time, - end_time=end_time, - limit=limit, - ) - - print(f"{site_id}: {len(readings)} lectures") - - all_readings.extend(readings) - - print(f"Lectures récupérées : {len(all_readings)}") - - if dry_run: - print("Dry-run terminé : aucune donnée écrite.") - return + limit = limit_for_window(start_time, end_time) engine = create_async_engine( str(settings.database_url), @@ -374,9 +407,40 @@ async def import_mock_api_history( ) try: - async with engine.begin() as connection: + # Garde-fou d'abord, y compris en dry-run : il est en lecture seule, et annoncer un + # succès pour une fenêtre que l'import réel refusera serait trompeur. + async with engine.connect() as connection: await refuse_if_overlaps_historical_dataset(connection, start_time, end_time) + async with create_mock_api_client() as client: + sites = await fetch_sites(client) + + print(f"Sites récupérés : {len(sites)}") + + all_readings: list[dict[str, Any]] = [] + + for site in sites: + site_id = read_text(site, "site_id") + + readings = await fetch_readings( + client=client, + site_id=site_id, + start_time=start_time, + end_time=end_time, + limit=limit, + ) + + print(f"{site_id}: {len(readings)} lectures") + + all_readings.extend(readings) + + print(f"Lectures récupérées : {len(all_readings)}") + + if dry_run: + print("Dry-run terminé : aucune donnée écrite.") + return + + async with engine.begin() as connection: await upsert_sites( connection, sites, @@ -397,7 +461,14 @@ async def import_mock_api_history( def parse_datetime(value: str) -> datetime: - return datetime.fromisoformat(value.replace("Z", "+00:00")) + # Sans fuseau, l'API le traite comme reçu, telle quelle, mais l'encodeur `timestamptz` + # d'asyncpg lirait un datetime naif dans le fuseau *local du processus* (correct dans le + # conteneur Airflow en UTC, décalé de 1-2h pour un import manuel lancé depuis un poste en + # Europe/Paris). Poser `tzinfo=UTC` explicitement, même pattern que `_vers_utc()` dans + # `app/services/reading.py`, garantit que la borne envoyée à l'API et celle comparée en SQL + # (refuse_if_overlaps_historical_dataset) désignent le même instant. + instant = datetime.fromisoformat(value.replace("Z", "+00:00")) + return instant if instant.tzinfo is not None else instant.replace(tzinfo=UTC) def parse_args() -> argparse.Namespace: @@ -415,12 +486,6 @@ def parse_args() -> argparse.Namespace: type=parse_datetime, ) - parser.add_argument( - "--limit", - type=int, - default=MAX_LIMIT, - ) - parser.add_argument( "--dry-run", action="store_true", @@ -432,9 +497,6 @@ def parse_args() -> argparse.Namespace: def main() -> None: args = parse_args() - if args.limit < 1 or args.limit > MAX_LIMIT: - raise ValueError(f"--limit doit être compris entre 1 et {MAX_LIMIT}.") - if args.start_time >= args.end_time: raise ValueError("--start-time doit être antérieur à --end-time.") @@ -442,7 +504,6 @@ def main() -> None: import_mock_api_history( start_time=args.start_time, end_time=args.end_time, - limit=args.limit, dry_run=args.dry_run, ) ) diff --git a/apps/backend/tests/etl/test_mock_api_import.py b/apps/backend/tests/etl/test_mock_api_import.py index 562699f..132a718 100644 --- a/apps/backend/tests/etl/test_mock_api_import.py +++ b/apps/backend/tests/etl/test_mock_api_import.py @@ -1,6 +1,6 @@ import json import sys -from datetime import datetime +from datetime import UTC, datetime from types import SimpleNamespace from typing import Any from unittest.mock import AsyncMock, MagicMock @@ -110,8 +110,8 @@ async def test_fetch_readings_sends_expected_query_parameters() -> None: transport = MockTransport(handler) - start_time = datetime.fromisoformat("2024-06-15T12:00:00") - end_time = datetime.fromisoformat("2024-06-15T13:00:00") + start_time = datetime.fromisoformat("2024-06-15T12:00:00+00:00") + end_time = datetime.fromisoformat("2024-06-15T13:00:00+00:00") async with AsyncClient( transport=transport, @@ -127,11 +127,61 @@ async def test_fetch_readings_sends_expected_query_parameters() -> None: assert len(readings) == 1 assert captured_params["site_id"] == "SITE001" - assert captured_params["start_time"] == "2024-06-15T12:00:00" - assert captured_params["end_time"] == "2024-06-15T13:00:00" + assert captured_params["start_time"] == "2024-06-15T12:00:00+00:00" + assert captured_params["end_time"] == "2024-06-15T13:00:00+00:00" assert captured_params["limit"] == "60" +async def test_fetch_readings_discards_a_reading_outside_the_requested_window() -> None: + # Le garde-fou `refuse_if_overlaps_historical_dataset` ne vérifie que la fenêtre demandée : + # une réponse dont un `timestamp` déborde de `[start_time, end_time)` (bug du mock, ou + # hostile) contournerait ce contrôle si elle atteignait la base telle quelle. + dans_la_fenetre = make_reading() + dans_la_fenetre["timestamp"] = "2024-06-15T12:00:00Z" + + hors_fenetre = make_reading() + hors_fenetre["timestamp"] = "2023-01-01T00:00:00Z" + + def handler(request: Request) -> Response: + return Response(status_code=200, json=[dans_la_fenetre, hors_fenetre]) + + async with AsyncClient( + transport=MockTransport(handler), + base_url="https://mock.test", + ) as client: + readings = await fetch_readings( + client=client, + site_id="SITE001", + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), + limit=2, + ) + + assert readings == [dans_la_fenetre] + + +async def test_fetch_readings_discards_a_reading_with_an_unparseable_timestamp() -> None: + invalide = make_reading() + invalide["timestamp"] = "pas une date" + + def handler(request: Request) -> Response: + return Response(status_code=200, json=[invalide]) + + async with AsyncClient( + transport=MockTransport(handler), + base_url="https://mock.test", + ) as client: + readings = await fetch_readings( + client=client, + site_id="SITE001", + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), + limit=1, + ) + + assert readings == [] + + async def test_fetch_readings_rejects_non_list_response() -> None: def handler(request: Request) -> Response: return Response( @@ -152,8 +202,8 @@ async def test_fetch_readings_rejects_non_list_response() -> None: await fetch_readings( client=client, site_id="SITE001", - start_time=datetime.fromisoformat("2024-06-15T12:00:00"), - end_time=datetime.fromisoformat("2024-06-15T13:00:00"), + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), limit=60, ) @@ -175,8 +225,8 @@ async def test_fetch_readings_raises_on_http_error() -> None: await fetch_readings( client=client, site_id="SITE999", - start_time=datetime.fromisoformat("2024-06-15T12:00:00"), - end_time=datetime.fromisoformat("2024-06-15T13:00:00"), + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), limit=60, ) @@ -363,8 +413,8 @@ async def test_refuse_if_overlaps_historical_dataset_lets_a_clear_window_through await mock_api_import.refuse_if_overlaps_historical_dataset( connection, - datetime.fromisoformat("2026-01-01T00:00:00"), - datetime.fromisoformat("2026-01-01T01:00:00"), + datetime.fromisoformat("2026-01-01T00:00:00+00:00"), + datetime.fromisoformat("2026-01-01T01:00:00+00:00"), ) connection.execute.assert_awaited_once() @@ -377,11 +427,57 @@ async def test_refuse_if_overlaps_historical_dataset_rejects_a_window_already_in with pytest.raises(ValueError, match="doublon inter-source"): await mock_api_import.refuse_if_overlaps_historical_dataset( connection, - datetime.fromisoformat("2023-06-15T12:00:00"), - datetime.fromisoformat("2023-06-15T13:00:00"), + datetime.fromisoformat("2023-06-15T12:00:00+00:00"), + datetime.fromisoformat("2023-06-15T13:00:00+00:00"), ) +def test_limit_for_window_returns_one_per_hour() -> None: + limite = mock_api_import.limit_for_window( + datetime.fromisoformat("2026-09-02T12:00:00+00:00"), + datetime.fromisoformat("2026-09-23T12:00:00+00:00"), + ) + + assert limite == 21 * 24 + + +def test_limit_for_window_rejects_a_partial_hour() -> None: + with pytest.raises(ValueError, match="nombre entier d'heures"): + mock_api_import.limit_for_window( + datetime.fromisoformat("2026-09-02T12:00:00+00:00"), + datetime.fromisoformat("2026-09-02T12:30:00+00:00"), + ) + + +def test_limit_for_window_rejects_a_window_above_the_api_cap() -> None: + with pytest.raises(ValueError, match="au-delà du plafond"): + mock_api_import.limit_for_window( + datetime.fromisoformat("2020-01-01T00:00:00+00:00"), + datetime.fromisoformat("2020-03-01T00:00:00+00:00"), + ) + + +def _mock_engine(*, overlap_count: int = 0) -> tuple[MagicMock, AsyncMock]: + """Engine dont `.connect()` (garde-fou) et `.begin()` (écriture) rendent tous deux la même + connexion, dont `scalar_one()` renvoie `overlap_count` : `import_mock_api_history` ouvre + désormais le garde-fou via `.connect()`, y compris en dry-run.""" + connection = AsyncMock() + connection.execute.return_value.scalar_one = MagicMock(return_value=overlap_count) + + def _context() -> MagicMock: + context = MagicMock() + context.__aenter__ = AsyncMock(return_value=connection) + context.__aexit__ = AsyncMock(return_value=None) + return context + + engine = MagicMock() + engine.connect.return_value = _context() + engine.begin.return_value = _context() + engine.dispose = AsyncMock() + + return engine, connection + + async def test_import_mock_api_history_dry_run_does_not_write( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -421,22 +517,24 @@ async def test_import_mock_api_history_dry_run_does_not_write( ), ) - create_engine_mock = MagicMock() + engine, connection = _mock_engine(overlap_count=0) monkeypatch.setattr( mock_api_import, "create_async_engine", - create_engine_mock, + MagicMock(return_value=engine), ) await mock_api_import.import_mock_api_history( - start_time=datetime.fromisoformat("2024-06-15T12:00:00"), - end_time=datetime.fromisoformat("2024-06-15T13:00:00"), - limit=60, + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), dry_run=True, ) - create_engine_mock.assert_not_called() + # Le garde-fou tourne quand même (lecture seule), mais aucune écriture n'a lieu. + connection.execute.assert_awaited_once() + engine.begin.assert_not_called() + engine.dispose.assert_awaited_once() async def test_import_mock_api_history_loads_data( @@ -478,24 +576,8 @@ async def test_import_mock_api_history_loads_data( ), ) - connection = AsyncMock() - connection.execute.return_value.scalar_one = MagicMock(return_value=0) - - transaction_context = MagicMock() - transaction_context.__aenter__ = AsyncMock( - return_value=connection, - ) - transaction_context.__aexit__ = AsyncMock( - return_value=None, - ) - - engine = MagicMock() - engine.begin.return_value = transaction_context - engine.dispose = AsyncMock() - - create_engine_mock = MagicMock( - return_value=engine, - ) + engine, connection = _mock_engine(overlap_count=0) + create_engine_mock = MagicMock(return_value=engine) upsert_sites_mock = AsyncMock() @@ -512,9 +594,8 @@ async def test_import_mock_api_history_loads_data( ) await mock_api_import.import_mock_api_history( - start_time=datetime.fromisoformat("2024-06-15T12:00:00"), - end_time=datetime.fromisoformat("2024-06-15T13:00:00"), - limit=60, + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), dry_run=False, ) @@ -528,6 +609,7 @@ async def test_import_mock_api_history_loads_data( [make_site()], ) + # Un appel pour le garde-fou (via .connect()), un pour READING_INSERT (via .begin()). assert connection.execute.await_count == 2 dernier_appel = connection.execute.await_args_list[-1] assert dernier_appel.args[0] is READING_INSERT @@ -537,49 +619,29 @@ async def test_import_mock_api_history_loads_data( async def test_import_mock_api_history_refuses_when_it_overlaps_the_historical_dataset( monkeypatch: pytest.MonkeyPatch, ) -> None: - def handler(request: Request) -> Response: - if request.url.path == "/api/v1/sites": - return Response(status_code=200, json=[make_site()]) - - if request.url.path == "/api/v1/readings": - return Response(status_code=200, json=[make_reading()]) - - return Response(status_code=404) - - client = AsyncClient(transport=MockTransport(handler), base_url="https://mock.test") - - monkeypatch.setattr(mock_api_import, "create_mock_api_client", lambda: client) monkeypatch.setattr( mock_api_import, "get_settings", lambda: SimpleNamespace(database_url="postgresql+asyncpg://test:test@localhost/test"), ) - connection = AsyncMock() - connection.execute.return_value.scalar_one = MagicMock(return_value=3) - - transaction_context = MagicMock() - transaction_context.__aenter__ = AsyncMock(return_value=connection) - transaction_context.__aexit__ = AsyncMock(return_value=None) - - engine = MagicMock() - engine.begin.return_value = transaction_context - engine.dispose = AsyncMock() - + engine, connection = _mock_engine(overlap_count=3) monkeypatch.setattr(mock_api_import, "create_async_engine", MagicMock(return_value=engine)) - upsert_sites_mock = AsyncMock() - monkeypatch.setattr(mock_api_import, "upsert_sites", upsert_sites_mock) + + # Le garde-fou tourne avant tout appel à l'API Mock : create_mock_api_client() ne doit + # jamais être invoqué pour une fenêtre refusée. + create_client_mock = MagicMock() + monkeypatch.setattr(mock_api_import, "create_mock_api_client", create_client_mock) with pytest.raises(ValueError, match="doublon inter-source"): await mock_api_import.import_mock_api_history( - start_time=datetime.fromisoformat("2023-06-15T12:00:00"), - end_time=datetime.fromisoformat("2023-06-15T13:00:00"), - limit=60, + start_time=datetime.fromisoformat("2023-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2023-06-15T13:00:00+00:00"), dry_run=False, ) connection.execute.assert_awaited_once() - upsert_sites_mock.assert_not_awaited() + create_client_mock.assert_not_called() engine.dispose.assert_awaited_once() @@ -593,6 +655,23 @@ def test_parse_datetime_accepts_z_suffix() -> None: ) +def test_parse_datetime_attaches_utc_to_a_naive_string() -> None: + # `--start-time`/`--end-time` du DAG sont formatés sans fuseau (Jinja `strftime`) : sans ce + # comportement, l'encodeur `timestamptz` d'asyncpg lirait le datetime naïf dans le fuseau + # *local du processus*, pas UTC, et le garde-fou comparerait une autre fenêtre que celle + # envoyée à l'API. + result = mock_api_import.parse_datetime("2024-06-15T12:00:00") + + assert result == datetime.fromisoformat("2024-06-15T12:00:00+00:00") + assert result.tzinfo is UTC + + +def test_parse_datetime_keeps_a_non_utc_offset_as_is() -> None: + result = mock_api_import.parse_datetime("2024-06-15T12:00:00+02:00") + + assert result == datetime.fromisoformat("2024-06-15T12:00:00+02:00") + + def test_parse_args_reads_cli_parameters( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -605,8 +684,6 @@ def test_parse_args_reads_cli_parameters( "2024-06-15T12:00:00Z", "--end-time", "2024-06-15T13:00:00Z", - "--limit", - "60", "--dry-run", ], ) @@ -619,34 +696,9 @@ def test_parse_args_reads_cli_parameters( assert args.end_time == datetime.fromisoformat( "2024-06-15T13:00:00+00:00", ) - assert args.limit == 60 assert args.dry_run is True -def test_main_rejects_limit_out_of_bounds( - monkeypatch: pytest.MonkeyPatch, -) -> None: - monkeypatch.setattr( - sys, - "argv", - [ - "mock_api_import", - "--start-time", - "2024-06-15T12:00:00Z", - "--end-time", - "2024-06-15T13:00:00Z", - "--limit", - "0", - ], - ) - - with pytest.raises( - ValueError, - match="--limit doit être compris entre 1 et 1000", - ): - mock_api_import.main() - - def test_main_rejects_invalid_period( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -659,8 +711,6 @@ def test_main_rejects_invalid_period( "2024-06-15T14:00:00Z", "--end-time", "2024-06-15T13:00:00Z", - "--limit", - "60", ], ) @@ -689,7 +739,6 @@ def test_main_runs_import( lambda: SimpleNamespace( start_time=start_time, end_time=end_time, - limit=60, dry_run=True, ), ) @@ -705,7 +754,6 @@ def test_main_runs_import( import_mock.assert_awaited_once_with( start_time=start_time, end_time=end_time, - limit=60, dry_run=True, ) @@ -750,8 +798,8 @@ async def test_fetch_readings_rejects_a_response_above_the_requested_limit() -> await fetch_readings( client=client, site_id="SITE001", - start_time=datetime.fromisoformat("2024-06-15T12:00:00"), - end_time=datetime.fromisoformat("2024-06-15T13:00:00"), + start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"), + end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"), limit=2, ) diff --git a/docs/architecture/00-vue-ensemble.md b/docs/architecture/00-vue-ensemble.md index bdef60d..7a0666d 100644 --- a/docs/architecture/00-vue-ensemble.md +++ b/docs/architecture/00-vue-ensemble.md @@ -80,7 +80,7 @@ API Mock est accepté comme définitivement perdu, aucune mesure réelle n'exist période. `mock_api_import` refuse toute fenêtre qui recouvrirait des lectures déjà importées du CSV plutôt que de laisser les deux sources dupliquer silencieusement un même instant, et le pipeline ML déduplique par construction (`DISTINCT ON`, source `csv` préférée) au cas où un -recouvrement se produirait malgré tout — voir [40-data.md](40-data.md). +recouvrement se produirait malgré tout, voir [40-data.md](40-data.md). Les liens de la supervision sont en trait plein depuis le 23/09 (issue #26) : Prometheus scrute `/metrics` avec un jeton, Grafana lit Prometheus et, par un rôle en lecture seule, les tables diff --git a/docs/architecture/40-data.md b/docs/architecture/40-data.md index 9ca6706..88dc21d 100644 --- a/docs/architecture/40-data.md +++ b/docs/architecture/40-data.md @@ -437,19 +437,23 @@ Les paramètres de ligne de commande disponibles pour l'import sont : ```text --start-time --end-time ---limit --dry-run ``` -**Piège sur `--limit`** : l'API ne renvoie pas un flux à un rythme naturel, elle répartit -exactement `limit` lectures, espacées uniformément, sur toute la fenêtre `[start_time, end_time)` -demandée. Une fenêtre d'une heure avec `limit=1000` renvoie donc 1000 lectures espacées de 3,6 -secondes à l'intérieur de cette heure, pas une lecture horaire — vérifié empiriquement en -interrogeant directement l'API. Le seul réglage qui produise une lecture par heure, alignée sur -l'heure et cohérente avec le grain horaire du reste du schéma (`period_minutes=60`, historique -CSV à une ligne par heure), est `limit` = nombre d'heures de la fenêtre. Le DAG `mock_api_import` -interroge toujours une fenêtre d'1h (`interval=timedelta(hours=1)`, voir -[10-infra.md](10-infra.md)), donc `limit=1`. +**Piège sur `limit`, corrigé dans le code plutôt que documenté** : l'API ne renvoie pas un flux à +un rythme naturel, elle répartit exactement `limit` lectures, espacées uniformément, sur toute la +fenêtre `[start_time, end_time)` demandée (vérifié empiriquement en interrogeant directement +l'API). Une fenêtre d'une heure avec `limit=1000`, le réglage d'origine, renvoyait donc 1000 +lectures espacées de 3,6 secondes à l'intérieur de cette heure, pas une lecture horaire, +incompatible avec les lags positionnels de `build_features`. Plutôt que documenter la règle +« `limit` = nombre d'heures de la fenêtre » et compter sur chaque appelant pour la respecter, +`limit_for_window()` la porte : `import_mock_api_history()` calcule `limit` depuis la fenêtre +reçue, refuse une fenêtre qui ne couvre pas un nombre entier d'heures, et refuse un intervalle de +plus de 1000 heures (le plafond `limit` de l'API). `--limit` n'existe donc plus côté CLI. Le DAG +`mock_api_import` interroge toujours une fenêtre d'1h (`interval=timedelta(hours=1)`, voir +[10-infra.md](10-infra.md)) : la fenêtre `[:45, :45)` place chaque lecture à :45, pas à :00 (la +première lecture atterrit au début de la fenêtre demandée), un décalage constant sans effet sur +les lags positionnels ni sur les jointures en aval. ### Flux d'ingestion API Mock @@ -528,22 +532,43 @@ deux dans `reading`. Trois décisions ferment cette réconciliation : - **Le trou temporel est accepté.** Le dataset historique s'arrête au 31/12/2024, et `mock_api_import` n'importe que l'heure précédant chaque déclenchement : rien ne comble - automatiquement la période intermédiaire, et rien ne le pourra jamais — aucune mesure réelle - n'existe pour ces instants. + automatiquement la période intermédiaire, et rien ne le pourra jamais, aucune mesure réelle + n'existe pour ces instants. Conséquence pour le ML, pas nouvelle mais que ce trou rend + définitive : `build_features()` calcule ses lags par `shift(n)` positionnel, et `train.py` + n'écarte que les lignes où `lag_168h` est `NaN`. Pour un site présent dans les deux sources, les + 168 premières lectures `api_history` qui suivent le trou héritent donc de lags et de moyennes + glissantes calculés sur décembre 2024 (et tant que l'ingestion a moins de 7 jours, c'est le cas + de toutes les lectures). Même effet, plus ponctuel, pour chaque heure que le DAG manque + (`mock_api_import` en échec, Airflow arrêté). Aucun garde-fou ne détecte aujourd'hui un lag + calculé sur un écart réel différent de celui attendu ; issue de suivi à ouvrir. - **Le recouvrement est refusé à l'ingestion.** `uq_reading_source` autorise deux lignes au même `(site_id, timestamp)` dès que `source` diffère : rien dans le schéma n'empêche donc un import Mock API manuel avec une fenêtre passée (le script accepte `--start-time`/`--end-time` arbitraires) de dupliquer un point déjà couvert par le CSV. `import_mock_api_history()` appelle - `refuse_if_overlaps_historical_dataset()` avant toute écriture : si la fenêtre demandée recouvre - au moins une lecture `source='csv'`, l'import est refusé (`ValueError`) plutôt que d'écrire un - doublon inter-source silencieux. + `refuse_if_overlaps_historical_dataset()` avant toute écriture, y compris en `--dry-run` (le + contrôle est en lecture seule) et avant le moindre appel à l'API Mock : si la fenêtre demandée + recouvre au moins une lecture `source='csv'`, l'import est refusé (`ValueError`) plutôt que + d'écrire un doublon inter-source silencieux. Le contrôle ne porte que sur la fenêtre demandée, + pas sur les lectures reçues : `fetch_readings()` écarte donc toute lecture dont le `timestamp` + déborde de `[start_time, end_time)`, pour qu'une réponse hors fenêtre (bug du mock, ou hostile) + ne puisse pas le contourner. Ce contrôle compare des instants, pas des chaînes : `parse_datetime()` + pose `tzinfo=UTC` sur une entrée sans fuseau (même pattern que `_vers_utc()` dans + `app/services/reading.py`), sans quoi l'encodeur `timestamptz` d'asyncpg lirait un datetime naïf + dans le fuseau local du **processus**, correct dans le conteneur Airflow (UTC) mais décalé pour + un import manuel lancé depuis un poste en Europe/Paris. - **Le pipeline ML déduplique en défense.** Le garde-fou ci-dessus protège l'ingestion, pas la lecture : si un recouvrement se produisait malgré tout (import direct en base, contournement du script), `ml/enervision_ml/data.py` ne doit pas casser silencieusement l'hypothèse de `build_features` (« une ligne par `(site_id, timestamp)` »). `load_from_database()` et `load_recent_from_database()` utilisent donc `SELECT DISTINCT ON (site_id, timestamp)`, `source - = 'csv'` gagnant sur `'api_history'` en cas d'égalité — l'historique étant une source vérifiée, - l'API Mock une entrée hostile (cf. ci-dessus). + = 'csv'` gagnant sur `'api_history'` en cas d'égalité, l'historique étant une source vérifiée, + l'API Mock une entrée hostile (cf. ci-dessus). **Cette préférence est spécifique au chargeur + ML.** `GET /readings` renvoie les deux lignes sans les fusionner, et `DriftRepository` / + `ReadingRepository.latest_by_site()` / `.latest_for_site()` départagent par `reading_id` le plus + grand (en pratique la ligne insérée en dernier, pas forcément `csv`) : en cas de recouvrement, la + dérive comparerait alors une prévision à une valeur différente de celle sur laquelle le modèle a + appris. Pas d'incohérence aujourd'hui tant que le recouvrement reste refusé à l'ingestion ; à + aligner si ce garde-fou devait un jour être contourné. ### Qualité des données de l'API Mock diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index e649be5..f81bc60 100644 --- a/etl/airflow/dags/mock_api_import.py +++ b/etl/airflow/dags/mock_api_import.py @@ -18,16 +18,6 @@ from airflow.timetables.trigger import CronTriggerTimetable # Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" -# L'API Mock ne renvoie pas un flux au rythme naturel : elle répartit exactement `limit` -# lectures, espacées uniformément, sur toute la fenêtre demandée (vérifié empiriquement : -# une fenêtre d'1h avec `limit=1000` renvoie 1000 lectures espacées de 3,6s à l'intérieur de -# cette heure, pas une lecture horaire). Comme la fenêtre de ce DAG est toujours 1h -# (`interval=timedelta(hours=1)` ci-dessous), `limit=1` est ce qui produit une lecture par -# heure, alignée sur l'heure, cohérente avec le grain horaire du reste du schéma -# (`period_minutes=60`, historique CSV à 1 ligne/heure). Un `limit` plus grand ici fabriquerait -# des lectures infra-horaires, incompatibles avec les lags positionnels de `build_features`. -LIMITE_LECTURES = 1 - # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # inférieure au pas horaire du DAG. NOMBRE_REPRISES = 2 @@ -36,7 +26,11 @@ PLAFOND_PAR_TENTATIVE = timedelta(minutes=10) # L'intervalle est déclaré explicitement pour ne pas dépendre de la valeur du paramètre Airflow # `create_cron_data_intervals`. Le déclenchement à :45 laisse quinze minutes avant `ml_score`, -# exécuté à l'heure pile, puis avant `alertes`, exécuté à :15. +# exécuté à l'heure pile, puis avant `alertes`, exécuté à :15. Conséquence vérifiée empiriquement +# sur l'API Mock (cf. `limit_for_window()` dans `app.etl.mock_api_import`) : la fenêtre importée +# est `[:45, :45)`, donc chaque lecture atterrit à :45, pas à :00, un décalage constant sur +# toute la série, sans effet sur les lags positionnels ni sur les jointures en aval (`ml_score` +# vise la dernière lecture + 1h, `derive` joint à l'égalité), cf. 40-data.md. PLANIFICATION = CronTriggerTimetable( "45 * * * *", timezone="UTC", @@ -58,8 +52,11 @@ with DAG( bash_command=( f"{COMMANDE_BACKEND} app.etl.mock_api_import " "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\" " - "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" " - f"--limit {LIMITE_LECTURES}" + "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" + # Pas de --limit : app.etl.mock_api_import.limit_for_window() le dérive de la + # fenêtre (une lecture/heure), et refuse une fenêtre qui ne couvre pas un nombre + # entier d'heures ou dépasse le plafond de l'API. Porter la règle dans le code, + # pas dans ce DAG, évite qu'un appel manuel oublie de la respecter. ), retries=NOMBRE_REPRISES, retry_delay=DELAI_ENTRE_REPRISES, diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 4e4e077..7227486 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -118,7 +118,8 @@ def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande - assert "--limit 1" in commande + # Pas de --limit : app.etl.mock_api_import.limit_for_window() le dérive de la fenêtre. + assert "--limit" not in commande @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) diff --git a/ml/tests/test_data_integration.py b/ml/tests/test_data_integration.py index 9555d02..0273c19 100644 --- a/ml/tests/test_data_integration.py +++ b/ml/tests/test_data_integration.py @@ -45,11 +45,13 @@ def test_load_from_database_deduplicates_two_sources_at_the_same_instant( # `uq_reading_source` autorise deux lignes au meme (site_id, timestamp) des que `source` # differe : le garde-fou vit dans `mock_api_import.py`, pas dans le schema. Le chargeur ML # doit donc imposer lui-meme "une ligne par (site_id, timestamp)", pas la supposer. + # + # `csv` est inseree en premier (reading_id le plus bas) et `api_history` en second (le plus + # haut) : un depart par `reading_id DESC` seul choisirait `api_history` a tort. Seule la + # preference explicite pour `source='csv'` fait gagner le bon reading_id ici, et le test + # cesserait de proteger cette regle si l'ordre d'insertion etait inverse. site_id = insere_site(connexion_ml) dataset_id = insere_dataset(connexion_ml) - insere_lecture( - connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history" - ) insere_lecture( connexion_ml, site_id, @@ -58,6 +60,9 @@ def test_load_from_database_deduplicates_two_sources_at_the_same_instant( source="csv", dataset_id=dataset_id, ) + insere_lecture( + connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history" + ) frame = load_from_database(connexion_ml) @@ -69,11 +74,10 @@ def test_load_from_database_deduplicates_two_sources_at_the_same_instant( def test_load_recent_from_database_prefers_csv_when_two_sources_share_an_instant( connexion_ml: Connection, ) -> None: + # Meme ordre d'insertion que ci-dessus, et pour la meme raison : `csv` doit gagner malgre un + # `reading_id` plus bas que celui d'`api_history`. site_id = insere_site(connexion_ml) dataset_id = insere_dataset(connexion_ml) - insere_lecture( - connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history" - ) insere_lecture( connexion_ml, site_id, @@ -82,6 +86,9 @@ def test_load_recent_from_database_prefers_csv_when_two_sources_share_an_instant source="csv", dataset_id=dataset_id, ) + insere_lecture( + connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history" + ) frame = load_recent_from_database( connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3) From cbbfaf4910f52896498d980b772b3327091071ab Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Wed, 23 Sep 2026 15:31:51 +0200 Subject: [PATCH 3/6] fix(etl): importer une seule mesure par heure depuis l'API Mock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit L'API Mock ne renvoie pas les mesures d'une période : elle génère `limit` points répartis sur l'intervalle demandé (1 000 par heure avec --limit 1000, un toutes les 3,6 s). Le DAG aurait écrit 7 000 lignes par heure et par environnement, alors que le dataset historique a une mesure horaire et que les features ML décalent par ligne : `shift(168)`, le retard d'une semaine, serait devenu un retard de dix minutes, sans erreur visible au scoring ni au réentraînement. Avec --limit 1, l'API renvoie la mesure de :00 de chaque heure, au pas du CSV. Constaté sur la recette le 23/09 avant la réactivation des DAGs. --- etl/README.md | 6 ++++-- etl/airflow/dags/mock_api_import.py | 6 +++--- etl/airflow/tests/test_dags.py | 2 +- 3 files changed, 8 insertions(+), 6 deletions(-) diff --git a/etl/README.md b/etl/README.md index 94f8ba2..cabf674 100644 --- a/etl/README.md +++ b/etl/README.md @@ -669,8 +669,10 @@ des recommandations (`alertes`, issue #116), l'import historique (`historical_im issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15). Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`. Il appelle -`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite de 1 000 lectures -par site. Les deux pipelines normalisent leurs données vers les tables communes `site` et +`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite d'une lecture +par site. L'API Mock génère autant de points que la limite demandée, répartis sur l'intervalle : +un seul donne la mesure de :00, au pas horaire du dataset historique, que les features ML +supposent en décalant par ligne. Les deux pipelines normalisent leurs données vers les tables communes `site` et `reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale des deux sources reste à compléter dans l'issue #15. diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index f0043d0..ca50254 100644 --- a/etl/airflow/dags/mock_api_import.py +++ b/etl/airflow/dags/mock_api_import.py @@ -18,9 +18,9 @@ from airflow.timetables.trigger import CronTriggerTimetable # Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" -# Le pipeline backend et l'API acceptent au maximum 1 000 lectures par site. -# Cette marge évite de perdre silencieusement une lecture si une heure en contient plus de 60. -LIMITE_LECTURES = 1000 +# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure +# de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes). +LIMITE_LECTURES = 1 # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # inférieure au pas horaire du DAG. diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index d03dd97..eb45364 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -118,7 +118,7 @@ def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande - assert "--limit 1000" in commande + assert commande.endswith("--limit 1") @pytest.mark.parametrize("task_id", ["detection", "recommandations"]) From 14eed08ff5e4f9af380aaa19978a5b31d5dc460b Mon Sep 17 00:00:00 2001 From: Dorian Date: Wed, 23 Sep 2026 15:36:46 +0200 Subject: [PATCH 4/6] fix(mock_api): remove merge conflict markers --- apps/backend/app/etl/mock_api_import.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/apps/backend/app/etl/mock_api_import.py b/apps/backend/app/etl/mock_api_import.py index ca24d2b..71e2e32 100644 --- a/apps/backend/app/etl/mock_api_import.py +++ b/apps/backend/app/etl/mock_api_import.py @@ -184,9 +184,7 @@ async def upsert_sites( ) -def _timestamp_in_window( - reading: dict[str, Any], start_time: datetime, end_time: datetime -) -> bool: +def _timestamp_in_window(reading: dict[str, Any], start_time: datetime, end_time: datetime) -> bool: valeur = reading.get("timestamp") if not isinstance(valeur, str): return False From 6c09beeb3ca11b7463a19cb14502dac2446c948b Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Wed, 23 Sep 2026 15:39:04 +0200 Subject: [PATCH 5/6] =?UTF-8?q?fix(etl):=20demander=20la=20mesure=20de=20l?= =?UTF-8?q?'heure=20pile=20=C3=A0=20l'API=20Mock?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Correctif de cbbfaf4, dont le message affirmait à tort une mesure à :00. Avec --limit 1, l'API Mock renvoie le point de début d'intervalle : le DAG, déclenché à :45 sur [:45 - 1 h, :45], aurait écrit ses mesures à :45, au pas horaire mais hors de la grille du dataset historique. L'intervalle part désormais de l'heure pile du déclenchement : un run à 13:45 demande [13:00, 13:45] et importe la mesure de 13:00, avant ml_score à 14:00. Vérifié contre l'API Mock en recette et par le rendu du gabarit Jinja. --- docs/architecture/10-infra.md | 2 +- etl/README.md | 10 +++++----- etl/airflow/dags/mock_api_import.py | 8 ++++---- etl/airflow/tests/test_dags.py | 4 ++-- 4 files changed, 12 insertions(+), 12 deletions(-) diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index 8cea7bd..9b47cc7 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -89,7 +89,7 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). | `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` | | `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` | -| `mock_api_import` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'heure précédant son déclenchement depuis l'API Mock | +| `mock_api_import` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe depuis l'API Mock la mesure de l'heure pile précédant son déclenchement | | `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. diff --git a/etl/README.md b/etl/README.md index cabf674..00bf60b 100644 --- a/etl/README.md +++ b/etl/README.md @@ -669,15 +669,15 @@ des recommandations (`alertes`, issue #116), l'import historique (`historical_im issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15). Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`. Il appelle -`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite d'une lecture -par site. L'API Mock génère autant de points que la limite demandée, répartis sur l'intervalle : -un seul donne la mesure de :00, au pas horaire du dataset historique, que les features ML -supposent en décalant par ligne. Les deux pipelines normalisent leurs données vers les tables communes `site` et +`app.etl.mock_api_import` sur l'intervalle qui va de l'heure pile à son déclenchement, avec une +limite d'une lecture par site. L'API Mock génère autant de points que la limite demandée, +répartis sur l'intervalle et le premier à son début : un seul donne la mesure de :00, au pas +horaire du dataset historique, que les features ML supposent en décalant par ligne. Les deux pipelines normalisent leurs données vers les tables communes `site` et `reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale des deux sources reste à compléter dans l'issue #15. Le DAG `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution -traite l'intervalle Airflow précédent. Les deux pipelines normalisent leurs données vers les +importe la mesure de l'heure pile qui précède son déclenchement. Les deux pipelines normalisent leurs données vers les tables communes `site` et `reading`, tout en conservant leur source (`csv` ou `api_history`). 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 diff --git a/etl/airflow/dags/mock_api_import.py b/etl/airflow/dags/mock_api_import.py index ca50254..f62aa66 100644 --- a/etl/airflow/dags/mock_api_import.py +++ b/etl/airflow/dags/mock_api_import.py @@ -1,7 +1,7 @@ """DAG d'import périodique des données de l'API Mock EnerVision (issue #15). Orchestre le pipeline existant `app.etl.mock_api_import` sans dupliquer sa logique ETL. -Chaque exécution traite l'heure précédant son déclenchement. +Chaque exécution importe la mesure de l'heure pile qui précède son déclenchement. Le pipeline backend reste responsable de la validation, de la normalisation, du suivi de la qualité, de l'idempotence et du chargement dans PostgreSQL/TimescaleDB. @@ -18,8 +18,8 @@ from airflow.timetables.trigger import CronTriggerTimetable # Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m" -# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle. Un seul donne la mesure -# de :00, au pas horaire du CSV que les features ML supposent (`shift(168)` compte des lignes). +# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle, le premier à son début. +# Un seul, depuis l'heure pile, donne la mesure de :00 au pas du CSV que suppose `shift(168)`. LIMITE_LECTURES = 1 # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste @@ -51,7 +51,7 @@ with DAG( task_id="import_mock_api", bash_command=( f"{COMMANDE_BACKEND} app.etl.mock_api_import " - "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\" " + "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\" " "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" " f"--limit {LIMITE_LECTURES}" ), diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index eb45364..473bda3 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -113,10 +113,10 @@ def test_mock_api_import_calls_the_existing_backend_module(dagbag: DagBag) -> No assert "app.etl.mock_api_import" in commande -def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None: +def test_mock_api_import_asks_for_the_on_the_hour_reading(dagbag: DagBag) -> None: commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command - assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande + assert "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\"" in commande assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande assert commande.endswith("--limit 1") From 49175d8ff77c9e056c4153c02c693cd0a890810a Mon Sep 17 00:00:00 2001 From: Dorian Date: Wed, 23 Sep 2026 15:55:12 +0200 Subject: [PATCH 6/6] fix: conflict --- apps/backend/app/etl/mock_api_import.py | 42 +++++++++++++------ .../backend/tests/etl/test_mock_api_import.py | 31 ++++++++++++-- docs/architecture/10-infra.md | 6 ++- docs/architecture/40-data.md | 29 +++++++------ 4 files changed, 77 insertions(+), 31 deletions(-) diff --git a/apps/backend/app/etl/mock_api_import.py b/apps/backend/app/etl/mock_api_import.py index 71e2e32..dc1ddb6 100644 --- a/apps/backend/app/etl/mock_api_import.py +++ b/apps/backend/app/etl/mock_api_import.py @@ -367,28 +367,44 @@ def limit_for_window(start_time: datetime, end_time: datetime) -> int: """Nombre de lectures à demander pour que l'API Mock en rende une par heure, alignée. L'API ne renvoie pas un flux à un rythme naturel : elle répartit exactement `limit` lectures, - espacées uniformément, sur toute la fenêtre `[start_time, end_time)` demandée (vérifié - empiriquement). `limit` = nombre d'heures de la fenêtre est donc le seul réglage cohérent avec - le grain horaire du reste du schéma (`period_minutes=60`, historique CSV à une ligne/heure) ; - un `limit` plus grand fabriquerait des lectures infra-horaires, incompatibles avec les lags - positionnels de `build_features`. La fenêtre doit donc couvrir un nombre entier d'heures. + espacées uniformément, sur toute la fenêtre `[start_time, end_time)` demandée, la première + au tout début de la fenêtre (vérifié empiriquement). Deux façons d'obtenir une lecture + alignée sur l'heure : + + - une fenêtre d'exactement N heures (`start_time` sur l'heure) donne, avec `limit=N`, N + lectures espacées d'1h pile, la première à `start_time` : c'est le chemin du backfill + manuel (plusieurs jours d'historique en un seul appel). + - une fenêtre plus courte qu'une heure, ou qui n'est pas un multiple entier d'heure, ne peut + espacer plusieurs lectures d'1h pile (l'espacement de l'API vaut toujours + `durée / limit`) : seule `limit=1` reste alignée, la lecture unique atterrissant à + `start_time`. C'est le chemin du DAG horaire, dont la fenêtre part de l'heure pile qui + précède son déclenchement jusqu'à l'instant du déclenchement lui-même (`:45`), donc plus + courte qu'une heure. + + Dans les deux cas, `start_time` doit tomber pile sur l'heure : c'est elle qui ancre + l'alignement, jamais `end_time`. Un `limit` plus grand que celui rendu ici fabriquerait des + lectures infra-horaires, incompatibles avec les lags positionnels de `build_features`. """ + if start_time.minute or start_time.second or start_time.microsecond: + raise ValueError( + f"La fenêtre doit démarrer pile sur l'heure : {start_time.isoformat()} ne l'est pas." + ) + duree = end_time - start_time heures, reste = divmod(duree.total_seconds(), 3600) - if reste != 0: - raise ValueError( - "La fenêtre doit couvrir un nombre entier d'heures pour obtenir une lecture par " - f"heure alignée : [{start_time.isoformat()}, {end_time.isoformat()}) n'en couvre pas." - ) + # Fenêtre plus courte qu'une heure, ou pas un multiple entier : aucun `limit` supérieur à 1 + # n'espacerait ses lectures d'1h pile (l'espacement vaut toujours durée / limit). Seule la + # lecture unique, ancrée sur `start_time`, reste alignée. + limit = int(heures) if reste == 0 and heures >= 1 else 1 - if heures > MAX_LIMIT: + if limit > MAX_LIMIT: raise ValueError( - f"La fenêtre demandée couvre {int(heures)}h, au-delà du plafond de {MAX_LIMIT} " + f"La fenêtre demandée couvre {limit}h, au-delà du plafond de {MAX_LIMIT} " "lectures accepté par l'API Mock." ) - return int(heures) + return limit async def import_mock_api_history( diff --git a/apps/backend/tests/etl/test_mock_api_import.py b/apps/backend/tests/etl/test_mock_api_import.py index 132a718..5156077 100644 --- a/apps/backend/tests/etl/test_mock_api_import.py +++ b/apps/backend/tests/etl/test_mock_api_import.py @@ -441,11 +441,34 @@ def test_limit_for_window_returns_one_per_hour() -> None: assert limite == 21 * 24 -def test_limit_for_window_rejects_a_partial_hour() -> None: - with pytest.raises(ValueError, match="nombre entier d'heures"): +def test_limit_for_window_falls_back_to_one_reading_under_an_hour() -> None: + # Le DAG horaire (`:45`) demande desormais [heure pile precedente, instant du declenchement) : + # une fenetre plus courte qu'une heure, dont l'espacement `duree/limit` ne peut jamais valoir + # 1h pile pour plus d'une lecture. Seule `limit=1`, ancree sur `start_time`, reste alignee. + limite = mock_api_import.limit_for_window( + datetime.fromisoformat("2026-09-02T12:00:00+00:00"), + datetime.fromisoformat("2026-09-02T12:45:00+00:00"), + ) + + assert limite == 1 + + +def test_limit_for_window_falls_back_to_one_reading_for_a_non_whole_hour_span() -> None: + # Meme raisonnement pour une fenetre de plus d'une heure mais qui n'en est pas un multiple + # entier : aucun `limit > 1` ne donnerait un espacement d'1h pile. + limite = mock_api_import.limit_for_window( + datetime.fromisoformat("2026-09-02T12:00:00+00:00"), + datetime.fromisoformat("2026-09-02T13:30:00+00:00"), + ) + + assert limite == 1 + + +def test_limit_for_window_rejects_a_start_time_not_on_the_hour() -> None: + with pytest.raises(ValueError, match="pile sur l'heure"): mock_api_import.limit_for_window( - datetime.fromisoformat("2026-09-02T12:00:00+00:00"), - datetime.fromisoformat("2026-09-02T12:30:00+00:00"), + datetime.fromisoformat("2026-09-02T12:05:00+00:00"), + datetime.fromisoformat("2026-09-02T13:05:00+00:00"), ) diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index 9b47cc7..b35e56f 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -99,8 +99,10 @@ modifier. Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures, à la minute `:45`. Un `CronTriggerTimetable` explicite lui attribue un intervalle d'une heure, y compris lors d'un -déclenchement manuel. Il transmet cet intervalle au script backend et charge les mesures dans -les tables communes `site` et `reading`. Le décalage à `:45` laisse quinze minutes avant le +déclenchement manuel, mais la fenêtre transmise au script backend part de l'heure pile qui +précède le déclenchement (pas de l'intervalle Airflow tel quel), pour que la mesure importée +tombe à :00 et non à :45, voir [40-data.md](40-data.md). Le pipeline charge la mesure dans les +tables communes `site` et `reading`. Le décalage à `:45` laisse quinze minutes avant le scoring exécuté à l'heure pile, puis quinze minutes supplémentaires avant les alertes à `:15`. `max_active_runs=1` empêche deux exécutions du DAG de se chevaucher. diff --git a/docs/architecture/40-data.md b/docs/architecture/40-data.md index 88dc21d..91ae347 100644 --- a/docs/architecture/40-data.md +++ b/docs/architecture/40-data.md @@ -442,18 +442,23 @@ Les paramètres de ligne de commande disponibles pour l'import sont : **Piège sur `limit`, corrigé dans le code plutôt que documenté** : l'API ne renvoie pas un flux à un rythme naturel, elle répartit exactement `limit` lectures, espacées uniformément, sur toute la -fenêtre `[start_time, end_time)` demandée (vérifié empiriquement en interrogeant directement -l'API). Une fenêtre d'une heure avec `limit=1000`, le réglage d'origine, renvoyait donc 1000 -lectures espacées de 3,6 secondes à l'intérieur de cette heure, pas une lecture horaire, -incompatible avec les lags positionnels de `build_features`. Plutôt que documenter la règle -« `limit` = nombre d'heures de la fenêtre » et compter sur chaque appelant pour la respecter, -`limit_for_window()` la porte : `import_mock_api_history()` calcule `limit` depuis la fenêtre -reçue, refuse une fenêtre qui ne couvre pas un nombre entier d'heures, et refuse un intervalle de -plus de 1000 heures (le plafond `limit` de l'API). `--limit` n'existe donc plus côté CLI. Le DAG -`mock_api_import` interroge toujours une fenêtre d'1h (`interval=timedelta(hours=1)`, voir -[10-infra.md](10-infra.md)) : la fenêtre `[:45, :45)` place chaque lecture à :45, pas à :00 (la -première lecture atterrit au début de la fenêtre demandée), un décalage constant sans effet sur -les lags positionnels ni sur les jointures en aval. +fenêtre `[start_time, end_time)` demandée, la première au tout début de la fenêtre (vérifié +empiriquement en interrogeant directement l'API). Une fenêtre d'une heure avec `limit=1000`, le +réglage d'origine, renvoyait donc 1000 lectures espacées de 3,6 secondes à l'intérieur de cette +heure, pas une lecture horaire, incompatible avec les lags positionnels de `build_features`. +Plutôt que documenter la règle « `limit` = nombre d'heures de la fenêtre » et compter sur chaque +appelant pour la respecter, `limit_for_window()` la porte : `import_mock_api_history()` calcule +`limit` depuis la fenêtre reçue, refuse une fenêtre dont `start_time` ne tombe pas pile sur +l'heure (c'est elle qui ancre l'alignement), et refuse un intervalle de plus de 1000 heures (le +plafond `limit` de l'API). `--limit` n'existe donc plus côté CLI. Deux formes de fenêtre sont +gérées : un multiple entier d'heures (`limit` = ce nombre d'heures, une lecture par heure +espacée d'1h pile, chemin du backfill manuel) ou une fenêtre plus courte qu'une heure, ou qui +n'en est pas un multiple entier (`limit=1`, seule valeur qui reste alignée quand l'espacement +`durée / limit` ne peut valoir 1h pile). Le DAG `mock_api_import` est dans ce second cas : il +demande la fenêtre `[heure pile précédant le déclenchement, instant du déclenchement)`, plus +courte qu'une heure, plutôt que l'intervalle Airflow `[data_interval_start, data_interval_end)` +tel quel (`[:45, :45)`) qui aurait placé l'unique lecture à :45, hors de la grille horaire du +reste du schéma. ### Flux d'ingestion API Mock