From bde2025714e5412f455af88e870f912d8f8d5324 Mon Sep 17 00:00:00 2001 From: Antoine Jacquin Date: Sat, 12 Sep 2026 19:49:14 +0200 Subject: [PATCH] =?UTF-8?q?Peupler=20le=20cache=20de=20la=20webapp=20?= =?UTF-8?q?=C3=A0=20la=20visualisation,=20sans=20rsync?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Les images (visualisations, vignettes, sous-tuiles) sont rapatriées depuis la machine de traitement par HTTP au moment où elles sont affichées, puis servies depuis le disque : une image absente est téléchargée, une image périmée ( ?v= plus récent que sa mtime, tuile régénérée) est rafraîchie. /api/tiles fusionne l'index du worker (cache 60 s) avec le cache local, si bien que les tuiles non cachées apparaissent quand même sur la carte. Worker hors ligne : disjoncteur 2 min après 3 échecs, index local seul, /api/status répond running=null — aucune erreur cliente. --- lidar_pipeline/tests/test_webapp.py | 133 ++++++++++++++++ lidar_pipeline/webapp.py | 238 ++++++++++++++++++++++++++-- 2 files changed, 357 insertions(+), 14 deletions(-) diff --git a/lidar_pipeline/tests/test_webapp.py b/lidar_pipeline/tests/test_webapp.py index 8f1a950..302b2b7 100644 --- a/lidar_pipeline/tests/test_webapp.py +++ b/lidar_pipeline/tests/test_webapp.py @@ -293,6 +293,139 @@ def test_tiles_endpoint_stamp_diffing(tmp_path, monkeypatch): assert webapp.tiles_data(stamp=full["stamp"] - 10)["tiles"] is not None +def test_tiles_endpoint_merges_remote_index(tmp_path, monkeypatch): + """L'index du worker prime ; le cache local complète les positions perdues.""" + import json + import lidar_pipeline.webapp as webapp + monkeypatch.setattr(webapp, "OUTPUT_DIR", tmp_path) + monkeypatch.setattr(webapp, "GENERATION_URL", "http://distant:8973") + (tmp_path / "index_tiles.json").write_text(json.dumps({ + "tiles": [ + {"col": 1054, "row": 6882, "resolution": 0.2, "viz": {}}, # commun + {"col": 1050, "row": 6882, "resolution": 0.2, "viz": {}}, # local seul + ], + "viz_meta": {"local": {"label": "Locale"}}, + "stats": {}, + }), encoding="utf-8") + import os + remote_stamp = os.stat(tmp_path / "index_tiles.json").st_mtime + 1000 + monkeypatch.setattr(webapp, "_remote_tiles_data", lambda: { + "stamp": remote_stamp, + "tiles": [ + # 1054 régénéré côté worker : sa version (viz remplie) doit primer + {"col": 1054, "row": 6882, "resolution": 0.2, "viz": {"slope": {}}}, + {"col": 1055, "row": 6882, "resolution": 0.2, "viz": {}}, # distant seul + ], + "viz_meta": {"slope": {"label": "Pente"}}, + }) + d = webapp.tiles_data() + entries = {(t["col"], json.dumps(t["viz"], sort_keys=True)) for t in d["tiles"]} + assert (1054, '{"slope": {}}') in entries # version distante retenue + assert (1055, '{}') in entries # tuile non cachée visible + assert (1050, '{}') in entries # historique local conservé + assert d["viz_meta"]["slope"]["label"] == "Pente" # distant prime + assert d["viz_meta"]["local"]["label"] == "Locale" # local complété + assert d["stamp"] == remote_stamp # max des deux stamps + + +def test_tiles_endpoint_worker_offline_serves_local(tmp_path, monkeypatch): + """Worker injoignable : index local servi seul, sans erreur.""" + import json + import lidar_pipeline.webapp as webapp + monkeypatch.setattr(webapp, "OUTPUT_DIR", tmp_path) + monkeypatch.setattr(webapp, "GENERATION_URL", "http://distant:8973") + (tmp_path / "index_tiles.json").write_text( + json.dumps({"tiles": [{"col": 1054}], "viz_meta": {}, "stats": {}}), + encoding="utf-8") + monkeypatch.setattr(webapp, "_remote_tiles_data", lambda: None) + d = webapp.tiles_data() + assert d["tiles"][0]["col"] == 1054 + assert d["stamp"] is not None + + +def test_worker_offline_breaker(monkeypatch): + """3 échecs consécutifs → tentatives suspendues ; un succès réarme.""" + import lidar_pipeline.webapp as webapp + monkeypatch.setattr(webapp, "_WORKER_HEALTH", {"fails": 0, "until": 0.0}) + webapp._worker_mark(False) + webapp._worker_mark(False) + assert webapp._worker_offline() is False # pas encore armé + webapp._worker_mark(False) + assert webapp._worker_offline() is True # disjoncteur armé + until = webapp._WORKER_HEALTH["until"] + webapp._worker_mark(False) + assert webapp._WORKER_HEALTH["until"] == until # fenêtre non étendue + webapp._worker_mark(True) + assert webapp._worker_offline() is False # réarmé immédiatement + + +def test_safe_rel_path(): + """Seuls les chemins images des préfixes autorisés sont rapatrieables.""" + from lidar_pipeline.webapp import _safe_rel_path + assert _safe_rel_path("/visualisations/dalle/img.avif") == \ + "visualisations/dalle/img.avif" + assert _safe_rel_path("/index_thumbs/a_b.jpg") == "index_thumbs/a_b.jpg" + assert _safe_rel_path("/index_subtiles/x_0_1.avif") == "index_subtiles/x_0_1.avif" + assert _safe_rel_path("/DTM/x.tif") is None # préfixe interdit + assert _safe_rel_path("/assets/app.js") is None # préfixe interdit + assert _safe_rel_path("/visualisations/../x.avif") is None # traversal + assert _safe_rel_path("/visualisations/a%20b.avif") is None # caractères + + +def test_fetch_remote_to_cache(tmp_path, monkeypatch): + """Rapatriement atomique depuis le worker ; échec sans fichier ni erreur.""" + import io + import lidar_pipeline.webapp as webapp + monkeypatch.setattr(webapp, "OUTPUT_DIR", tmp_path) + monkeypatch.setattr(webapp, "GENERATION_URL", "http://distant:8973") + monkeypatch.setattr(webapp, "_WORKER_HEALTH", {"fails": 0, "until": 0.0}) + + class _Resp(io.BytesIO): + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + captured = {} + + def fake_urlopen(req, timeout=None): + captured["url"] = req.full_url + return _Resp(b"IMAGEDATA") + + monkeypatch.setattr(webapp.urllib.request, "urlopen", fake_urlopen) + assert webapp._fetch_remote_to_cache("visualisations/dalle/img.avif") is True + p = tmp_path / "visualisations" / "dalle" / "img.avif" + assert p.read_bytes() == b"IMAGEDATA" + assert captured["url"] == "http://distant:8973/visualisations/dalle/img.avif" + assert not list(tmp_path.rglob("*.part")) # pas de résidu temporaire + assert webapp._WORKER_HEALTH["fails"] == 0 # succès → disjoncteur sain + + def boom(req, timeout=None): + raise OSError("injoignable") + + monkeypatch.setattr(webapp.urllib.request, "urlopen", boom) + assert webapp._fetch_remote_to_cache("visualisations/dalle/x.avif") is False + assert not (tmp_path / "visualisations" / "dalle" / "x.avif").exists() + assert webapp._WORKER_HEALTH["fails"] == 1 + + +def test_status_worker_offline_no_error(monkeypatch): + """Worker hors ligne : running=null (ignoré par l'UI), pas de 503.""" + from fastapi import HTTPException + import lidar_pipeline.webapp as webapp + monkeypatch.setattr(webapp, "GENERATION_URL", "http://distant:8973") + + def unreachable(method, path, payload=None, timeout=30): + raise HTTPException(503, "machine de traitement injoignable") + + monkeypatch.setattr(webapp, "_proxy_api", unreachable) + d = webapp.status(_FakeRequest(host="192.168.1.7")) + assert d["running"] is None + assert d["worker_offline"] is True + assert d["regen_allowed"] is True + + def test_viz_step_names_match_pipeline(): """Les noms acceptés par l'API sont exactement les étapes VIZ_STEPS du pipeline.""" from lidar_pipeline.webapp import _viz_step_names diff --git a/lidar_pipeline/webapp.py b/lidar_pipeline/webapp.py index d71aac0..3444971 100644 --- a/lidar_pipeline/webapp.py +++ b/lidar_pipeline/webapp.py @@ -51,10 +51,13 @@ import json import logging import math import os +import re import subprocess import sys import threading import time +import urllib.parse +import urllib.request from pathlib import Path from typing import Optional @@ -132,11 +135,152 @@ if SYNC_CMD: # statique soit actif même avant la première génération de l'index. _assets_dir = OUTPUT_DIR / "assets" _assets_dir.mkdir(parents=True, exist_ok=True) + +# --- Cache à la demande : la visualisation peuple le cache local ------------- +# Sans LIDAR_SYNC_CMD (pas de rsync), les images viennent de la machine de +# traitement par HTTP : une image absente du disque — ou plus ancienne que la +# version ?v= référencée (mtime ms de la source sur le worker) — est téléchar- +# gée depuis GENERATION_URL, écrite sur disque puis servie. Requêtes concur- +# rentes dédupliquées (verrou par fichier), flux réseau plafonné (sémaphore). +from fastapi.staticfiles import StaticFiles as _StaticFiles +from starlette.concurrency import run_in_threadpool + +_ONDEMAND_PREFIXES = ("visualisations", "index_thumbs", "index_subtiles") +_FETCH_SEMAPHORE = threading.Semaphore(4) +_FETCH_LOCKS = {} +_FETCH_LOCKS_GUARD = threading.Lock() +_MAX_CACHE_BYTES = 256 * 1024 * 1024 # garde-fou par fichier (dalle AVIF ~35 Mo) + +# Disjoncteur : la machine de traitement n'est pas toujours allumée. Après +# quelques échecs consécutifs, les tentatives (index + images) sont suspendues +# 2 min — réponses locales instantanées, aucun martèlement ni erreur cliente. +_WORKER_HEALTH = {"fails": 0, "until": 0.0} +_WORKER_OFFLINE_AFTER = 3 +_WORKER_OFFLINE_SECONDS = 120.0 + + +def _worker_offline(): + return time.time() < _WORKER_HEALTH["until"] + + +def _worker_mark(ok): + if ok: + _WORKER_HEALTH.update(fails=0, until=0.0) + else: + _WORKER_HEALTH["fails"] += 1 + if (_WORKER_HEALTH["fails"] >= _WORKER_OFFLINE_AFTER + and _WORKER_HEALTH["until"] == 0.0): + _WORKER_HEALTH["until"] = time.time() + _WORKER_OFFLINE_SECONDS + logger.info("Machine de traitement injoignable — " + "rapatriement à la demande suspendu 2 min") + + +def _version_from_query(scope): + """Version ?v= (mtime ms de la source) demandée par le client, ou None.""" + try: + qs = urllib.parse.parse_qs(scope.get("query_string", b"").decode("ascii")) + v = qs.get("v", [None])[0] + return int(v) if v and v.isdigit() else None + except Exception: + return None + + +def _safe_rel_path(url_path): + """Chemin relatif validé (préfixe autorisé, composants sains), ou None.""" + rel = url_path.lstrip("/") + prefix = rel.split("/", 1)[0] + if prefix not in _ONDEMAND_PREFIXES: + return None + parts = [p for p in rel.split("/") if p not in ("", ".")] + if not parts or any(p == ".." or not re.fullmatch(r"[A-Za-z0-9_.-]+", p) + for p in parts): + return None + return "/".join(parts) + + +def _fetch_remote_to_cache(rel_path): + """Télécharge GENERATION_URL/rel_path vers OUTPUT_DIR/rel_path (atomique).""" + if _worker_offline(): + return False + dest = OUTPUT_DIR / rel_path + dest.parent.mkdir(parents=True, exist_ok=True) + tmp = dest.with_name(dest.name + ".part") + url = f"{GENERATION_URL}/{urllib.parse.quote(rel_path)}" + try: + with _FETCH_SEMAPHORE: + req = urllib.request.Request(url, headers={"User-Agent": "lidar-webapp-cache"}) + with urllib.request.urlopen(req, timeout=60) as r, open(tmp, "wb") as f: + total = 0 + while True: + chunk = r.read(1 << 20) + if not chunk: + break + total += len(chunk) + if total > _MAX_CACHE_BYTES: + raise IOError(f"fichier trop volumineux : {rel_path}") + f.write(chunk) + os.replace(tmp, dest) + _worker_mark(True) + logger.info(f"Cache à la demande : {rel_path} rapatrié " + f"({total / 1e6:.1f} Mo)") + return True + except Exception as e: + _worker_mark(False) + tmp.unlink(missing_ok=True) + logger.debug(f"Cache à la demande : échec {rel_path} : {e}") + return False + + +class _OnDemandStaticFiles(_StaticFiles): + """Montage statique qui peuple le cache depuis la machine de traitement. + + Fichier local absent, ou périmé (?v= plus récent que sa mtime) → téléchar- + gement puis service ; sans GENERATION_URL ou en cas d'échec, comportement + statique normal (404, ou version locale périmée plutôt que rien). + + url_prefix : préfixe du montage (ex. visualisations) — dans une app + montée, scope["path"] est le sous-chemin sans préfixe, on le reconstruit. + """ + + def __init__(self, *, directory, url_prefix): + super().__init__(directory=directory) + self.url_prefix = url_prefix + + async def get_response(self, path, scope): + want = None + if GENERATION_URL: + rel = _safe_rel_path(f"/{self.url_prefix}/{path}") + if rel is not None: + local = OUTPUT_DIR / rel + if not local.is_file(): + want = rel + else: + v = _version_from_query(scope) + if v is not None and int(local.stat().st_mtime * 1000) + 500 < v: + want = rel # tuile régénérée sur le worker depuis le cache + if want is not None: + with _FETCH_LOCKS_GUARD: + lock = _FETCH_LOCKS.setdefault(want, threading.Lock()) + # Téléchargement bloquant hors de la boucle d'événements ; le + # verrou par fichier déduplique les requêtes concurrentes. + await run_in_threadpool(self._fetch_under_lock, want, lock) + return await super().get_response(path, scope) + + @staticmethod + def _fetch_under_lock(rel, lock): + with lock: + return _fetch_remote_to_cache(rel) + + for name in ("index_thumbs", "index_subtiles", "visualisations", "DTM"): _dir = OUTPUT_DIR / name - if _dir.exists(): - from fastapi.staticfiles import StaticFiles - app.mount(f"/{name}", StaticFiles(directory=str(_dir)), name=name) + _dir.mkdir(parents=True, exist_ok=True) + if name in _ONDEMAND_PREFIXES and GENERATION_URL: + app.mount(f"/{name}", + _OnDemandStaticFiles(directory=str(_dir), url_prefix=name), + name=name) + else: + app.mount(f"/{name}", _StaticFiles(directory=str(_dir)), name=name) @app.get("/assets/{file_path:path}") @@ -492,8 +636,15 @@ def status(request: Request = None): regen_allowed = _ip_in_regen_cidr(_client_ip(request)) if GENERATION_URL: # Webapp légère : l'état vient de la machine de traitement ; le droit - # de génération reste jugé sur l'IP du navigateur (localement). - data = _proxy_api("GET", "/api/status", timeout=10) + # de génération reste jugé sur l'IP du navigateur (localement). Hors + # ligne : running=null (l'interface ignore, aucun faux « terminé »). + try: + data = _proxy_api("GET", "/api/status", timeout=10) + except HTTPException as e: + if e.status_code != 503: + raise + return {"running": None, "distant": True, "worker_offline": True, + "regen_allowed": regen_allowed} data["distant"] = True data["regen_allowed"] = regen_allowed return data @@ -535,25 +686,84 @@ def available_layers(): return found +_REMOTE_INDEX = {"data": None, "fetched": 0.0} +_REMOTE_INDEX_TTL = 60.0 + + +def _tile_entry_key(t): + """Clé d'unicité d'une tuile affichée (position, résolution, quadrant).""" + return (t.get("col"), t.get("row"), t.get("resolution"), + t.get("sub_i"), t.get("sub_j")) + + +def _remote_tiles_data(): + """Index de la machine de traitement (cache 60 s), ou None si injoignable. + + Une panne ne remonte pas : dernière réponse connue servie telle quelle + (stale-while-error), et le disjoncteur suspend les tentatives. + """ + if not GENERATION_URL or _worker_offline(): + return _REMOTE_INDEX["data"] + now = time.time() + if _REMOTE_INDEX["fetched"] > now - _REMOTE_INDEX_TTL: + return _REMOTE_INDEX["data"] + try: + with urllib.request.urlopen(f"{GENERATION_URL}/api/tiles", timeout=4) as r: + data = json.loads(r.read().decode("utf-8")) + _worker_mark(True) + _REMOTE_INDEX.update({"data": data, "fetched": now}) + except Exception as e: + _worker_mark(False) + _REMOTE_INDEX["fetched"] = now # ne pas marteler à chaque requête + logger.debug(f"Index distant injoignable : {e}") + return _REMOTE_INDEX["data"] + + @app.get("/api/tiles") def tiles_data(stamp: Optional[float] = None): - """Données de carte (index_tiles.json) — tuiles visibles sans recharger. + """Données de carte — cache local + index de la machine de traitement. Le pipeline lancé par /api/generate tourne avec --incremental-index : il - réécrit ce fichier après chaque tuile terminée. La carte sonde cet - endpoint pendant un run et fusionne les nouveautés. Avec ?stamp=X : la - réponse est allégée ({tiles: null}) tant que le fichier n'a pas changé. + réécrit index_tiles.json après chaque tuile terminée. La carte sonde cet + endpoint et fusionne les nouveautés. Sur la webapp légère, l'index du + worker (cache 60 s) fait autorité : les tuiles absentes du cache local y + figurent quand même — leurs images sont rapatriées à la visualisation par + le montage statique à la demande. Worker hors ligne : cache local seul, + sans erreur. Avec ?stamp=X : réponse allégée ({tiles: null}) si rien n'a + changé (le stamp combine les deux index). """ path = OUTPUT_DIR / "index_tiles.json" try: - current = path.stat().st_mtime - if stamp is not None and abs(current - stamp) < 1e-4: - return {"stamp": current, "tiles": None} + local_stamp = path.stat().st_mtime data = json.loads(path.read_text(encoding="utf-8")) except (OSError, ValueError): + local_stamp = None + data = None + remote = _remote_tiles_data() + remote_stamp = (remote or {}).get("stamp") + stamps = [s for s in (local_stamp, remote_stamp) if s is not None] + current = max(stamps) if stamps else None + if stamp is not None and current is not None and abs(current - stamp) < 1e-4: + return {"stamp": current, "tiles": None} + if data is None and remote is None: return {"stamp": None, "tiles": None} - data["stamp"] = current - return data + # L'index distant prime (tuiles régénérées : ?v= neuf → re-téléchargement) ; + # le local complète les positions que le worker n'a plus (historique). + merged = list((remote or {}).get("tiles") or []) + keys = {_tile_entry_key(t) for t in merged} + for t in (data or {}).get("tiles") or []: + k = _tile_entry_key(t) + if k not in keys: + merged.append(t) + keys.add(k) + viz_meta = dict((data or {}).get("viz_meta") or {}) + viz_meta.update((remote or {}).get("viz_meta") or {}) + return { + "tiles": merged, + "viz_meta": viz_meta, + "stats": (remote or data or {}).get("stats") or {}, + "stamp": current, + } # --- Synchronisation + rebuild de l'index en arrière-plan -------------------