"""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 /.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)