La carte à tuiles XYZ devient la seule interface : l'API de génération (preview, generate, status, stop, file, cell), le job unique + file persistante, la délégation LIDAR_GENERATION_URL et les garde-fous (token, LIDAR_REGEN_CIDR) passent de webapp.py à mapserve.py ; l'interface gagne les boutons + Zone / ⤒ Compléter / régénération à la dalle, les options de run et la progression cadre par cadre. mapserve sert aussi l'inventaire /api/tiles et les dalles en statique versionné : le worker image complète remplace la webapp sur le 8973 du générateur, le Pi n'exécute plus que lidar-maps. index.py se réduit aux registres partagés + vignettes/sous-tuiles/inventaire (l'UI HTML/JS et les mosaïques d'overview partent avec export.py et les compose/scripts de la webapp).
168 lines
6.3 KiB
Python
168 lines
6.3 KiB
Python
"""Progression structurée d'un run de génération (contrat pipeline → interface).
|
|
|
|
Le pipeline et le téléchargement IGN écrivent des événements JSON ligne à
|
|
ligne dans <sortie>/.generation.events.jsonl ; mapserve les agrège via
|
|
/api/status pour afficher les tuiles une par une (nom, type, état) dans la
|
|
file de génération, à la place du journal brut.
|
|
|
|
Module volontairement découplé du serveur : le rendu des tuiles pourra être
|
|
déporté hors du serveur web — le producteur d'événements ne dépend que du
|
|
dossier de sortie. Les workers parallèles écrivent chacun leur ligne en mode
|
|
ajout (une ligne < 4 Ko avec O_APPEND est atomique sur POSIX).
|
|
"""
|
|
|
|
import json
|
|
import os
|
|
import re
|
|
import time
|
|
from pathlib import Path
|
|
|
|
EVENTS_FILENAME = ".generation.events.jsonl"
|
|
|
|
# Phases d'une tuile, dans l'ordre d'exécution.
|
|
PHASES = ("download", "classif", "dtm", "viz", "tile")
|
|
|
|
_TILE_RE = re.compile(r"LHD_FXX_(\d+)_(\d+)_")
|
|
|
|
|
|
def events_path(output_dir):
|
|
"""Chemin du journal d'événements d'un run."""
|
|
return Path(output_dir) / EVENTS_FILENAME
|
|
|
|
|
|
def tile_short_name(name):
|
|
"""Nom court d'une tuile : LHD_FXX_1054_6882_PTS_LAMB93_IGN69 → 1054-6882."""
|
|
m = _TILE_RE.search(name or "")
|
|
if not m:
|
|
return name
|
|
return f"{int(m.group(1))}-{int(m.group(2))}"
|
|
|
|
|
|
def reset_events(output_dir):
|
|
"""Prépare un journal vide pour un nouveau run (tronque l'ancien)."""
|
|
try:
|
|
path = events_path(output_dir)
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
path.write_text("", encoding="utf-8")
|
|
except OSError:
|
|
pass # progression best-effort : ne jamais bloquer le traitement
|
|
|
|
|
|
def report_event(output_dir, tile, phase, state, detail=None, res=None):
|
|
"""Ajoute un événement (une ligne JSON). Ne lève jamais.
|
|
|
|
Ouverture en mode ajout à chaque appel : les workers parallélisés
|
|
(ProcessPoolExecutor, 'spawn') ne partagent aucun descripteur de fichier.
|
|
"""
|
|
event = {"ts": round(time.time(), 3), "tile": tile, "phase": phase, "state": state}
|
|
run_id = os.environ.get("LIDAR_RUN_ID")
|
|
if run_id:
|
|
event["run"] = run_id
|
|
if detail is not None:
|
|
event["detail"] = detail
|
|
if res is not None:
|
|
event["res"] = res
|
|
try:
|
|
with open(events_path(output_dir), "a", encoding="utf-8") as fh:
|
|
fh.write(json.dumps(event, ensure_ascii=False) + "\n")
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
def read_events(output_dir, limit=5000):
|
|
"""Lit les événements du journal (dernières `limit` lignes, invalides ignorées)."""
|
|
try:
|
|
lines = events_path(output_dir).read_text(
|
|
encoding="utf-8", errors="replace").splitlines()
|
|
except OSError:
|
|
return []
|
|
events = []
|
|
for line in lines[-limit:]:
|
|
try:
|
|
events.append(json.loads(line))
|
|
except ValueError:
|
|
continue # ligne partielle en cours d'écriture par le run
|
|
return events
|
|
|
|
|
|
def _step_label(phase, detail, res, viz_labels):
|
|
"""Libellé français d'une étape, affiché tel quel dans la file."""
|
|
if phase == "download":
|
|
return "Téléchargement"
|
|
if phase == "classif":
|
|
return "Classification"
|
|
if phase == "dtm":
|
|
return f"DTM {str(res).replace('.', ',')} m"
|
|
if phase == "viz":
|
|
label = (viz_labels or {}).get(detail, detail)
|
|
if res is not None and float(res) != 0.5:
|
|
label += f" {str(res).replace('.', ',')} m"
|
|
return label
|
|
return detail or phase
|
|
|
|
|
|
def aggregate_tiles(events, viz_labels=None):
|
|
"""Regroupe les événements en tuiles : une entrée par tuile, étapes ordonnées.
|
|
|
|
Retourne [{short, name, state, steps: [{key, label, state, detail}]}] dans
|
|
l'ordre de première apparition. État d'une étape : running | ok | fail |
|
|
skip. État d'une tuile : pending | running | done | failed.
|
|
"""
|
|
tiles = {}
|
|
for ev in events:
|
|
if not isinstance(ev, dict) or ev.get("phase") not in PHASES:
|
|
continue
|
|
name = ev.get("tile") or ""
|
|
short = tile_short_name(name)
|
|
tile = tiles.setdefault(
|
|
short, {"short": short, "name": name, "state": "pending", "steps": {}})
|
|
if name:
|
|
tile["name"] = name
|
|
state = ev.get("state")
|
|
if state == "start":
|
|
state = "running" # état d'affichage normalisé
|
|
detail = ev.get("detail")
|
|
res = ev.get("res")
|
|
if ev.get("phase") == "tile":
|
|
if state == "ok":
|
|
tile["state"] = "done"
|
|
elif state == "fail":
|
|
tile["state"] = "failed"
|
|
continue
|
|
# Clé d'étape : le détail n'identifie une étape que pour les viz
|
|
# (leur nom) — les sous-étapes du téléchargement fusionnent en une.
|
|
phase = ev.get("phase")
|
|
key = f"viz:{detail}:{res}" if phase == "viz" else f"{phase}:{res}"
|
|
step = tile["steps"].setdefault(
|
|
key, {"key": key, "detail": detail, "res": res, "state": state,
|
|
"label": _step_label(phase, detail, res, viz_labels)})
|
|
# ok/fail sont terminaux — on ne régresse pas vers running/skip
|
|
if step["state"] not in ("ok", "fail") or state == "fail":
|
|
step["state"] = state
|
|
|
|
result = []
|
|
for tile in tiles.values():
|
|
steps = list(tile["steps"].values())
|
|
if tile["state"] == "pending":
|
|
if any(s["state"] == "fail" for s in steps):
|
|
tile["state"] = "failed"
|
|
elif steps:
|
|
tile["state"] = "running" # pas encore d'événement tuile final
|
|
result.append({"short": tile["short"], "name": tile["name"],
|
|
"state": tile["state"], "steps": steps})
|
|
return result
|
|
|
|
|
|
def progress_snapshot(output_dir, viz_labels=None, run_id=None):
|
|
"""État consolidé des tuiles du dernier run (prêt pour /api/status).
|
|
|
|
run_id (LIDAR_RUN_ID posé par mapserve) : ne garde que les événements de
|
|
CE run — les workers d'un run annulé peuvent continuer d'écrire (O_APPEND)
|
|
après la troncature du journal par le run suivant, et pollueraient sinon
|
|
la file affichée. Sans run_id (run CLI autonome) : tout le journal.
|
|
"""
|
|
events = read_events(output_dir)
|
|
if run_id is not None:
|
|
events = [ev for ev in events if isinstance(ev, dict) and ev.get("run") == run_id]
|
|
return aggregate_tiles(events, viz_labels=viz_labels)
|