From d0dc8d90e9b3ec45515ebe2a2d85a6d6cae09ccf Mon Sep 17 00:00:00 2001 From: Jacquin Antoine Date: Mon, 28 Sep 2026 00:03:52 +0200 Subject: [PATCH] Make the worker batch timeout configurable via LIDAR_BATCH_TIMEOUT MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pool of tile workers used to cancel every remaining tile after a hardcoded 2-hour wall clock, silently truncating large batches (a 670-tile completion run lost its last 348 tiles that way). The timeout now defaults to unlimited and can be capped per deployment with the LIDAR_BATCH_TIMEOUT environment variable (seconds); the local worker compose sets it to 6 hours. 💘 Generated with Crush Assisted-by: Crush:glm-5.2 --- AGENTS.md | 2 +- docker-compose.worker.yml | 2 ++ lidar_pipeline/pipeline.py | 35 +++++++++++++++++++++------ lidar_pipeline/tests/test_pipeline.py | 18 ++++++++++++++ 4 files changed, 48 insertions(+), 9 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index f4959ea..750b997 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -84,7 +84,7 @@ - `_d8_accumulate_numba` uses numba with `argsort` top-down sweep. Python fallback exists. - Ray-tracing (SVF, openness): processes one direction at a time to limit VRAM. Auto-falls back to CPU on OOM via `_ray_trace_horizons`. - Multi-resolution: 0.5 m has no filename suffix whatever its position; every other resolution (including the 0.2 m default) uses a `_r0p2` style suffix. Ground classification done once, shared across resolutions. -- `ProcessPoolExecutor` has a 2-hour wall-clock safety timeout (prevents indefinite hang from stuck workers). +- `ProcessPoolExecutor` batch runs are unlimited by default; set `LIDAR_BATCH_TIMEOUT` (seconds) to cap a batch's wall clock and cancel the remaining workers past it. ### Numba usage pattern - Defined at function scope with `@njit(cache=True)` — first call compiles (~2-3s), subsequent calls hit disk cache. diff --git a/docker-compose.worker.yml b/docker-compose.worker.yml index df03f7d..2fc98d9 100644 --- a/docker-compose.worker.yml +++ b/docker-compose.worker.yml @@ -38,6 +38,8 @@ services: # Generations started from a remote map use the GPU - LIDAR_GPU=1 - LIDAR_WORKERS=auto + # Wall-clock cap for one batch of tiles, in seconds (unset or 0 = unlimited) + - LIDAR_BATCH_TIMEOUT=21600 # Workers per GPU capped by the free VRAM when the run starts # ((free - reserve) / per-worker peak); the excess runs on the CPU. # Peak estimated at 2048 MiB: adjust after measuring (nvidia-smi during a run). diff --git a/lidar_pipeline/pipeline.py b/lidar_pipeline/pipeline.py index aca9502..0968669 100644 --- a/lidar_pipeline/pipeline.py +++ b/lidar_pipeline/pipeline.py @@ -47,6 +47,22 @@ def resolve_workers(value): return 1 +def _batch_timeout_s(): + """Wall-clock budget for one batch, from LIDAR_BATCH_TIMEOUT (seconds). + + Unset, empty or 0 means unlimited (default). A negative or non-numeric + value is ignored with a warning. + """ + raw = (os.environ.get("LIDAR_BATCH_TIMEOUT") or "").strip() + if not raw: + return 0.0 + try: + return max(0.0, float(raw)) + except ValueError: + logger.warning(f"LIDAR_BATCH_TIMEOUT={raw!r} is not a number — ignored (unlimited)") + return 0.0 + + def _file_basename(path): """Extract base name from a LAZ/LAS file, removing all known extensions. @@ -852,15 +868,17 @@ class LidarArchaeoPipeline: for laz_file in files } done = 0 - t_deadline = time.time() + 7200 + budget_s = _batch_timeout_s() + deadline = time.time() + budget_s if budget_s > 0 else None try: - # timeout=: without it, as_completed blocks between two - # completions and the 2 h deadline is never evaluated if - # no worker returns (run stuck forever). + # timeout= lets the deadline be checked between two + # completions; without it a stuck worker would block + # as_completed forever. No deadline (LIDAR_BATCH_TIMEOUT + # unset or 0) → wait indefinitely. + as_completed_kwargs = ({"timeout": max(1.0, deadline - time.time())} + if deadline is not None else {}) try: - futures_iter = as_completed(future_to_file, - timeout=max(1.0, t_deadline - time.time())) - for future in futures_iter: + for future in as_completed(future_to_file, **as_completed_kwargs): laz_file = future_to_file[future] done += 1 try: @@ -877,7 +895,8 @@ class LidarArchaeoPipeline: "tile", "fail", detail=str(e)) results[laz_file.name] = False except FuturesTimeoutError: - logger.error("Timeout exceeded (2 h) — cancelling remaining workers") + logger.error(f"Batch exceeded LIDAR_BATCH_TIMEOUT " + f"({budget_s:g} s) — cancelling remaining workers") for f in future_to_file: f.cancel() except KeyboardInterrupt: diff --git a/lidar_pipeline/tests/test_pipeline.py b/lidar_pipeline/tests/test_pipeline.py index 0de739e..9501657 100644 --- a/lidar_pipeline/tests/test_pipeline.py +++ b/lidar_pipeline/tests/test_pipeline.py @@ -440,3 +440,21 @@ class TestResolveWorkers: for _ in range(4): # 4th call: empty queue pipeline._init_worker_slot(q) assert calls == [("gpu", 1), ("cpu",)] + + +class TestBatchTimeout: + def test_env_var_parsing(self, monkeypatch): + """LIDAR_BATCH_TIMEOUT: unset/0/invalid = unlimited, N seconds = N.""" + from lidar_pipeline.pipeline import _batch_timeout_s + monkeypatch.delenv("LIDAR_BATCH_TIMEOUT", raising=False) + assert _batch_timeout_s() == 0.0 + monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "") + assert _batch_timeout_s() == 0.0 + monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "0") + assert _batch_timeout_s() == 0.0 + monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "3600") + assert _batch_timeout_s() == 3600.0 + monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "-5") + assert _batch_timeout_s() == 0.0 + monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "abc") + assert _batch_timeout_s() == 0.0