diff --git a/lidar_pipeline/pipeline.py b/lidar_pipeline/pipeline.py index 533be08..510deb5 100644 --- a/lidar_pipeline/pipeline.py +++ b/lidar_pipeline/pipeline.py @@ -11,7 +11,7 @@ import logging import multiprocessing import shutil import time -from concurrent.futures import ProcessPoolExecutor, as_completed +from concurrent.futures import ProcessPoolExecutor, as_completed, TimeoutError as FuturesTimeoutError from datetime import datetime from pathlib import Path import subprocess @@ -589,27 +589,32 @@ class LidarArchaeoPipeline: done = 0 t_deadline = time.time() + 7200 try: - for future in as_completed(future_to_file): - if time.time() > t_deadline: - logger.error("Délai dépassé (2h) — annulation des workers restants") - for f in future_to_file: - f.cancel() - break - laz_file = future_to_file[future] - done += 1 - try: - success = future.result() - results[laz_file.name] = success - status = "✓" if success else "✗" - logger.info(f" [{done}/{len(files)}] {status} {laz_file.name}") - if success and self.incremental_index: - self._rebuild_index_incremental() - except Exception as e: - logger.error(f" [{done}/{len(files)}] ✗ {laz_file.name}: {e}") - logger.debug(f" Traceback:", exc_info=True) - report_event(self.output_dir, _file_basename(laz_file), - "tile", "fail", detail=str(e)) - results[laz_file.name] = False + # timeout= : sans lui, as_completed bloque entre deux + # complétions et le délai de 2 h n'est jamais évalué si + # aucun worker ne rend la main (run figé pour toujours). + try: + futures_iter = as_completed(future_to_file, + timeout=max(1.0, t_deadline - time.time())) + for future in futures_iter: + laz_file = future_to_file[future] + done += 1 + try: + success = future.result() + results[laz_file.name] = success + status = "✓" if success else "✗" + logger.info(f" [{done}/{len(files)}] {status} {laz_file.name}") + if success and self.incremental_index: + self._rebuild_index_incremental() + except Exception as e: + logger.error(f" [{done}/{len(files)}] ✗ {laz_file.name}: {e}") + logger.debug(f" Traceback:", exc_info=True) + report_event(self.output_dir, _file_basename(laz_file), + "tile", "fail", detail=str(e)) + results[laz_file.name] = False + except FuturesTimeoutError: + logger.error("Délai dépassé (2h) — annulation des workers restants") + for f in future_to_file: + f.cancel() except KeyboardInterrupt: logger.info("Interruption — annulation des travaux en cours...") for f in future_to_file: