From 147a469ddc0a39db74e2b6ce41b66606c3f0bea3 Mon Sep 17 00:00:00 2001 From: Antoine Jacquin Date: Sun, 27 Sep 2026 11:23:51 +0200 Subject: [PATCH] =?UTF-8?q?Rendre=20la=20carte=20autonome,=20t=C3=A9l?= =?UTF-8?q?=C3=A9charger=20la=20pyramide=20en=20fond,=20refermer=20les=20t?= =?UTF-8?q?rous=20de=20zoom=20et=20borner=20les=20workers=20GPU=20par=20la?= =?UTF-8?q?=20VRAM?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 5.5 --- AGENTS.md | 4 + docker-compose.worker.yml | 5 + lidar_pipeline/gpu.py | 70 +++++++++++++ lidar_pipeline/mapserve.py | 75 ++++++++++---- lidar_pipeline/pipeline.py | 46 +++++++-- lidar_pipeline/tests/test_gpu.py | 55 ++++++++++- lidar_pipeline/tests/test_mapserve.py | 137 ++++++++++++++++++++++++++ lidar_pipeline/tests/test_pipeline.py | 15 +++ lidar_pipeline/tests/test_tiles.py | 116 ++++++++++++++++++++++ lidar_pipeline/tiles.py | 100 +++++++++++++++++-- 10 files changed, 586 insertions(+), 37 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 20bb433..779cf76 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -37,7 +37,11 @@ - **Pyramide mise à jour pendant le rendu** : le pipeline régénère l'inventaire après chaque dalle **par défaut** (`incremental_index` sauf `--no-index`, quel que soit le lanceur), avec passe différée quand l'anti-rebond de 3 s ignore une dalle. Côté carte, `_bg_poller` surveille l'inventaire toutes les `BG_WATCH_S` = 10 s (fichier local ou charge utile amont) et lance un scan dès qu'il change ; les tuiles des dalles modifiées passent **en tête** de file (`_bg_enqueue(front=True)`), devant l'arriéré. - **Pyramide à stockage réduit** : `tiles.zoom_cached` — seuls les niveaux standard pairs ≤ `TILE_CACHE_MAX_Z` (16) vont sur disque (une tuile @2x de z vaut une 256 px de z+1) ; les autres sont rendus à la volée (`get_tile`, cache mémoire `_mem_tiles`), même en `LIDAR_TILE_CACHE_ONLY`. Interface en AVIF (`@2x.avif`) avec `EvenLevelTileLayer` (niveau de tuile arrondi au pair supérieur) ; maintenance de fond en AVIF sur les seuls niveaux stockés. ~24 Mo → 1,2 Mo sur 8 dalles. - **Mémoire bornée de la carte (Pi, `mem_limit: 1g`)** : `_open_source` (`tiles.py`) met en cache une **copie détachée** de chaque source (une image AVIF ouverte garde son décodeur, ~18 Mo de plus par quadrant 2500²) et compte 4 octets/pixel (PIL stocke le RGB sur 32 bits) ; `Dockerfile.maps` fixe `MALLOC_MMAP_THRESHOLD_=1048576` pour que glibc rende les grands tampons au système. Sans ces deux points, la navigation à fort zoom (niveaux rendus à la volée) montait à ~700 Mo de RSS et le conteneur était tué en boucle par l'OOM killer (502 côté Traefik). +- **Carte autonome (Pi sans worker)** : amont (`LIDAR_SOURCE_URL`/`LIDAR_MAPS_URL`) éteint ⇒ la carte reste servie depuis le disque. `_remote_payload` mémorise aussi l'échec (TTL 60 s, délai 10 s : sinon chaque requête repayait le délai réseau et `/api/map/meta` ne répondait plus), puis l'index local (`_build_index`) prend le relais ; coupe-circuits sur les sources (`_SOURCE_OFFLINE`, 60 s) et les tuiles (`_UPSTREAM`, réarmé à expiration) ; source rapatriée datée à sa version amont (`os.utime`) pour que l'index local garde les mêmes dates ; pas de requête amont pour une tuile sans dalle quand l'inventaire vient de l'amont. En cache seule, une tuile périmée reste servie (`no-store`, `X-Tile-Pending`) en attendant la nouvelle. +- **Pyramide téléchargée en tâche de fond** : avec `LIDAR_MAPS_URL`, la maintenance (`_bg_process_one`) TÉLÉCHARGE les tuiles des niveaux stockés depuis l'amont (stat `telechargees`), sans attendre l'affichage ; rendu local en repli si l'amont ne répond pas. +- **Première apparition des dalles (`_apply_seen`, `index_xyz/.sources_seen.json`)** : date effective d'une source = max(version, première entrée dans l'inventaire). Une dalle écrite avant mais inventoriée après une tuile (TTL de l'index amont, anti-rebond de l'inventaire) périme donc la tuile — sinon trou permanent à ce niveau, sur disque, en mémoire et dans le navigateur (stamp `?v=` inchangé). Registre persistant ; absent avec un cache existant (mise à jour) ⇒ tout est périmé une fois. - **Encodage AVIF rapide** : `AVIF_SPEED = 9` (`rendering.py`, `tiles.py`, `_SUBTILE_AVIF_SPEED` dans `index.py`) — dalle 5000 × 5000 px encodée en 0,6 s au lieu de 4 s (+3 % de taille, −0,3 dB) ; l'encodage était l'étape la plus longue du rendu d'une couche. +- **Workers bornés par la VRAM** : chaque processus du pool prend une place fixe à sa création (`_init_worker_slot` dans `pipeline.py`, file multiprocessing) calculée par `gpu_worker_slots` (`gpu.py`) : au plus (VRAM libre − `LIDAR_GPU_RESERVE_MIB` 512) / `LIDAR_GPU_WORKER_MIB` (2048, pic estimé : contexte CUDA + calage conjoint CuPy en float64 + EDT de `_fill_nans`) workers par GPU, au moins un, entrelacés ; l'excédent tourne en CPU (`force_cpu`). L'ancien round-robin par numéro de fichier mettait 6 workers par RTX 5060 de 8 Go avec `LIDAR_WORKERS=auto` (12) : OOM. VRAM illisible (nvidia-smi en échec) : round-robin sans borne. - **Préparation sur GPU** : `_fill_nans` (transformée de distance `cupyx`) et les gradients de `SharedDEM` passent sur GPU quand il est actif (repli CPU). `NUMBA_CACHE_DIR=/tmp/numba-cache` (Dockerfile) évite la recompilation des noyaux numba à chaque worker. - **Default output is AVIF**, not WebP. Use `--format webp` for WebP. Quality default is 60 (visually lossless on smooth color ramps, ~÷3 vs q98). - **Calage vertical des faisceaux de vol** : chaque tuile mélange plusieurs passes (1-2 `PointSourceId` par passe) parfois biaisées verticalement de quelques cm (±2,5 cm mesurés sur 1000_6882). `create_dtm_fast` mesure l'offset robuste de chaque faisceau (points sol, maille 1 m, surface médiane itérée 3×) et retranche les offsets ≥ 0,5 cm (`STRIP_ALIGN_THRESHOLD` dans `dtm.py`) avant rastérisation. Offsets calculés **par tuile** (ils dérivent le long d'une ligne de vol : jamais de table globale), mémoïsés par LAS sol, consignés dans `DTM/*_dtm*_stripalign.json` (sidecar de cache : absent, ou version/seuil/paramètres différents ⇒ régénération du DTM). Désactivable : `--no-strip-align`. diff --git a/docker-compose.worker.yml b/docker-compose.worker.yml index 0ded456..e87e043 100644 --- a/docker-compose.worker.yml +++ b/docker-compose.worker.yml @@ -36,6 +36,11 @@ services: # Les générations lancées depuis une carte distante utilisent le GPU - LIDAR_GPU=1 - LIDAR_WORKERS=auto + # Workers par GPU bornés par la VRAM libre au lancement du run + # ((libre − réserve) / pic d'un worker) ; l'excédent tourne en CPU. + # Pic estimé à 2048 Mo : à ajuster après mesure (nvidia-smi pendant un run). + # - LIDAR_GPU_WORKER_MIB=2048 + # - LIDAR_GPU_RESERVE_MIB=512 # Une ligne par worker : BLAS/OpenMP mono-thread, sinon 12 workers × N # threads écrasent les 14 cœurs (load 68+ observé pendant les runs) - OMP_NUM_THREADS=1 diff --git a/lidar_pipeline/gpu.py b/lidar_pipeline/gpu.py index f3dc6b4..7be410b 100644 --- a/lidar_pipeline/gpu.py +++ b/lidar_pipeline/gpu.py @@ -294,6 +294,76 @@ def set_active_gpu(gpu_id): _init_gpu() +def force_cpu(): + """Interdit le GPU à ce processus (worker en excédent de VRAM). + + Aucun candidat ne survit au filtre et CUDA_VISIBLE_DEVICES vide empêche + CuPy de créer un contexte (~300 Mo de VRAM par processus sinon). + """ + global _restricted_gpu_ids + _restricted_gpu_ids = [] + os.environ['CUDA_VISIBLE_DEVICES'] = '' + _init_gpu() + + +# Pic de VRAM d'un worker (Mo) : contexte CUDA (~300) + calage conjoint des +# lignes sur CuPy (~70-100 o par point sol en float64, ~15 M points par dalle +# 0,2 m + bande de 100 m) + transformée de distance de _fill_nans (indices +# int32 ×2 sur ~7000² px ≈ 400 Mo) ; le pool CuPy garde ses blocs entre deux +# étapes. Estimation depuis le code, à ajuster par LIDAR_GPU_WORKER_MIB. +GPU_WORKER_MIB = int(os.environ.get("LIDAR_GPU_WORKER_MIB", "2048") or 2048) +# VRAM laissée libre sur chaque GPU (affichage, autres processus). +GPU_RESERVE_MIB = int(os.environ.get("LIDAR_GPU_RESERVE_MIB", "512") or 512) + + +def gpu_free_mib(): + """VRAM libre par GPU (indice hôte → Mo) via nvidia-smi, sans contexte CUDA. + + Dictionnaire vide si la mesure échoue : l'appelant ne borne alors rien. + """ + try: + import subprocess + result = subprocess.run( + ['nvidia-smi', '--query-gpu=index,memory.free', + '--format=csv,noheader,nounits'], + capture_output=True, text=True, timeout=15, + ) + if result.returncode != 0: + return {} + free = {} + for line in result.stdout.strip().splitlines(): + parts = [p.strip() for p in line.split(',')] + if len(parts) >= 2: + free[int(parts[0])] = int(parts[1]) + return free + except Exception: + return {} + + +def gpu_worker_slots(gpu_ids, n_workers, free_mib=None): + """Place de chaque worker du pool : indice GPU, -1 (CPU forcé) ou None. + + Chaque GPU reçoit au plus (VRAM libre − réserve) / GPU_WORKER_MIB + workers (au moins un), entrelacés entre GPU ; l'excédent tourne en CPU + au lieu de saturer la VRAM (OOM). Sans GPU : None partout (choix laissé + au worker). VRAM inconnue : round-robin sans borne (comportement + historique). + """ + if not gpu_ids: + return [None] * n_workers + if free_mib is None: + free_mib = gpu_free_mib() + if not all(g in free_mib for g in gpu_ids): + return [gpu_ids[i % len(gpu_ids)] for i in range(n_workers)] + capacity = {g: max(1, (free_mib[g] - GPU_RESERVE_MIB) // GPU_WORKER_MIB) + for g in gpu_ids} + slots = [] + for rank in range(max(capacity.values())): + slots += [g for g in gpu_ids if capacity[g] > rank] + slots = slots[:n_workers] + return slots + [-1] * (n_workers - len(slots)) + + def _gpu_available(): """Check if GPU is usable right now.""" try: diff --git a/lidar_pipeline/mapserve.py b/lidar_pipeline/mapserve.py index 95ea596..2311846 100644 --- a/lidar_pipeline/mapserve.py +++ b/lidar_pipeline/mapserve.py @@ -278,27 +278,46 @@ def background_scan(output_dir=None): def _bg_process_one(): - """Traite une source à rapatrier, sinon une tuile de la file. Vrai si - un travail a été fait (rapatriement ou rendu).""" - try: - src = _bg["prefetch"].popleft() - except IndexError: - src = None - if src is not None: + """Traite une tuile de la file, sinon une source à rapatrier. Vrai si + un travail a été fait (téléchargement, rendu ou rapatriement). + + Tuiles d'abord : légères et visibles tout de suite, elles ne doivent pas + attendre des milliers de sources (utiles aux seuls forts zooms rendus à + la volée). Les tuiles déjà à jour sont passées sans pause. + """ + conf = _bg_conf() + while True: + item = _bg_pop() + if item is None: + break + layer, z, x, y = item + _data, state = tiles_mod.cached_tile(OUTPUT_DIR, layer, z, x, y, + conf["scale"], conf["fmt"]) + if state == "pending": + break + _bg["stats"]["a_jour" if state == "fresh" else "vides"] += 1 + if item is None: + try: + src = _bg["prefetch"].popleft() + except IndexError: + return False if src.ensure(): _bg["stats"]["sources_locales"] = _bg["stats"].get("sources_locales", 0) + 1 return True return False # amont éteint : repris au prochain scan - item = _bg_pop() - if item is None: - return False - layer, z, x, y = item - conf = _bg_conf() - _data, state = tiles_mod.cached_tile(OUTPUT_DIR, layer, z, x, y, - conf["scale"], conf["fmt"]) - if state != "pending": - _bg["stats"]["a_jour" if state == "fresh" else "vides"] += 1 - return False + if MAPS_URL: + # Serveur de tuiles amont : la pyramide est TÉLÉCHARGÉE en tâche de + # fond (quelques dizaines de Ko, aucun décodage de dalle) ; la carte + # se remplit sans visite et reste servie ensuite sans l'amont. + suffix = "@2x" if conf["scale"] == 2 else "" + data = _fetch_upstream(f"tiles/{layer}/{z}/{x}/{y}{suffix}.{conf['fmt']}") + if data: + tiles_mod._write_atomic( + tiles_mod.tile_cache_path(OUTPUT_DIR, layer, z, x, y, + conf["scale"], conf["fmt"]), data) + _bg["stats"]["telechargees"] = _bg["stats"].get("telechargees", 0) + 1 + return True + # Amont absent ou éteint : rendu local depuis les sources gardées (autonomie). tiles_mod.get_tile(OUTPUT_DIR, layer, z, x, y, conf["scale"], conf["fmt"]) _bg["stats"]["rendues"] += 1 return True @@ -539,7 +558,8 @@ def _upstream_mark(ok): _UPSTREAM.update(fails=0, until=0.0) return _UPSTREAM["fails"] += 1 - if _UPSTREAM["fails"] >= _UPSTREAM_OFFLINE_AFTER and _UPSTREAM["until"] == 0.0: + # Réarmé après expiration : un amont resté éteint est resuspendu. + if _UPSTREAM["fails"] >= _UPSTREAM_OFFLINE_AFTER and not _upstream_offline(): _UPSTREAM["until"] = time.time() + _UPSTREAM_OFFLINE_SECONDS logger.info("Serveur de tuiles amont injoignable — tentatives suspendues 2 min") @@ -602,7 +622,11 @@ def _tile_bytes(layer, z, x, y, scale, fmt): except Exception as e: # noqa: BLE001 — une tuile ne casse pas la carte logger.warning(f"Rendu de tuile impossible ({layer} {z}/{x}/{y}) : {e}") data = None - if data is None and MAPS_URL: + # Inventaire venu de l'amont sans aucune dalle sous la tuile : + # l'amont n'en sait pas plus, inutile de l'interroger. + if data is None and MAPS_URL and not ( + tiles_mod.REMOTE_SOURCE_URL + and not tiles_mod._contributing(OUTPUT_DIR, layer, z, x, y, scale)): suffix = "@2x" if scale == 2 else "" data = _fetch_upstream(f"tiles/{layer}/{z}/{x}/{y}{suffix}.{fmt}") if data: @@ -664,11 +688,20 @@ async def tile(layer: str, z: int, x: int, name: str, v: Optional[str] = None): empty = False finally: _release_lock(key, entry) + if pending: + # Tuile périmée (dalle régénérée) : l'ancienne reste affichée en + # attendant la nouvelle, plutôt qu'un trou transparent. + try: + data = tiles_mod.tile_cache_path(OUTPUT_DIR, layer, z, x, y, + scale, fmt).read_bytes() + except OSError: + data = None else: data = await run_in_threadpool(_tile_bytes, layer, z, x, y, scale, fmt) pending = False empty = data is None - if empty or pending: + stale = pending and bool(data) + if empty or (pending and not stale): data = _transparent_tile(scale, fmt) headers = dict(_cors_headers()) if pending: @@ -679,7 +712,7 @@ async def tile(layer: str, z: int, x: int, name: str, v: Optional[str] = None): else: headers["Cache-Control"] = ("public, max-age=31536000, immutable" if v else "public, max-age=300, must-revalidate") - headers["X-Tile-Empty"] = "1" if (empty or pending) else "0" + headers["X-Tile-Empty"] = "1" if (empty or (pending and not stale)) else "0" return Response(content=data, media_type=_MEDIA[fmt], headers=headers) diff --git a/lidar_pipeline/pipeline.py b/lidar_pipeline/pipeline.py index eba7a27..eaf0326 100644 --- a/lidar_pipeline/pipeline.py +++ b/lidar_pipeline/pipeline.py @@ -91,7 +91,7 @@ from .visualizations import ( generate_anomaly_mask, generate_relief_oriente, ) -from .gpu import gpu_cleanup, num_gpus, available_gpu_ids, restrict_gpus, safe_gpu_call +from .gpu import gpu_cleanup, num_gpus, available_gpu_ids, restrict_gpus, safe_gpu_call, gpu_worker_slots from .ign import generate_ign_overlay from .rendering import tif_to_crop @@ -750,13 +750,27 @@ class LidarArchaeoPipeline: logger.info(f"Traitement parallèle avec {self.workers} workers...") logger.info(f"Fichiers: {len(files)}") - with ProcessPoolExecutor(max_workers=self.workers) as executor: - # Round-robin assign each file to a real GPU host index - active_ids = self.gpu_ids if self.gpu_ids else available_gpu_ids() + # Une place fixe par processus du pool (GPU ou CPU), prise à sa + # création : au plus « VRAM libre / pic d'un worker » processus par + # GPU. L'affectation par numéro de fichier (round-robin) laissait + # 6 workers sur chaque GPU de 8 Go avec LIDAR_WORKERS=auto : OOM. + active_ids = self.gpu_ids if self.gpu_ids else available_gpu_ids() + slots = gpu_worker_slots(active_ids, self.workers) + if active_ids: + per_gpu = ", ".join(f"GPU {g} : {slots.count(g)}" for g in active_ids) + n_cpu = slots.count(-1) + logger.info(f"Répartition des workers : {per_gpu}" + + (f", CPU : {n_cpu} (VRAM insuffisante)" if n_cpu else "")) + slot_queue = multiprocessing.Queue() + for slot in slots: + slot_queue.put(slot) + with ProcessPoolExecutor(max_workers=self.workers, + initializer=_init_worker_slot, + initargs=(slot_queue,)) as executor: resolutions_str = ','.join(str(r) for r in self.resolutions) future_to_file = { - executor.submit(_process_file_standalone, str(laz_file), str(self.input_dir), str(self.output_dir), resolutions_str, self.force, self.ground_method, self.ign_classes, self.force_classify, self.keep_tif, self.quality, self.only_viz, self.skip_viz, self.output_format, active_ids[file_idx % len(active_ids)] if active_ids else None, self.openness_downsample, self.edge_buffer): laz_file - for file_idx, laz_file in enumerate(files) + executor.submit(_process_file_standalone, str(laz_file), str(self.input_dir), str(self.output_dir), resolutions_str, self.force, self.ground_method, self.ign_classes, self.force_classify, self.keep_tif, self.quality, self.only_viz, self.skip_viz, self.output_format, None, self.openness_downsample, self.edge_buffer): laz_file + for laz_file in files } done = 0 t_deadline = time.time() + 7200 @@ -857,6 +871,26 @@ class LidarArchaeoPipeline: logger.warning(f" Note: Impossible de supprimer les fichiers temporaires: {e}") +def _init_worker_slot(slot_queue): + """Initialiseur du pool : le processus prend sa place une fois pour toutes. + + Indice GPU → set_active_gpu ; -1 → CPU forcé (VRAM insuffisante) ; + None ou file vide → choix laissé au worker (pas de GPU détecté). + """ + import queue + from . import gpu + try: + slot = slot_queue.get_nowait() + except queue.Empty: + return + if slot is None: + return + if slot < 0: + gpu.force_cpu() + else: + gpu.set_active_gpu(slot) + + def _process_file_standalone(laz_file_str, input_dir, output_dir, resolution, force=False, ground_method='auto', ign_classes="sol", force_classify=False, keep_tif=False, quality=60, only_viz=None, skip_viz=None, output_format='avif', gpu_id=None, openness_downsample=None, edge_buffer=0.0): """Standalone function for multiprocessing — creates its own pipeline instance. diff --git a/lidar_pipeline/tests/test_gpu.py b/lidar_pipeline/tests/test_gpu.py index 0d65ed5..d2d3053 100644 --- a/lidar_pipeline/tests/test_gpu.py +++ b/lidar_pipeline/tests/test_gpu.py @@ -199,4 +199,57 @@ def test_safe_gpu_call_retries_on_cpu_for_any_error(monkeypatch): gpu.safe_gpu_call(g, 1) raise AssertionError("devait relancer l'erreur") except ValueError: - pass \ No newline at end of file + pass + +def test_gpu_worker_slots_bounded_by_free_vram(monkeypatch): + """Places GPU par la VRAM libre : l'excédent de workers passe en CPU. + + Cas réel : LIDAR_WORKERS=auto = 12 workers sur 2 RTX 5060 (8 Go) — + 6 workers par GPU, soit bien plus que la VRAM libre ne tient au pic + (calage des lignes + comblement GPU) : OOM. + """ + from lidar_pipeline import gpu + monkeypatch.setattr(gpu, "GPU_WORKER_MIB", 2000) + monkeypatch.setattr(gpu, "GPU_RESERVE_MIB", 500) + free = {0: 7500, 1: 4600} + slots = gpu.gpu_worker_slots([0, 1], 12, free_mib=free) + assert len(slots) == 12 + assert slots.count(0) == 3 # (7500 - 500) // 2000 + assert slots.count(1) == 2 # (4600 - 500) // 2000 + assert slots.count(-1) == 7 # reste en CPU + # GPU d'abord, entrelacés : les premiers workers créés se répartissent + assert slots[:4] == [0, 1, 0, 1] + # Moins de workers que de places : aucun CPU forcé + assert gpu.gpu_worker_slots([0, 1], 3, free_mib=free) == [0, 1, 0] + + +def test_gpu_worker_slots_without_gpu_or_vram_info(): + """Sans GPU : tout en CPU implicite (None). VRAM inconnue : round-robin + historique (pas de bornage sans mesure).""" + from lidar_pipeline import gpu + assert gpu.gpu_worker_slots([], 4, free_mib={}) == [None] * 4 + assert gpu.gpu_worker_slots([0, 1], 4, free_mib={}) == [0, 1, 0, 1] + + +def test_gpu_worker_slots_keeps_one_gpu_worker_when_tight(monkeypatch): + """GPU presque plein : au moins une place pour que le GPU serve encore + (le repli CPU de safe_gpu_call couvre l'OOM éventuel).""" + from lidar_pipeline import gpu + monkeypatch.setattr(gpu, "GPU_WORKER_MIB", 2000) + monkeypatch.setattr(gpu, "GPU_RESERVE_MIB", 500) + assert gpu.gpu_worker_slots([0], 3, free_mib={0: 900}) == [0, -1, -1] + + +def test_force_cpu_disables_gpu_selection(monkeypatch): + """force_cpu() : aucun candidat GPU, CuPy jamais initialisé.""" + import os + from lidar_pipeline import gpu + monkeypatch.setattr(gpu, "_restricted_gpu_ids", None) + monkeypatch.setattr(gpu, "_gpu_candidates", [(0, "X", "12.0", 8000, 1, 12)]) + monkeypatch.setattr(gpu, "_gpu_initialized", False) + monkeypatch.setattr(gpu, "_env_set_by_init", False) + monkeypatch.delenv("CUDA_VISIBLE_DEVICES", raising=False) + gpu.force_cpu() + assert gpu.available_gpu_ids() == [] + assert gpu.HAS_GPU is False + assert os.environ.get("CUDA_VISIBLE_DEVICES") == "" diff --git a/lidar_pipeline/tests/test_mapserve.py b/lidar_pipeline/tests/test_mapserve.py index 2589501..9048f81 100644 --- a/lidar_pipeline/tests/test_mapserve.py +++ b/lidar_pipeline/tests/test_mapserve.py @@ -1032,3 +1032,140 @@ def test_remote_fed_map_keeps_sources_locally(tmp_path, monkeypatch): while mapserve._bg["prefetch"]: mapserve._bg_process_one() assert all(s_.path.exists() for s_ in [thumb] + quads) and not full.path.exists() + + +def _bg_reset(monkeypatch, mapserve): + import collections + for key, val in (("queue", collections.deque()), ("queued", set()), + ("snapshot", None), ("prefetch", collections.deque()), + ("stats", {"rendues": 0, "a_jour": 0, "vides": 0, "scans": 0, + "dernier_scan": None, "ajouts_dernier_scan": 0})): + monkeypatch.setitem(mapserve._bg, key, val) + + +def test_background_downloads_tiles_from_upstream(tmp_path, monkeypatch): + """Maintenance de fond + LIDAR_MAPS_URL : les tuiles de la pyramide sont + TÉLÉCHARGÉES en tâche de fond (pas seulement à l'affichage), sans rendu + local — la carte se remplit même sans visite et reste servie ensuite + sans l'amont.""" + import lidar_pipeline.mapserve as mapserve + from lidar_pipeline import tiles + _setup(tmp_path, monkeypatch) + _bg_reset(monkeypatch, mapserve) + monkeypatch.setattr(mapserve, "MAPS_URL", "http://amont:8975") + AVIF = b"\x00\x00\x00\x20ftypavif-TUILE-AMONT" + calls = [] + monkeypatch.setattr(mapserve, "_fetch_upstream", + lambda rel: calls.append(rel) or AVIF) + + def no_render(*a, **k): + raise AssertionError("rendu local inutile : la tuile vient de l'amont") + + monkeypatch.setattr(tiles, "get_tile", no_render) + z = 13 + x, y = _tile_of_cell(1054, 6882, z) + mapserve._bg_enqueue([("aspect", z, x, y)]) + assert mapserve._bg_process_one() is True + assert calls == [f"tiles/aspect/{z}/{x}/{y}@2x.avif"] + cache = tiles.tile_cache_path(tmp_path, "aspect", z, x, y, 2, "avif") + assert cache.read_bytes() == AVIF + assert mapserve._bg["stats"]["telechargees"] == 1 + # Fraîche : servie en cache seule sans amont ni rendu + data, state = tiles.cached_tile(tmp_path, "aspect", z, x, y, 2, "avif") + assert state == "fresh" and data == AVIF + + +def test_background_renders_locally_when_upstream_down(tmp_path, monkeypatch): + """Amont éteint : la maintenance rend la tuile localement (autonomie).""" + import lidar_pipeline.mapserve as mapserve + from lidar_pipeline import tiles + _setup(tmp_path, monkeypatch) + _bg_reset(monkeypatch, mapserve) + monkeypatch.setattr(mapserve, "MAPS_URL", "http://amont:8975") + monkeypatch.setattr(mapserve, "_fetch_upstream", lambda rel: None) + z = 13 + x, y = _tile_of_cell(1054, 6882, z) + mapserve._bg_enqueue([("aspect", z, x, y)]) + assert mapserve._bg_process_one() is True + assert tiles.tile_cache_path(tmp_path, "aspect", z, x, y, 2, "avif").is_file() + assert mapserve._bg["stats"]["rendues"] == 1 + + +def test_upstream_breaker_rearms_after_expiry(monkeypatch): + """La suspension expirée, de nouveaux échecs la relancent (sinon chaque + tuile repaie le délai réseau d'un amont éteint, indéfiniment).""" + import lidar_pipeline.mapserve as mapserve + monkeypatch.setattr(mapserve, "_UPSTREAM", {"fails": 0, "until": 0.0}) + for _ in range(3): + mapserve._upstream_mark(False) + assert mapserve._upstream_offline() is True + mapserve._UPSTREAM["until"] = 1.0 # suspension expirée + assert mapserve._upstream_offline() is False + mapserve._upstream_mark(False) # nouvel échec : resuspendu + assert mapserve._upstream_offline() is True + + +def test_empty_tile_skips_upstream_when_inventory_is_upstream(tmp_path, monkeypatch): + """Inventaire venu de l'amont (LIDAR_SOURCE_URL) et aucune dalle sous la + tuile : l'amont n'en sait pas plus — pas de requête (20 s par tuile vide + quand il est éteint).""" + import lidar_pipeline.mapserve as mapserve + from lidar_pipeline import tiles + _setup(tmp_path, monkeypatch) + monkeypatch.setattr(mapserve, "MAPS_URL", "http://amont:8975") + monkeypatch.setattr(tiles, "REMOTE_SOURCE_URL", "http://amont:8973") + monkeypatch.setattr(tiles, "_contributing", lambda *a, **k: []) + calls = [] + monkeypatch.setattr(mapserve, "_fetch_upstream", lambda rel: calls.append(rel)) + resp = asyncio.run(mapserve.tile("aspect", 15, 1, "1.png")) + assert calls == [] + assert resp.headers["X-Tile-Empty"] == "1" + + +def test_cache_only_serves_stale_tile_while_pending(tmp_path, monkeypatch): + """Cache seule, tuile périmée (dalle régénérée) et amont muet : l'ancienne + tuile reste affichée en attendant la nouvelle — pas un trou transparent.""" + import os + import lidar_pipeline.mapserve as mapserve + from lidar_pipeline import tiles + _setup(tmp_path, monkeypatch) + monkeypatch.setenv("LIDAR_TILE_CACHE_ONLY", "1") + monkeypatch.setattr(mapserve, "MAPS_URL", "") + z = 13 + x, y = _tile_of_cell(1054, 6882, z) + cache = tiles.tile_cache_path(tmp_path, "aspect", z, x, y, 2, "avif") + cache.parent.mkdir(parents=True, exist_ok=True) + cache.write_bytes(b"ANCIENNE-TUILE") + os.utime(cache, (1_000_000, 1_000_000)) # plus vieille que la dalle + assert tiles.cached_tile(tmp_path, "aspect", z, x, y, 2, "avif")[1] == "pending" + resp = asyncio.run(mapserve.tile("aspect", z, x, f"{y}@2x.avif")) + assert resp.body == b"ANCIENNE-TUILE" + assert resp.headers["Cache-Control"] == "no-store" # redemandée plus tard + assert resp.headers["X-Tile-Pending"] == "1" + assert resp.headers["X-Tile-Empty"] == "0" + + +def test_background_tiles_before_sources(tmp_path, monkeypatch): + """Tuiles d'abord (légères, visibles tout de suite), sources ensuite : + des milliers de sources à rapatrier ne retardent pas la pyramide.""" + import lidar_pipeline.mapserve as mapserve + from lidar_pipeline import tiles + _setup(tmp_path, monkeypatch) + _bg_reset(monkeypatch, mapserve) + monkeypatch.setattr(mapserve, "MAPS_URL", "http://amont:8975") + monkeypatch.setattr(mapserve, "_fetch_upstream", lambda rel: b"TUILE") + fetched = [] + + class Src: + def ensure(self): + fetched.append(1) + return True + + mapserve._bg["prefetch"].append(Src()) + z = 13 + x, y = _tile_of_cell(1054, 6882, z) + mapserve._bg_enqueue([("aspect", z, x, y)]) + assert mapserve._bg_process_one() is True + assert fetched == [] and mapserve._bg["stats"]["telechargees"] == 1 + assert mapserve._bg_process_one() is True + assert fetched == [1] diff --git a/lidar_pipeline/tests/test_pipeline.py b/lidar_pipeline/tests/test_pipeline.py index 123b577..d0ca77b 100644 --- a/lidar_pipeline/tests/test_pipeline.py +++ b/lidar_pipeline/tests/test_pipeline.py @@ -308,3 +308,18 @@ class TestResolveWorkers: from lidar_pipeline.pipeline import resolve_workers assert resolve_workers('abc') == 1 assert resolve_workers(None) == 1 + + def test_worker_slot_initializer(self, monkeypatch): + """Chaque processus du pool prend UNE place à sa création : GPU + (set_active_gpu) ou CPU forcé (-1) ; None ou file vide = libre.""" + import queue + from lidar_pipeline import gpu, pipeline + calls = [] + monkeypatch.setattr(gpu, "set_active_gpu", lambda i: calls.append(("gpu", i))) + monkeypatch.setattr(gpu, "force_cpu", lambda: calls.append(("cpu",))) + q = queue.Queue() + for slot in (1, -1, None): + q.put(slot) + for _ in range(4): # 4e appel : file vide + pipeline._init_worker_slot(q) + assert calls == [("gpu", 1), ("cpu",)] diff --git a/lidar_pipeline/tests/test_tiles.py b/lidar_pipeline/tests/test_tiles.py index 56f9d29..afaf62f 100644 --- a/lidar_pipeline/tests/test_tiles.py +++ b/lidar_pipeline/tests/test_tiles.py @@ -571,3 +571,119 @@ def test_unstored_level_rendered_on_the_fly_without_disk(tmp_path, monkeypatch): assert data and data[:4] == b"\x89PNG" assert not tiles.tile_cache_path(tmp_path, "aspect", z, x, y, 1, "png").exists() assert tiles.get_tile(tmp_path, "aspect", z, x, y, 1, "png") == data and len(calls) == 1 + + +def test_remote_payload_failure_not_retried_each_call(monkeypatch): + """Amont injoignable sans inventaire connu : l'échec est mémorisé (TTL). + + Sinon chaque appel (plusieurs par /api/map/meta) repaie le délai de + connexion : l'interface ne se chargeait plus quand le worker était éteint. + """ + import urllib.request + from lidar_pipeline import tiles + monkeypatch.setattr(tiles, "REMOTE_SOURCE_URL", "http://amont:8973") + monkeypatch.setattr(tiles, "_remote_cache", + {"payload": None, "at": 0.0, "index": None, "root": None}) + calls = [] + + def boom(*a, **k): + calls.append(1) + raise OSError("hôte injoignable") + + monkeypatch.setattr(urllib.request, "urlopen", boom) + for _ in range(3): + assert tiles._remote_payload() is None + assert len(calls) == 1 + + +def test_fetch_source_offline_breaker(tmp_path, monkeypatch): + """Source amont en échec : les suivantes échouent sans attendre le délai + réseau pendant la suspension (la maintenance passe au rendu local).""" + import urllib.request + from lidar_pipeline import tiles + monkeypatch.setattr(tiles, "_SOURCE_OFFLINE", {"until": 0.0}) + calls = [] + + def boom(*a, **k): + calls.append(1) + raise OSError("hôte injoignable") + + monkeypatch.setattr(urllib.request, "urlopen", boom) + assert tiles._fetch_source("http://amont/a", tmp_path / "a.avif") is False + assert tiles._fetch_source("http://amont/b", tmp_path / "b.avif") is False + assert len(calls) == 1 + + +def test_fetched_source_dated_to_upstream_version(tmp_path, monkeypatch): + """Une source rapatriée porte la date de version amont : quand l'amont + s'éteint, l'index local retrouve les mêmes dates et les tuiles déjà + faites restent fraîches (pas de pyramide entière à refaire).""" + from lidar_pipeline import tiles + + def fake_fetch(url, dest): + dest.write_bytes(b"x") + return True + + monkeypatch.setattr(tiles, "_fetch_source", fake_fetch) + src = tiles._Source(tmp_path / "q.avif", (0, 0, 1, 1), 0.2, + url="http://amont/q.avif", version=1700000000000) + assert src.ensure() is True + assert (tmp_path / "q.avif").stat().st_mtime == 1700000000.0 + + +def test_tile_refreshed_when_older_dalle_appears(tmp_path): + """Dalle entrée dans l'inventaire APRÈS le rendu d'une tuile qui la couvre, + mais avec une date de version plus ancienne (écrite avant, inventoriée + après) : la tuile doit être périmée — sinon trou permanent à ce niveau, + sur le disque, en mémoire et dans le navigateur (stamp inchangé).""" + import io + import os + from PIL import Image + from lidar_pipeline import tiles + tiles.clear_source_cache() + tiles._mem_tiles.clear() + z = 10 + assert _tile_of_cell(1054, 6882, z) == _tile_of_cell(1055, 6882, z) + x, y = _tile_of_cell(1054, 6882, z) + _make_dalle(tmp_path, 1054, 6882, ["slope"], color=(200, 30, 30)) + tiles.source_index(tmp_path, force=True) + first = tiles.get_tile(tmp_path, "slope", z, x, y) + assert first is not None + stamp = tiles.tiles_stamp(tmp_path) + + _make_dalle(tmp_path, 1055, 6882, ["slope"], color=(30, 200, 30)) + for f in tmp_path.rglob("*1055_6882*"): + os.utime(f, (1_000_000, 1_000_000)) # version antérieure à la tuile + tiles.source_index(tmp_path, force=True) + assert tiles.cached_tile(tmp_path, "slope", z, x, y)[1] == "pending" + second = tiles.get_tile(tmp_path, "slope", z, x, y) + img = Image.open(io.BytesIO(second)).convert("RGBA") + assert any(g > 150 and r < 80 and a == 255 + for r, g, b, a in img.getdata()), "nouvelle dalle absente de la tuile" + assert tiles.tiles_stamp(tmp_path) > stamp + # Registre persistant : un redémarrage ne réinvalide rien + tiles._index_cache.clear() + tiles.source_index(tmp_path, force=True) + assert tiles.cached_tile(tmp_path, "slope", z, x, y)[1] == "fresh" + tiles.clear_source_cache() + + +def test_existing_cache_without_registry_is_refreshed_once(tmp_path): + """Mise à jour : cache de tuiles présent mais pas de registre — les + tuiles existantes (peut-être trouées) sont périmées une seule fois.""" + from lidar_pipeline import tiles + tiles._seen_cache.clear() + tiles._mem_tiles.clear() + z = 10 + x, y = _tile_of_cell(1054, 6882, z) + _make_dalle(tmp_path, 1054, 6882, ["slope"]) + tiles.source_index(tmp_path, force=True) + assert tiles.get_tile(tmp_path, "slope", z, x, y) is not None + (tmp_path / tiles.TILE_DIRNAME / tiles._SEEN_FILE).unlink() # version précédente + tiles._seen_cache.clear() + tiles.source_index(tmp_path, force=True) + assert tiles.cached_tile(tmp_path, "slope", z, x, y)[1] == "pending" + tiles.get_tile(tmp_path, "slope", z, x, y) + tiles._seen_cache.clear() + tiles.source_index(tmp_path, force=True) + assert tiles.cached_tile(tmp_path, "slope", z, x, y)[1] == "fresh" diff --git a/lidar_pipeline/tiles.py b/lidar_pipeline/tiles.py index 0f12473..bacd69b 100644 --- a/lidar_pipeline/tiles.py +++ b/lidar_pipeline/tiles.py @@ -69,6 +69,12 @@ FETCH_WORKERS = max(1, int(os.environ.get("LIDAR_TILE_FETCH_WORKERS", "2") or 2) _fetch_sem = threading.Semaphore(FETCH_WORKERS) _fetch_locks = {} _fetch_guard = threading.Lock() +# Coupe-circuit des sources : après un échec, plus aucune tentative pendant +# _SOURCE_OFFLINE_S — sans lui, chaque source d'une file de rapatriement +# repaierait le délai réseau d'un amont éteint et la maintenance ne rendrait +# plus rien (la carte doit rester autonome). +_SOURCE_OFFLINE_S = 60.0 +_SOURCE_OFFLINE = {"until": 0.0} # Index des sources reconstruit au plus toutes les _INDEX_TTL secondes (ou dès # qu'un dossier change de mtime : ajout/suppression de dalle). @@ -201,7 +207,7 @@ class _Source: pour la péremption des tuiles, avant même tout téléchargement. """ - __slots__ = ("path", "bounds", "res", "url", "version") + __slots__ = ("path", "bounds", "res", "url", "version", "seen") def __init__(self, path, bounds, res, url=None, version=None): self.path = path @@ -209,20 +215,35 @@ class _Source: self.res = res # résolution nominale (m/px) self.url = url self.version = version + self.seen = 0.0 # 1re apparition de la dalle dans l'inventaire def mtime(self): + """Date de référence pour la péremption des tuiles : version de la + source, repoussée à sa première apparition dans l'inventaire (une + dalle écrite avant mais inventoriée après une tuile la périme).""" if self.version is not None: - return self.version / 1000.0 - try: - return self.path.stat().st_mtime - except OSError: - return None + base = self.version / 1000.0 + else: + try: + base = self.path.stat().st_mtime + except OSError: + return None + return max(base, self.seen) def ensure(self): """Garantit la présence locale de l'image (rapatriement si besoin).""" if self.url is None or self.path.is_file(): return self.path.is_file() - return _fetch_source(self.url, self.path) + if not _fetch_source(self.url, self.path): + return False + if self.version is not None: + # Datée à sa version amont : l'index local (amont éteint) retrouve + # la même date et les tuiles déjà faites restent fraîches. + try: + os.utime(self.path, (self.version / 1000.0, self.version / 1000.0)) + except OSError: + pass + return True def _remote_thumb_px(k): @@ -238,6 +259,8 @@ def _fetch_source(url, dest): with lock: if dest.is_file(): return True + if time.time() < _SOURCE_OFFLINE["until"]: + return False try: req = urllib.request.Request( url, headers={"User-Agent": "lidar-maps-source"}) @@ -247,6 +270,7 @@ def _fetch_source(url, dest): data = r.read() except Exception as e: # noqa: BLE001 — amont éteint : tuile partielle logger.debug(f"Source amont indisponible ({url}) : {e}") + _SOURCE_OFFLINE["until"] = time.time() + _SOURCE_OFFLINE_S return False _write_atomic(dest, data) logger.info(f"Source rapatriée : {dest.name} ({len(data) / 1e6:.1f} Mo)") @@ -258,13 +282,15 @@ def _remote_payload(force=False): import json import urllib.request now = time.time() - if not force and _remote_cache["payload"] is not None \ + # Échec mémorisé lui aussi (TTL) : amont éteint sans inventaire connu, un + # appel par TTL au plus — pas un délai réseau à chaque requête de la carte. + if not force and _remote_cache["at"] \ and now - _remote_cache["at"] < _REMOTE_TTL: return _remote_cache["payload"] try: req = urllib.request.Request(f"{REMOTE_SOURCE_URL}/api/tiles", headers={"User-Agent": "lidar-maps-source"}) - with urllib.request.urlopen(req, timeout=30) as r: + with urllib.request.urlopen(req, timeout=10) as r: payload = json.loads(r.read().decode("utf-8")) except Exception as e: # noqa: BLE001 — on garde le dernier index connu logger.warning(f"Index amont injoignable ({REMOTE_SOURCE_URL}) : {e}") @@ -481,6 +507,61 @@ def _build_index(output_dir): return layers +_SEEN_FILE = ".sources_seen.json" +_seen_cache = {} +_seen_lock = threading.Lock() + + +def _apply_seen(output_dir, layers): + """Date chaque dalle de sa première apparition dans l'inventaire. + + Une dalle peut entrer dans l'inventaire APRÈS le rendu d'une tuile qui la + couvre tout en portant une date plus ancienne (écrite avant, inventoriée + après : TTL de l'index amont, anti-rebond de l'inventaire). Comparer les + seules dates de version laisserait la tuile « fraîche » — trou permanent + à ce niveau, jusque dans le navigateur (stamp inchangé). Registre + persistant (index_xyz/.sources_seen.json) : un redémarrage ne périme + rien ; au tout premier inventaire, rien n'est daté (tout est à rendre). + """ + import json + path = Path(output_dir) / TILE_DIRNAME / _SEEN_FILE + with _seen_lock: + seen = _seen_cache.get(str(path)) + first = False + if seen is None: + try: + seen = {k: float(v) for k, v in + json.loads(path.read_text(encoding="utf-8")).items()} + except (OSError, ValueError, AttributeError): + # Pas de registre : tout premier inventaire (rien à périmer), + # sauf si un cache de tuiles existe déjà (mise à jour) — ses + # tuiles ont pu être faites avant l'arrivée de dalles : toutes + # sont périmées une fois (l'ancienne reste servie en attendant). + cache_dir = path.parent + upgraded = cache_dir.is_dir() and any( + p.is_dir() for p in cache_dir.iterdir()) + seen, first = {}, not upgraded + _seen_cache[str(path)] = seen + now = time.time() + added = False + for layer, per_cell in layers.items(): + for (col, row), tiers in per_cell.items(): + key = f"{layer}|{col}|{row}" + at = seen.get(key) + if at is None: + at = seen[key] = 0.0 if first else now + added = True + if at: + for group in tiers: + for src in group: + src.seen = at + if added or first: + try: + _write_atomic(path, json.dumps(seen).encode("utf-8")) + except Exception as e: # noqa: BLE001 — registre perdu : péremption par version seule + logger.debug(f"Registre des dalles non écrit ({path}) : {e}") + + def source_index(output_dir, force=False): """Index des sources, mémoïsé (TTL + mtime des dossiers surveillés).""" output_dir = Path(output_dir) @@ -504,6 +585,7 @@ def source_index(output_dir, force=False): layers = _build_index(output_dir) # amont muet : cache local seul else: layers = _build_index(output_dir) + _apply_seen(output_dir, layers) _index_cache[key] = {"layers": layers, "stamp": stamp, "at": now} return layers