From b78322bd6178ce57cf215b63e5ff6352741e7096 Mon Sep 17 00:00:00 2001 From: Dorian Date: Wed, 23 Sep 2026 15:19:39 +0200 Subject: [PATCH] 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)