diff --git a/lidar_pipeline/cli.py b/lidar_pipeline/cli.py index 8e0855e..7562578 100644 --- a/lidar_pipeline/cli.py +++ b/lidar_pipeline/cli.py @@ -107,6 +107,16 @@ Exemples: default=1, help="Nombre de workers pour traitement parallèle (défaut: 1)" ) + parser.add_argument( + "-g", "--gpu", + nargs="?", + const="all", + default=None, + metavar="LISTE", + help="Sélectionner le(s) GPU(s) à utiliser : un index (ex: -g 0), " + "une liste (ex: -g 0,2), ou 'all' pour tous (ex: -g). " + "Sans -g : GPU auto (le premier si disponible)." + ) parser.add_argument( "-f", "--force", action="store_true", @@ -190,6 +200,21 @@ Exemples: logger.info("Pipeline LiDAR Archéologique") logger.info("=" * 60) + # Parse --gpu into a list of GPU IDs + gpu_ids = None + gpu_arg = args.gpu + if gpu_arg is not None: + if gpu_arg == 'all': + gpu_ids = None # Don't restrict — use all GPUs + else: + try: + gpu_ids = [int(g.strip()) for g in gpu_arg.split(',')] + except ValueError: + parser.error(f"GPU invalide: {gpu_arg!r}. Utilisez un index, une liste (0,2) ou 'all'.") + if gpu_ids is not None: + from .gpu import restrict_gpus + restrict_gpus(gpu_ids) + # Kill orphan PDAL processes on interrupt or termination signal.signal(signal.SIGINT, _kill_orphan_pdal) signal.signal(signal.SIGTERM, _kill_orphan_pdal) @@ -220,6 +245,7 @@ Exemples: only_viz=only_viz, skip_viz=skip_viz, output_format=args.format, + gpu_ids=gpu_ids, ) # If --file is specified, process only matching files diff --git a/lidar_pipeline/gpu.py b/lidar_pipeline/gpu.py index 442f85b..ad9e72c 100644 --- a/lidar_pipeline/gpu.py +++ b/lidar_pipeline/gpu.py @@ -25,6 +25,9 @@ _NUM_GPUS = 0 HAS_GPU = False _gpu_name = None _gpu_mem_gb = 0 +# System-level GPU IDs that are currently visible (after restrict_gpus). +# Populated by restrict_gpus() or auto-detected at import time. +_available_gpu_ids: list[int] = [] try: import subprocess @@ -44,6 +47,8 @@ try: except (ValueError, IndexError): pass HAS_GPU = True + # All GPUs are visible by default + _available_gpu_ids = list(range(_NUM_GPUS)) except (FileNotFoundError, subprocess.TimeoutExpired, Exception): pass @@ -61,8 +66,12 @@ def _init_gpu(): Import CuPy only when needed, so CUDA_VISIBLE_DEVICES can be set before the CUDA context is created. + + Validates that the GPU can actually execute kernels by running + a small computation. This catches CUDA_ERROR_NO_BINARY_FOR_GPU + and other compute capability mismatches before they crash visualizations. """ - global _xp, _cp, _cp_ndimage, _gpu_initialized + global _xp, _cp, _cp_ndimage, _gpu_initialized, HAS_GPU if _gpu_initialized: return _gpu_initialized = True @@ -71,41 +80,80 @@ def _init_gpu(): import cupyx.scipy.ndimage as _real_cupy_ndimage # Verify GPU is actually accessible _real_cupy.cuda.runtime.getDevice() + # Warm-up: run a small computation to verify kernel execution works. + # This catches CUDA_ERROR_NO_BINARY_FOR_GPU (compute capability + # mismatch) and driver errors before we commit to GPU mode. + _test = _real_cupy.array([1.0, 2.0, 3.0], dtype=_real_cupy.float32) + _result = _real_cupy.sum(_test * _test) + # Force execution (CuPy is lazy — .get() ensures the kernel ran) + _ = _result.get() + del _test, _result _xp = _real_cupy _cp = _real_cupy _cp_ndimage = _real_cupy_ndimage except (ImportError, Exception) as e: - logger.debug(f"CuPy non disponible: {e}") + logger.warning(f"GPU non disponible — mode CPU: {e}") _xp = np _cp = None _cp_ndimage = None + HAS_GPU = False + + +def restrict_gpus(gpu_ids: list[int]): + """Restrict which GPUs are visible to the process. + + Sets CUDA_VISIBLE_DEVICES so only the specified system GPU IDs + are accessible. Also updates _available_gpu_ids and _NUM_GPUS. + Must be called before any GPU operation. + + Args: + gpu_ids: List of system-level GPU indices to make visible. + """ + global _NUM_GPUS, HAS_GPU, _available_gpu_ids + if not gpu_ids or not HAS_GPU: + return + + # Validate IDs against total GPU count from nvidia-smi + total_count = _NUM_GPUS or 1 + valid_ids = [gid % total_count for gid in gpu_ids] + _available_gpu_ids = valid_ids + _NUM_GPUS = len(valid_ids) + + os.environ['CUDA_VISIBLE_DEVICES'] = ','.join(str(g) for g in valid_ids) + logger.info(f"GPU visibles: {_available_gpu_ids}") def num_gpus(): - """Return the total number of CUDA GPUs in the system.""" + """Return the number of visible (available) GPUs.""" return _NUM_GPUS def set_active_gpu(gpu_id): """Set the active GPU for the current process via CUDA_VISIBLE_DEVICES. + gpu_id is an index into the currently visible GPU list + (_available_gpu_ids), not a system-level ID. + MUST be called before any GPU operation (to_gpu, etc.) to ensure CuPy creates its CUDA context on the correct device. With lazy initialization, CuPy is imported AFTER this call, so it only sees the assigned GPU. Args: - gpu_id: 0-based GPU index (referring to the system GPU numbering). + gpu_id: 0-based index into the visible GPU list. """ if not HAS_GPU or _NUM_GPUS <= 1: return # Nothing to do for single GPU or no GPU gpu_id = gpu_id % _NUM_GPUS - # Set CUDA_VISIBLE_DEVICES before CuPy context creation - os.environ['CUDA_VISIBLE_DEVICES'] = str(gpu_id) + # Map visible-GPU index back to the real system GPU ID + system_gpu_id = _available_gpu_ids[gpu_id] - logger.info(f" GPU {gpu_id} sélectionnée pour ce worker") + # Set CUDA_VISIBLE_DEVICES before CuPy context creation + os.environ['CUDA_VISIBLE_DEVICES'] = str(system_gpu_id) + + logger.info(f" GPU {system_gpu_id} sélectionnée pour ce worker") def _gpu_available(): @@ -210,4 +258,48 @@ def gpu_cleanup(): try: _cp.get_default_memory_pool().free_all_blocks() except Exception: - pass \ No newline at end of file + pass + + +def disable_gpu(): + """Disable GPU acceleration for the rest of this process. + + Called when a CUDA error indicates the GPU is unusable (e.g. + CUDA_ERROR_NO_BINARY_FOR_GPU). Falls back to numpy for all + subsequent operations. + """ + global HAS_GPU, _xp, _cp, _cp_ndimage + if not HAS_GPU: + return # Already disabled + logger.warning("GPU désactivé — passage en mode CPU pour la suite du processus") + HAS_GPU = False + _xp = np + _cp = None + _cp_ndimage = None + + +def is_gpu_active(): + """Check if GPU acceleration is currently active. + + Unlike the HAS_GPU module-level variable (which can go stale if + imported directly), this always reflects the current runtime state. + Use this in logging tags and conditional GPU paths. + """ + return HAS_GPU + + +def safe_gpu_call(func, *args, **kwargs): + """Call a function with GPU arrays, retrying on CPU if GPU fails. + + Usage: + result = safe_gpu_call(generate_svf, dem_file, basename, vis_dir, resolution, shared=shared) + """ + try: + return func(*args, **kwargs) + except Exception as e: + err_msg = str(e) + if _cp is not None and ('CUDA' in err_msg or 'cuda' in err_msg or 'GPU' in err_msg): + logger.warning(f"Erreur GPU ({e.__class__.__name__}), retry en CPU...") + disable_gpu() + return func(*args, **kwargs) + raise \ No newline at end of file diff --git a/lidar_pipeline/pipeline.py b/lidar_pipeline/pipeline.py index d4660cc..8e46497 100644 --- a/lidar_pipeline/pipeline.py +++ b/lidar_pipeline/pipeline.py @@ -63,7 +63,7 @@ from .visualizations import ( generate_roughness, generate_wavelet, generate_svf, generate_aniso_open, generate_paths, ) -from .gpu import gpu_cleanup, num_gpus, safe_gpu_call +from .gpu import gpu_cleanup, num_gpus, restrict_gpus, safe_gpu_call from .ign import generate_ign_overlay from .rendering import tif_to_png @@ -107,7 +107,7 @@ VIZ_STEPS = [ class LidarArchaeoPipeline: """Orchestrates the LiDAR archaeological analysis pipeline.""" - def __init__(self, input_dir, output_dir, resolution=0.5, workers=1, force=False, ground_method='auto', force_classify=False, keep_tif=False, quality=98, only_viz=None, skip_viz=None, output_format='avif'): + def __init__(self, input_dir, output_dir, resolution=0.5, workers=1, force=False, ground_method='auto', force_classify=False, keep_tif=False, quality=98, only_viz=None, skip_viz=None, output_format='avif', gpu_ids=None): self.input_dir = Path(input_dir) self.output_dir = Path(output_dir) # Accept single float or comma-separated string for multi-resolution @@ -127,6 +127,7 @@ class LidarArchaeoPipeline: self.only_viz = only_viz self.skip_viz = skip_viz self.output_format = output_format + self.gpu_ids = gpu_ids self.temp_dir = self.output_dir / "temp" if not self.input_dir.exists(): @@ -448,6 +449,10 @@ class LidarArchaeoPipeline: results = {} t_pipeline_start = time.time() + # Restrict visible GPUs in the main process before spawning workers + if self.gpu_ids is not None: + restrict_gpus(self.gpu_ids) + if self.workers > 1 and len(files) > 1: n_gpus = num_gpus() or 1 if n_gpus > 1: @@ -462,7 +467,7 @@ class LidarArchaeoPipeline: # Pass resolutions as comma-separated string for multiprocessing serialization resolutions_str = ','.join(str(r) for r in self.resolutions) future_to_file = { - executor.submit(_process_file_standalone, str(laz_file), str(self.input_dir), str(self.output_dir), resolutions_str, self.force, self.ground_method, self.force_classify, self.keep_tif, self.quality, self.only_viz, self.skip_viz, self.output_format, gpu_id % n_gpus): laz_file + executor.submit(_process_file_standalone, str(laz_file), str(self.input_dir), str(self.output_dir), resolutions_str, self.force, self.ground_method, self.force_classify, self.keep_tif, self.quality, self.only_viz, self.skip_viz, self.output_format, gpu_id % n_gpus, gpu_ids=self.gpu_ids): laz_file for gpu_id, laz_file in enumerate(files) } done = 0 @@ -531,14 +536,18 @@ class LidarArchaeoPipeline: logger.warning(f" Note: Impossible de supprimer les fichiers temporaires: {e}") -def _process_file_standalone(laz_file_str, input_dir, output_dir, resolution, force=False, ground_method='auto', force_classify=False, keep_tif=False, quality=98, only_viz=None, skip_viz=None, output_format='avif', gpu_id=None): +def _process_file_standalone(laz_file_str, input_dir, output_dir, resolution, force=False, ground_method='auto', force_classify=False, keep_tif=False, quality=98, only_viz=None, skip_viz=None, output_format='avif', gpu_id=None, gpu_ids=None): """Standalone function for multiprocessing — creates its own pipeline instance. Each worker gets its own temp directory to avoid file conflicts. When multiple GPUs are available, each worker is assigned a GPU via - CuPy's Device API to balance load across GPUs. + CUDA_VISIBLE_DEVICES to balance load across GPUs. """ - # Assign GPU to this worker using CuPy's Device API + # Restrict visible GPUs first, then pick one for this worker + if gpu_ids is not None: + from .gpu import restrict_gpus + restrict_gpus(gpu_ids) + if gpu_id is not None and gpu_id >= 0: from .gpu import set_active_gpu set_active_gpu(gpu_id)