1 Commits

Author SHA1 Message Date
d0dc8d90e9 Make the worker batch timeout configurable via LIDAR_BATCH_TIMEOUT
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
2026-09-28 00:03:52 +02:00
4 changed files with 48 additions and 9 deletions

View File

@ -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.

View File

@ -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).

View File

@ -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:

View File

@ -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