Split webapp for Raspberry Pi deployment, remote generation API and sync
La webapp (carte + vignettes) et la génération de tuiles se déploient sur deux machines : image légère Dockerfile.webapp (FastAPI + Pillow AVIF natif + pyproj) sur Raspberry Pi, pipeline complet sur la machine de traitement. LIDAR_GENERATION_URL délègue /api/generate, /api/preview et /api/status ; /api/sync ramène les tuiles par rsync puis régénère vignettes et index localement. Token partagé optionnel (LIDAR_API_TOKEN/LIDAR_REMOTE_TOKEN). Retire du dépôt les journaux internes (.swival, audit-findings) et les données (data/, notebooks/). Doc : docs/DEPLOY_WEBAPP.md.
This commit is contained in:
154
lidar_pipeline/progress.py
Normal file
154
lidar_pipeline/progress.py
Normal file
@ -0,0 +1,154 @@
|
||||
"""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 ; la webapp 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é de la webapp : 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 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}
|
||||
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):
|
||||
"""État consolidé des tuiles du dernier run (prêt pour /api/status)."""
|
||||
return aggregate_tiles(read_events(output_dir), viz_labels=viz_labels)
|
||||
Reference in New Issue
Block a user