diff --git a/lidar_pipeline/progress.py b/lidar_pipeline/progress.py index b8c3e8e..d5d97e5 100644 --- a/lidar_pipeline/progress.py +++ b/lidar_pipeline/progress.py @@ -12,6 +12,7 @@ ajout (une ligne < 4 Ko avec O_APPEND est atomique sur POSIX). """ import json +import os import re import time from pathlib import Path @@ -54,6 +55,9 @@ def report_event(output_dir, tile, phase, state, detail=None, res=None): (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: @@ -149,6 +153,15 @@ def aggregate_tiles(events, viz_labels=None): 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) +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 la webapp) : 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) diff --git a/lidar_pipeline/tests/test_webapp.py b/lidar_pipeline/tests/test_webapp.py index 1cec64e..e4a0c19 100644 --- a/lidar_pipeline/tests/test_webapp.py +++ b/lidar_pipeline/tests/test_webapp.py @@ -831,18 +831,31 @@ def test_rebuild_index_background(tmp_path, monkeypatch): def test_status_exposes_tiles_from_events(tmp_path, monkeypatch): - """/api/status agrège les événements en tuiles (nom court + étapes).""" + """/api/status agrège les événements en tuiles (nom court + étapes). + + Les événements sont filtrés par identifiant de run : ceux d'un run + annulé (workers survivants écrivant après la troncature) n'apparaissent + pas dans la file du run courant.""" import lidar_pipeline.webapp as webapp from lidar_pipeline.progress import report_event, reset_events monkeypatch.setattr(webapp, "OUTPUT_DIR", tmp_path) - reset_events(tmp_path) - report_event(tmp_path, "LHD_FXX_1054_6882_PTS_LAMB93_IGN69", "viz", - "start", "aspect", res=0.5) - s = webapp.status() + with monkeypatch.context() as m: + m.setitem(webapp._job, "run_id", "run-courant") + reset_events(tmp_path) + m.setenv("LIDAR_RUN_ID", "run-courant") + report_event(tmp_path, "LHD_FXX_1054_6882_PTS_LAMB93_IGN69", "viz", + "start", "aspect", res=0.5) + # Worker retardataire d'un run annulé : réappend après le reset + m.setenv("LIDAR_RUN_ID", "run-annule") + report_event(tmp_path, "LHD_FXX_1055_6882_PTS_LAMB93_IGN69", "viz", + "start", "aspect", res=0.5) + m.delenv("LIDAR_RUN_ID") + s = webapp.status() match = [t for t in s["tiles"] if t["short"] == "1054-6882"] assert len(match) == 1 assert match[0]["state"] == "running" assert any(step["label"] == "Aspect" for step in match[0]["steps"]) + assert all(t["short"] != "1055-6882" for t in s["tiles"]) def test_generate_resets_progress_events(tmp_path, monkeypatch): diff --git a/lidar_pipeline/webapp.py b/lidar_pipeline/webapp.py index c9fc753..2f28bdd 100644 --- a/lidar_pipeline/webapp.py +++ b/lidar_pipeline/webapp.py @@ -72,6 +72,7 @@ import threading import time import urllib.parse import urllib.request +import uuid from pathlib import Path from typing import Optional @@ -537,7 +538,7 @@ class ExportRequest(BaseModel): # --- État du job de génération ------------------------------------------- -_job = {"proc": None, "started": None, "returncode": None, "cmd": None, "finished": None, "qid": None} +_job: dict = {"proc": None, "started": None, "returncode": None, "cmd": None, "finished": None, "qid": None, "run_id": None} _job_lock = threading.Lock() @@ -548,7 +549,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", "qid")} + state = {k: _job.get(k) for k in ("started", "returncode", "cmd", "finished", "qid", "run_id")} try: JOB_FILE.write_text(json.dumps(state), encoding="utf-8") except OSError: @@ -567,6 +568,7 @@ def _load_job_state(): _job["cmd"] = state.get("cmd") _job["finished"] = state.get("finished") _job["qid"] = state.get("qid") + _job["run_id"] = state.get("run_id") _load_job_state() @@ -881,7 +883,8 @@ def status(request: Request = None): "regen_allowed": regen_allowed, # Tuiles une par une (nom court, type, état) — agrégées depuis les # événements écrits par le pipeline (.generation.events.jsonl) - "tiles": progress_snapshot(OUTPUT_DIR, viz_labels=_viz_step_labels()), + "tiles": progress_snapshot(OUTPUT_DIR, viz_labels=_viz_step_labels(), + run_id=_job.get("run_id")), } @@ -1272,11 +1275,14 @@ def _launch_job(tiles, viz, req, qid=None): 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) + # dans la file correspondent au run qui démarre, pas au précédent) et + # identifiant de run transmis au pipeline (LIDAR_RUN_ID) : les événements + # tardifs d'un run annulé (workers survivants) sont filtrés à la lecture. from .progress import reset_events reset_events(OUTPUT_DIR) + run_id = uuid.uuid4().hex[:12] _job.update({"proc": None, "started": time.time(), "returncode": None, - "cmd": cmd, "finished": None, "qid": qid}) + "cmd": cmd, "finished": None, "qid": qid, "run_id": run_id}) _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 @@ -1284,6 +1290,7 @@ def _launch_job(tiles, viz, req, qid=None): # reste confiné à son groupe. p = subprocess.Popen(cmd, stdout=log_fh, stderr=subprocess.STDOUT, cwd="/app" if Path("/app").exists() else None, + env=dict(os.environ, LIDAR_RUN_ID=run_id), start_new_session=True) def _watch():