Étiqueter les événements de progression par run pour filtrer les retardataires
This commit is contained in:
@ -12,6 +12,7 @@ ajout (une ligne < 4 Ko avec O_APPEND est atomique sur POSIX).
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
import os
|
||||||
import re
|
import re
|
||||||
import time
|
import time
|
||||||
from pathlib import Path
|
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.
|
(ProcessPoolExecutor, 'spawn') ne partagent aucun descripteur de fichier.
|
||||||
"""
|
"""
|
||||||
event = {"ts": round(time.time(), 3), "tile": tile, "phase": phase, "state": state}
|
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:
|
if detail is not None:
|
||||||
event["detail"] = detail
|
event["detail"] = detail
|
||||||
if res is not None:
|
if res is not None:
|
||||||
@ -149,6 +153,15 @@ def aggregate_tiles(events, viz_labels=None):
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
||||||
def progress_snapshot(output_dir, viz_labels=None):
|
def progress_snapshot(output_dir, viz_labels=None, run_id=None):
|
||||||
"""État consolidé des tuiles du dernier run (prêt pour /api/status)."""
|
"""État consolidé des tuiles du dernier run (prêt pour /api/status).
|
||||||
return aggregate_tiles(read_events(output_dir), viz_labels=viz_labels)
|
|
||||||
|
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)
|
||||||
|
|||||||
@ -831,18 +831,31 @@ def test_rebuild_index_background(tmp_path, monkeypatch):
|
|||||||
|
|
||||||
|
|
||||||
def test_status_exposes_tiles_from_events(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
|
import lidar_pipeline.webapp as webapp
|
||||||
from lidar_pipeline.progress import report_event, reset_events
|
from lidar_pipeline.progress import report_event, reset_events
|
||||||
monkeypatch.setattr(webapp, "OUTPUT_DIR", tmp_path)
|
monkeypatch.setattr(webapp, "OUTPUT_DIR", tmp_path)
|
||||||
reset_events(tmp_path)
|
with monkeypatch.context() as m:
|
||||||
report_event(tmp_path, "LHD_FXX_1054_6882_PTS_LAMB93_IGN69", "viz",
|
m.setitem(webapp._job, "run_id", "run-courant")
|
||||||
"start", "aspect", res=0.5)
|
reset_events(tmp_path)
|
||||||
s = webapp.status()
|
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"]
|
match = [t for t in s["tiles"] if t["short"] == "1054-6882"]
|
||||||
assert len(match) == 1
|
assert len(match) == 1
|
||||||
assert match[0]["state"] == "running"
|
assert match[0]["state"] == "running"
|
||||||
assert any(step["label"] == "Aspect" for step in match[0]["steps"])
|
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):
|
def test_generate_resets_progress_events(tmp_path, monkeypatch):
|
||||||
|
|||||||
@ -72,6 +72,7 @@ import threading
|
|||||||
import time
|
import time
|
||||||
import urllib.parse
|
import urllib.parse
|
||||||
import urllib.request
|
import urllib.request
|
||||||
|
import uuid
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
@ -537,7 +538,7 @@ class ExportRequest(BaseModel):
|
|||||||
|
|
||||||
|
|
||||||
# --- État du job de génération -------------------------------------------
|
# --- É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()
|
_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.
|
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).
|
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:
|
try:
|
||||||
JOB_FILE.write_text(json.dumps(state), encoding="utf-8")
|
JOB_FILE.write_text(json.dumps(state), encoding="utf-8")
|
||||||
except OSError:
|
except OSError:
|
||||||
@ -567,6 +568,7 @@ def _load_job_state():
|
|||||||
_job["cmd"] = state.get("cmd")
|
_job["cmd"] = state.get("cmd")
|
||||||
_job["finished"] = state.get("finished")
|
_job["finished"] = state.get("finished")
|
||||||
_job["qid"] = state.get("qid")
|
_job["qid"] = state.get("qid")
|
||||||
|
_job["run_id"] = state.get("run_id")
|
||||||
|
|
||||||
|
|
||||||
_load_job_state()
|
_load_job_state()
|
||||||
@ -881,7 +883,8 @@ def status(request: Request = None):
|
|||||||
"regen_allowed": regen_allowed,
|
"regen_allowed": regen_allowed,
|
||||||
# Tuiles une par une (nom court, type, état) — agrégées depuis les
|
# Tuiles une par une (nom court, type, état) — agrégées depuis les
|
||||||
# événements écrits par le pipeline (.generation.events.jsonl)
|
# é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_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||||
log_fh = open(LOG_FILE, "w", encoding="utf-8")
|
log_fh = open(LOG_FILE, "w", encoding="utf-8")
|
||||||
# Nouveau run : journal d'événements remis à zéro (les tuiles affichées
|
# 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
|
from .progress import reset_events
|
||||||
reset_events(OUTPUT_DIR)
|
reset_events(OUTPUT_DIR)
|
||||||
|
run_id = uuid.uuid4().hex[:12]
|
||||||
_job.update({"proc": None, "started": time.time(), "returncode": None,
|
_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()
|
_save_job_state()
|
||||||
# start_new_session : le pipeline et ses workers/PDAL forment leur
|
# start_new_session : le pipeline et ses workers/PDAL forment leur
|
||||||
# propre groupe de processus — /api/stop peut le tuer en bloc sans
|
# 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.
|
# reste confiné à son groupe.
|
||||||
p = subprocess.Popen(cmd, stdout=log_fh, stderr=subprocess.STDOUT,
|
p = subprocess.Popen(cmd, stdout=log_fh, stderr=subprocess.STDOUT,
|
||||||
cwd="/app" if Path("/app").exists() else None,
|
cwd="/app" if Path("/app").exists() else None,
|
||||||
|
env=dict(os.environ, LIDAR_RUN_ID=run_id),
|
||||||
start_new_session=True)
|
start_new_session=True)
|
||||||
|
|
||||||
def _watch():
|
def _watch():
|
||||||
|
|||||||
Reference in New Issue
Block a user