"""Structured progress of a generation run (pipeline → interface contract). The pipeline and the IGN download write JSON events, one per line, to /.generation.events.jsonl; mapserve aggregates them through /api/status to show tiles one by one (name, type, state) in the generation queue, instead of the raw log. This module is deliberately decoupled from the server: tile rendering can be moved out of the web server — the event producer depends only on the output directory. Parallel workers each write their own line in append mode (a line < 4 KB written with O_APPEND is atomic on POSIX). """ import json import os import re import time from pathlib import Path EVENTS_FILENAME = ".generation.events.jsonl" # Phases of a tile, in execution order. PHASES = ("download", "classif", "dtm", "viz", "tile") _TILE_RE = re.compile(r"LHD_FXX_(\d+)_(\d+)_") def events_path(output_dir): """Path of a run's event log.""" return Path(output_dir) / EVENTS_FILENAME def tile_short_name(name): """Short name of a tile: 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): """Prepare an empty log for a new run (truncates the previous one).""" try: path = events_path(output_dir) path.parent.mkdir(parents=True, exist_ok=True) path.write_text("", encoding="utf-8") except OSError: pass # best-effort progress: never block processing def report_event(output_dir, tile, phase, state, detail=None, res=None): """Append an event (one JSON line). Never raises. The file is opened in append mode on every call: parallel workers (ProcessPoolExecutor, 'spawn') share no file descriptor. """ 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): """Read the log events (last `limit` lines, invalid ones skipped).""" 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 # partial line still being written by the run return events def _step_label(phase, detail, res, viz_labels): """Label of a step, displayed as is in the queue.""" if phase == "download": return "Download" if phase == "classif": return "Classification" if phase == "dtm": return f"DTM {res} m" if phase == "viz": label = (viz_labels or {}).get(detail, detail) if res is not None and float(res) != 0.5: label += f" {res} m" return label return detail or phase def aggregate_tiles(events, viz_labels=None): """Group events into tiles: one entry per tile, ordered steps. Returns [{short, name, state, steps: [{key, detail, res, state, label}]}] in order of first appearance. Step state: running | ok | fail | skip. Tile state: 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" # normalized display state 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 # Step key: the detail identifies a step only for viz (their name) # — download sub-steps merge into a single one. 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 are terminal — never regress to running/skip (fail still overrides ok) 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" # no final tile event yet 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): """Consolidated tile state of the last run (ready for /api/status). run_id (LIDAR_RUN_ID set by mapserve): keeps only the events of THIS run — workers of a cancelled run may keep writing (O_APPEND) after the next run truncated the log, and would otherwise pollute the displayed queue. Without run_id (standalone CLI run): the whole log. """ 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)