Sélection de tuiles au clic, file d'attente et rafraîchissement fiable

- Clic sur la carte pour choisir des dalles précises, même sans données
  existantes (via /api/cell) ; les zones dessinées s'ajoutent à la
  sélection au lieu de la remplacer
- Une demande lancée pendant un run part en file d'attente côté serveur
  et démarre à la fin du travail en cours (plus de refus 409), file
  vidable depuis l'interface
- Tuiles et interface toujours fraîches : images servies sans cache
  navigateur (revalidation 304), rechargement de la carte seulement une
  fois le rebuild de l'index terminé
- serve-webapp.sh : sous-commande update retirée, documentation de
  déploiement corrigée en conséquence
This commit is contained in:
Antoine Jacquin
2026-09-16 22:37:53 +02:00
parent 26c05319fd
commit f024514427
6 changed files with 636 additions and 128 deletions

View File

@ -85,6 +85,7 @@ OUTPUT_DIR = Path(os.environ.get("LIDAR_OUTPUT_DIR", "/data/output"))
PORT = int(os.environ.get("LIDAR_PORT", "8973"))
LOG_FILE = OUTPUT_DIR / ".generation.log"
JOB_FILE = OUTPUT_DIR / ".generation.job.json"
QUEUE_FILE = OUTPUT_DIR / ".generation.queue.json"
MAX_CELLS = 400 # garde-fou : ~400 km² max par demande (zones dessinées)
MAX_CELLS_ALL = 10000 # garde-fou : passe globale « tout compléter » (dalles déjà présentes)
MAX_EXPORT_TILES = 64 # garde-fou : dalles assemblables par export (mémoire)
@ -245,7 +246,25 @@ def _fetch_remote_to_cache(rel_path):
return False
class _OnDemandStaticFiles(_StaticFiles):
class _NoCacheStaticFiles(_StaticFiles):
"""Fichiers statiques servis sans cache navigateur (revalidation 304).
Les images (visualisations, vignettes, sous-tuiles, DTM) n'ont pas
toujours une URL invalidée après régénération : le ?v= suit la mtime de
la SOURCE, pas celle du fichier servi (vignette recalculée après coup,
cache webapp rapatrié…). Sans Cache-Control, le cache heuristique du
navigateur peut continuer d'afficher l'ancien rendu sur la même URL.
no-cache force la revalidation à chaque affichage — ETag/Last-Modified
la rend quasi gratuite quand le fichier n'a pas changé (réponse 304).
"""
async def get_response(self, path, scope):
resp = await super().get_response(path, scope)
resp.headers["Cache-Control"] = "no-cache, must-revalidate"
return resp
class _OnDemandStaticFiles(_NoCacheStaticFiles):
"""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-
@ -294,7 +313,7 @@ for name in ("index_thumbs", "index_subtiles", "visualisations", "DTM"):
_OnDemandStaticFiles(directory=str(_dir), url_prefix=name),
name=name)
else:
app.mount(f"/{name}", _StaticFiles(directory=str(_dir)), name=name)
app.mount(f"/{name}", _NoCacheStaticFiles(directory=str(_dir)), name=name)
@app.get("/assets/{file_path:path}")
@ -507,7 +526,7 @@ class ExportRequest(BaseModel):
# --- État du job de génération -------------------------------------------
_job = {"proc": None, "started": None, "returncode": None, "cmd": None, "finished": None}
_job = {"proc": None, "started": None, "returncode": None, "cmd": None, "finished": None, "qid": None}
_job_lock = threading.Lock()
@ -518,7 +537,7 @@ def _save_job_state():
plus le dernier run et la file de génération reste vide au rechargement.
Best-effort : n'acquiert pas le verrou (appelé aussi sous verrou).
"""
state = {k: _job.get(k) for k in ("started", "returncode", "cmd", "finished")}
state = {k: _job.get(k) for k in ("started", "returncode", "cmd", "finished", "qid")}
try:
JOB_FILE.write_text(json.dumps(state), encoding="utf-8")
except OSError:
@ -536,10 +555,54 @@ def _load_job_state():
_job["returncode"] = state.get("returncode")
_job["cmd"] = state.get("cmd")
_job["finished"] = state.get("finished")
_job["qid"] = state.get("qid")
_load_job_state()
# --- File d'attente des demandes de génération -----------------------------
# Une demande arrivée pendant un run n'est plus rejetée (ancien 409) : elle
# part en file et démarre automatiquement à la fin du run en cours — le
# travail déjà lancé n'est jamais coupé. La file survit aux redémarrages du
# serveur (fichier .generation.queue.json) et est vidable via /api/queue/clear.
_queue = []
_queue_seq = 0
def _save_queue():
"""Écrit la file d'attente sur disque (atomique, best-effort)."""
try:
tmp = QUEUE_FILE.with_suffix(".tmp")
tmp.write_text(json.dumps(_queue), encoding="utf-8")
os.replace(tmp, QUEUE_FILE)
except OSError:
pass
def _load_queue():
"""Recharge la file d'attente au démarrage (et son compteur d'ids)."""
global _queue_seq
try:
data = json.loads(QUEUE_FILE.read_text(encoding="utf-8"))
if isinstance(data, list):
_queue[:] = [it for it in data
if isinstance(it, dict) and isinstance(it.get("req"), dict)]
except (OSError, ValueError):
pass
_queue_seq = max((it.get("id", 0) for it in _queue), default=0)
def _queue_summary():
"""Aperçu de la file pour /api/status (sans la requête brute)."""
return [{"id": it.get("id"),
"tuiles": len(it["req"].get("tiles") or []),
"toutes": bool(it["req"].get("all_missing")),
"en_attente_depuis": it.get("queued_at")}
for it in _queue]
_load_queue()
# --- Presets de couches (partagés entre navigateurs) --------------------------
# Stockés dans OUTPUT_DIR/.presets.json : un preset est un jeu nommé de
@ -665,6 +728,28 @@ def bbox_to_cells(w, s, e, n):
return [(c, r) for r in rows for c in cols]
def point_to_cell(lon, lat):
"""Cellule L93 de 1 km (col, row) contenant un point WGS84.
Cohérent avec bbox_to_cells : une cellule (col, row) couvre
X ∈ [col, col+1] km, Y ∈ [row-1, row] km — sert au clic carte pour
sélectionner une dalle précise, même sans données existantes.
"""
from .index import _approx_wgs84_to_l93
try:
from rasterio.warp import transform as warp_transform
x, y = warp_transform("EPSG:4326", "EPSG:2154", [lon], [lat])
x, y = x[0], y[0]
except ImportError:
try:
from pyproj import Transformer
transformer = Transformer.from_crs("EPSG:4326", "EPSG:2154", always_xy=True)
x, y = transformer.transform(lon, lat)
except ImportError:
x, y = _approx_wgs84_to_l93(lon, lat)
return int(math.floor(x / 1000)), int(math.floor(y / 1000)) + 1
def processed_cells(output_dir):
"""Ensemble des cellules (col, row) ayant déjà des visualisations."""
from .index import scan_tiles
@ -764,6 +849,12 @@ def status(request: Request = None):
data["regen_allowed"] = regen_allowed
return data
from .progress import progress_snapshot
with _job_lock:
pending = bool(_queue)
if pending:
# File d'attente non vide et serveur libre (ex. redémarrage du
# conteneur en cours de file) : la demande suivante repart.
_start_next_queued()
with _job_lock:
proc = _job["proc"]
running = proc is not None and proc.poll() is None
@ -772,6 +863,8 @@ def status(request: Request = None):
"started": _job["started"],
"returncode": _job["returncode"],
"cmd": _job["cmd"],
"qid": _job.get("qid"),
"queue": _queue_summary(),
"log": _tail_log(40),
# L'interface masque les boutons de génération hors réseau local
"regen_allowed": regen_allowed,
@ -896,15 +989,24 @@ _rebuild = {"running": False, "error": None, "done": None, "phase": None}
def _start_rebuild():
"""Lance en arrière-plan le rebuild de l'index (vignettes + carte)."""
"""Lance en arrière-plan le rebuild de l'index (vignettes + carte).
L'état running est posé de façon SYNCHRONE, avant le fil : un GET
/api/sync juste après le POST doit voir le rebuild en cours — sinon il
lirait le done du rebuild précédent et l'interface rechargerait une
carte périmée. Le timestamp started retourné permet au client
d'attendre la fin de CE rebuild précis (done >= started).
"""
if _rebuild["running"]:
raise HTTPException(409, "un rebuild est déjà en cours")
started = time.time()
_rebuild.update({"running": True, "error": None, "done": None,
"phase": "index"})
# Index distant (mode deux machines) : cache 60 s vidé pour que le
# premier /api/tiles suivant reflète l'état du worker sans attendre.
_REMOTE_INDEX["fetched"] = 0.0
def _run():
_rebuild["running"] = True
_rebuild["error"] = None
_rebuild["done"] = None
_rebuild["phase"] = "index"
try:
from .index import build_index
build_index(OUTPUT_DIR)
@ -916,7 +1018,7 @@ def _start_rebuild():
_rebuild["done"] = time.time()
threading.Thread(target=_run, daemon=True).start()
return {"demarré": True}
return {"demarré": True, "started": started}
@app.post("/api/sync", dependencies=[Depends(_require_token)])
@ -984,6 +1086,24 @@ def preview(req: PreviewRequest):
return {"count": len(todo), "capped": capped, "cells": todo}
@app.get("/api/cell")
def cell_at_point(lat: float, lng: float):
"""Dalle L93 de 1 km contenant un point WGS84 (clic sur la carte).
Retourne col, row et les coins GPS de la dalle — permet de sélectionner
une tuile individuelle même sans données existantes (téléchargement +
génération à la validation).
"""
if GENERATION_URL:
# Webapp légère : la conversion précise vit sur la machine de traitement
return _proxy_api("GET", f"/api/cell?lat={lat}&lng={lng}")
from .index import attach_gps_bounds
col, row = point_to_cell(lng, lat)
tile = {"col": col, "row": row}
attach_gps_bounds([tile])
return {"col": col, "row": row, "corners": tile.get("corners") or []}
def _build_command(tiles, regenerate=False, ground_class="ign", bare_earth=False, ign_classes="sol", viz=None, reclassify=False):
"""Commande de génération : téléchargement IGN + traitement des fichiers.
@ -1029,15 +1149,13 @@ def _build_command(tiles, regenerate=False, ground_class="ign", bare_earth=False
return cmd
@app.post("/api/generate", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def generate(req: GenerateRequest):
if GENERATION_URL:
# Webapp légère : téléchargement + traitement sur la machine distante,
# qui valide les paramètres et renvoie son propre état de job.
data = _proxy_api("POST", "/api/generate", json.loads(req.json()))
data["distant"] = True
return data
def _resolve_request(req):
"""Valide une demande de génération et résout sa liste de tuiles.
Retourne (tuiles [(col, row)], viz [noms d'étapes]). Réexécutée au
démarrage de chaque demande sortie de file : les dalles présentes dans
input/ peuvent avoir changé entre la mise en file et le départ du run.
"""
names = _viz_step_names()
viz = [v for v in (req.viz or []) if v] or ["aspect"]
invalid = [v for v in viz if v not in names]
@ -1067,41 +1185,124 @@ def generate(req: GenerateRequest):
raise HTTPException(
400, f"méthode de classification invalide : {req.ground_class!r} "
f"(attendu : {', '.join(GROUND_CLASS_METHODS)})")
return tiles, viz
def _launch_job(tiles, viz, req, qid=None):
"""Démarre un run du pipeline. À appeler sous _job_lock, serveur libre.
Retourne la commande lancée. qid : identifiant de file d'attente quand
la demande vient de la file (suivi côté interface).
"""
cmd = _build_command(tiles, regenerate=req.regenerate, ground_class=req.ground_class,
bare_earth=req.bare_earth, ign_classes=req.ign_classes,
viz=viz, reclassify=req.reclassify)
LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
log_fh = open(LOG_FILE, "w", encoding="utf-8")
# Nouveau run : journal d'événements remis à zéro (les tuiles affichées
# dans la file correspondent au run qui démarre, pas au précédent)
from .progress import reset_events
reset_events(OUTPUT_DIR)
_job.update({"proc": None, "started": time.time(), "returncode": None,
"cmd": cmd, "finished": None, "qid": qid})
_save_job_state()
# start_new_session : le pipeline et ses workers/PDAL forment leur
# propre groupe de processus — /api/stop peut le tuer en bloc sans
# toucher la webapp, et le gestionnaire SIGTERM du pipeline (killpg)
# reste confiné à son groupe.
p = subprocess.Popen(cmd, stdout=log_fh, stderr=subprocess.STDOUT,
cwd="/app" if Path("/app").exists() else None,
start_new_session=True)
def _watch():
rc = p.wait()
log_fh.close()
with _job_lock:
_job["returncode"] = rc
_job["finished"] = time.time()
_save_job_state()
# Le run est fini : la demande suivante de la file peut partir.
_start_next_queued()
threading.Thread(target=_watch, daemon=True).start()
_job["proc"] = p
return cmd
def _start_next_queued():
"""Démarre la première demande de la file dès que le serveur est libre.
Les demandes devenues invalides (ex. plus aucune tuile à compléter)
sont ignorées ; une demande repoussée en tête de file si un run part
entre-temps. Ne coupe jamais le run en cours.
"""
while True:
with _job_lock:
proc = _job["proc"]
if proc is not None and proc.poll() is None:
return False
if not _queue:
return False
item = _queue.pop(0)
_save_queue()
try:
req = GenerateRequest(**item["req"])
tiles, viz = _resolve_request(req)
except HTTPException as e:
logger.warning("Demande en file ignorée : %s", e.detail)
continue
with _job_lock:
proc = _job["proc"]
if proc is not None and proc.poll() is None:
_queue.insert(0, item) # un autre run est parti : remise en tête
_save_queue()
return False
_launch_job(tiles, viz, req, qid=item.get("id"))
return True
@app.post("/api/generate", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def generate(req: GenerateRequest):
if GENERATION_URL:
# Webapp légère : téléchargement + traitement sur la machine distante,
# qui valide les paramètres, met en file et renvoie son propre état.
data = _proxy_api("POST", "/api/generate", json.loads(req.json()))
data["distant"] = True
return data
tiles, viz = _resolve_request(req)
with _job_lock:
proc = _job["proc"]
if proc is not None and proc.poll() is None:
raise HTTPException(409, "une génération est déjà en cours")
cmd = _build_command(tiles, regenerate=req.regenerate, ground_class=req.ground_class, bare_earth=req.bare_earth, ign_classes=req.ign_classes, viz=viz, reclassify=req.reclassify)
LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
log_fh = open(LOG_FILE, "w", encoding="utf-8")
# Nouveau run : journal d'événements remis à zéro (les tuiles affichées
# dans la file correspondent au run qui démarre, pas au précédent)
from .progress import reset_events
reset_events(OUTPUT_DIR)
_job.update({"proc": None, "started": time.time(), "returncode": None, "cmd": cmd, "finished": None})
_save_job_state()
# start_new_session : le pipeline et ses workers/PDAL forment leur
# propre groupe de processus — /api/stop peut le tuer en bloc sans
# toucher la webapp, et le gestionnaire SIGTERM du pipeline (killpg)
# reste confiné à son groupe.
p = subprocess.Popen(cmd, stdout=log_fh, stderr=subprocess.STDOUT,
cwd="/app" if Path("/app").exists() else None,
start_new_session=True)
def _watch():
rc = p.wait()
log_fh.close()
with _job_lock:
_job["returncode"] = rc
_job["finished"] = time.time()
_save_job_state()
threading.Thread(target=_watch, daemon=True).start()
_job["proc"] = p
running = proc is not None and proc.poll() is None
if running:
# Un run est en cours : la demande part en file d'attente et
# démarrera à sa fin — jamais de coupure du travail en place.
global _queue_seq
_queue_seq += 1
item = {"id": _queue_seq, "req": json.loads(req.json()),
"queued_at": time.time()}
_queue.append(item)
_save_queue()
return {"demarré": False, "en_file": len(_queue), "qid": item["id"],
"tuiles": len(tiles)}
cmd = _launch_job(tiles, viz, req)
return {"demarré": True, "tuiles": len(tiles), "commande": " ".join(cmd)}
@app.post("/api/queue/clear", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def queue_clear():
"""Retire les demandes en attente (le run en cours n'est pas touché)."""
if GENERATION_URL:
# Webapp légère : la file vit sur la machine de traitement
return _proxy_api("POST", "/api/queue/clear")
with _job_lock:
n = len(_queue)
_queue.clear()
_save_queue()
return {"retirées": n}
@app.post("/api/stop", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def stop_generation():