Peupler le cache de la webapp à la visualisation, sans rsync

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.
This commit is contained in:
Antoine Jacquin
2026-09-12 19:49:14 +02:00
parent 68272cabf0
commit bde2025714
2 changed files with 357 additions and 14 deletions

View File

@ -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

View File

@ -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 -------------------