"""Pipeline orchestration for LiDAR archaeological analysis. LidarArchaeoPipeline coordinates the full processing chain: 1. Ground classification (PDAL/SMRF) 2. DTM generation 3. Visualization generation (17 products) 4. Rendering (AVIF/WebP conversion) """ import logging import multiprocessing import os import shutil import threading import time from concurrent.futures import ProcessPoolExecutor, as_completed, TimeoutError as FuturesTimeoutError from datetime import datetime from pathlib import Path import subprocess # Use 'spawn' to avoid CUDA context corruption in forked subprocesses try: multiprocessing.set_start_method('spawn') except RuntimeError: pass # Already set (e.g. in tests or when called multiple times) logger = logging.getLogger("lidar") def resolve_workers(value): """Traduit l'option -w en nombre effectif de workers. Résolu au lancement de chaque run (pas au démarrage du serveur, ni par tuile : le pool de processus vit le temps du run). « auto » = cœurs logiques - 2 (un pour l'OS/serveur, un pour les phases I/O et l'indexation), borné [2, 16] — chaque worker traite une tuile et peut lancer un processus PDAL en flux, le compte reste raisonnable même sur une grosse machine. """ if isinstance(value, str) and value.strip().lower() == "auto": cpus = os.cpu_count() or 4 return max(2, min(cpus - 2, 16)) try: return max(1, int(value)) except (TypeError, ValueError): return 1 def _file_basename(path): """Extract base name from a LAZ/LAS file, removing all known extensions. Handles double extensions like .copc.laz correctly: 'file.copc.laz' -> 'file', not 'file.copc' """ name = Path(path).name # Remove known LiDAR extensions (order matters: longest first) for ext in ['.copc.laz', '.copc.las', '.laz', '.las']: if name.lower().endswith(ext): return name[:-len(ext)] return Path(path).stem class FilePrefixFilter(logging.Filter): """Adds a file prefix to log messages when processing a specific file.""" def __init__(self): super().__init__() self.basename = None def filter(self, record): if self.basename: record.msg = f"[{self.basename}] {record.msg}" return True # Module-level filter instance so process_file can set it _file_filter = FilePrefixFilter() from .progress import report_event from .dtm import (classify_ground, create_dtm_fast, STRIP_ALIGN_VERSION, STRIP_ALIGN_THRESHOLD, STRIP_JITTER_BIN, STRIP_JITTER_SMOOTH, STRIP_LINE_WINDOW, STRIP_LINE_CELL, STRIP_LINE_MODEL) from .visualizations import ( SharedDEM, generate_hillshade, generate_slope, generate_aspect, generate_openness, generate_mslrm, generate_sailore, generate_roughness, generate_wavelet, generate_solar, generate_svf, generate_flow_accumulation, generate_anomaly_mask, generate_relief_oriente, generate_densite_sol, ) from .gpu import gpu_cleanup, num_gpus, available_gpu_ids, restrict_gpus, safe_gpu_call, gpu_worker_slots from .ign import generate_ign_overlay from .rendering import tif_to_crop, LOSSLESS_GRAY_KEYWORDS # Ordered list of visualization steps. # Each entry: (name, function_or_lambda) # Adding a new visualization = add a generate_* function + register here. VIZ_STEPS = [ ('hillshade', generate_hillshade), ('slope', generate_slope), ('aspect', generate_aspect), ('mslrm', generate_mslrm), ('sailore', generate_sailore), ('pos_open', lambda d, b, v, r, shared=None: generate_openness(d, b, v, r, positive=True, shared=shared)), ('neg_open', lambda d, b, v, r, shared=None: generate_openness(d, b, v, r, positive=False, shared=shared)), ('svf', generate_svf), ('roughness', generate_roughness), ('wavelet', generate_wavelet), ('flow_acc', generate_flow_accumulation), ('solar', generate_solar), ('anomaly', generate_anomaly_mask), ('relief_oriente', generate_relief_oriente), ('densite_sol', generate_densite_sol), ('ortho', lambda d, b, v, r: generate_ign_overlay( d, b, v, r, layer='ORTHOIMAGERY.ORTHOPHOTOS', title='Photographie Aérienne IGN', legend_label='Orthophotographie\nImage aérienne', description='Photographie aérienne IGN (Orthophoto)', out_suffix='ortho')), ('topo', lambda d, b, v, r: generate_ign_overlay( d, b, v, r, layer='GEOGRAPHICALGRIDSYSTEMS.PLANIGNV2', title='Carte Topographique IGN', legend_label='Carte IGN\nPlan topographique', description='Carte topographique IGN (Plan IGN)', out_suffix='topo')), ] 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', ign_classes="sol", force_classify=False, keep_tif=False, quality=60, only_viz=None, skip_viz=None, output_format='avif', gpu_ids=None, no_index=False, incremental_index=False, strip_align=True, openness_downsample=None, edge_buffer=0.0): self.input_dir = Path(input_dir) self.output_dir = Path(output_dir) # Accept single float or comma-separated string for multi-resolution if isinstance(resolution, str): self.resolutions = [float(r.strip()) for r in resolution.split(',')] elif isinstance(resolution, (list, tuple)): self.resolutions = [float(r) for r in resolution] else: self.resolutions = [float(resolution)] self.resolution = self.resolutions[0] # Primary resolution (backward compat) self.workers = workers self.force = force self.ground_method = ground_method self.ign_classes = ign_classes self.force_classify = force_classify self.keep_tif = keep_tif self.quality = quality self.only_viz = only_viz self.skip_viz = skip_viz self.output_format = output_format self.gpu_ids = gpu_ids self.no_index = no_index # Inventaire mis à jour après chaque dalle par défaut : la carte (et sa # maintenance de pyramide) suit le rendu en cours, quel que soit le # lanceur (carte, compose, run.sh). --no-index le désactive. self.incremental_index = bool(incremental_index) or not no_index self.strip_align = strip_align self.openness_downsample = openness_downsample self.edge_buffer = float(edge_buffer) self._last_index_rebuild = 0.0 self._index_timer = None self._index_lock = threading.Lock() self.temp_dir = self.output_dir / "temp" if not self.input_dir.exists(): raise ValueError(f"Répertoire introuvable: {self.input_dir}") self.output_dir.mkdir(parents=True, exist_ok=True) self.temp_dir.mkdir(exist_ok=True) self.dtm_dir = self.output_dir / "DTM" self.vis_dir = self.output_dir / "visualisations" for d in [self.dtm_dir, self.vis_dir]: d.mkdir(exist_ok=True) # Filter visualizations based on --only / --skip all_viz_names = [name for name, _ in VIZ_STEPS] if only_viz: invalid = set(only_viz) - set(all_viz_names) if invalid: raise ValueError(f"Visualisations inconnues: {', '.join(invalid)}. Disponibles: {', '.join(all_viz_names)}") self.viz_steps = [(n, f) for n, f in VIZ_STEPS if n in only_viz] elif skip_viz: invalid = set(skip_viz) - set(all_viz_names) if invalid: raise ValueError(f"Visualisations inconnues: {', '.join(invalid)}. Disponibles: {', '.join(all_viz_names)}") self.viz_steps = [(n, f) for n, f in VIZ_STEPS if n not in skip_viz] else: # Sans --only/--skip : seules les couches affichées par la carte # (PANEL_VIZ, index.py) sont produites. from .index import panel_steps default = panel_steps() self.viz_steps = [(n, f) for n, f in VIZ_STEPS if default is None or n in default] logger.info("Pipeline initialisé") logger.info(f" Entrée : {self.input_dir}") logger.info(f" Sortie : {self.output_dir}") if len(self.resolutions) > 1: logger.info(f" Résolutions : {', '.join(f'{r}m/px' for r in self.resolutions)}") else: logger.info(f" Résolution : {self.resolution}m/px") logger.info(f" Workers : {workers}") logger.info(f" Force : {'OUI' if self.force else 'non (skip existing)'}") logger.info(f" Classification sol : {self.ground_method}") logger.info(f" Force classif.: {'OUI' if self.force_classify else 'non'}") logger.info(f" Keep TIFF : {'OUI' if self.keep_tif else 'non'}") if self.edge_buffer > 0: logger.info(f" Raccord bords: {self.edge_buffer:g} m (points sol des tuiles voisines)") logger.info(f" Qualité {self.output_format.upper()}: {self.quality if self.quality < 100 else 'lossless'}") if only_viz: logger.info(f" Visualisations: uniquement {', '.join(only_viz)}") elif skip_viz: logger.info(f" Visualisations: tout sauf {', '.join(skip_viz)}") logger.info(f" Visualisations: {len(self.viz_steps)}/{len(VIZ_STEPS)}") def find_laz_files(self): """Find all LAZ/LAS files in input directory, triés du nord au sud. Les tuiles LHD_FXX_{col}_{row} sont ordonnées par ligne décroissante (row = nord en km) puis colonne croissante : les workers prennent les fichiers dans l'ordre de soumission, la carte se remplit ainsi du nord vers le sud lors des passes globales. Les fichiers hors pattern LHD restent triés par nom, en fin de liste. """ from .index import parse_basename_coords files = list(self.input_dir.glob("*.laz")) + list(self.input_dir.glob("*.las")) def _north_key(f): coords = parse_basename_coords(_file_basename(f)) if coords: col, row = coords return (0, -row, col, f.name) return (1, 0, 0, f.name) files.sort(key=_north_key) logger.info(f"{len(files)} fichier(s) LiDAR trouvé(s) — triés du nord au sud") for f in files: logger.debug(f" {f.name}") return files def check_tools(self): """Check that required external tools are available.""" for name, cmd in [('pdal', 'pdal --version'), ('gdal', 'gdalinfo --version')]: try: result = subprocess.run(cmd.split(), capture_output=True, check=True, text=True) version = result.stdout.strip().split('\n')[0] logger.info(f" ✓ {name}: {version}") except (subprocess.CalledProcessError, FileNotFoundError): logger.error(f" ✗ {name} non disponible") return False return True @staticmethod def _expected_output_path(name, basename, file_vis_dir, output_format='avif'): """Return the expected output filename for a visualization step.""" ext = 'avif' if output_format == 'avif' else 'webp' if name in LOSSLESS_GRAY_KEYWORDS: ext = 'webp' # aplats de niveaux : WebP sans perte (rendering.tif_to_crop) if name == 'pos_open': return file_vis_dir / f"{basename}_positive_openness.{ext}" elif name == 'neg_open': return file_vis_dir / f"{basename}_negative_openness.{ext}" elif name == 'hillshade': return file_vis_dir / f"{basename}_hillshade_multi.{ext}" else: return file_vis_dir / f"{basename}_{name}.{ext}" def generate_all_visualizations(self, dtm_file, basename, resolution=None, vis_dir=None, force_images=None): """Generate all archaeological visualizations for one DTM file. Optimisation: SharedDEM is only computed if at least one visualization needs to be generated. When all WebP outputs exist, SharedDEM is skipped entirely (saves ~2min per file on re-runs). """ if resolution is None: resolution = self.resolution logger.info(" Génération visualisations:") # Use provided vis_dir (for multi-resolution subdirectories) or default file_vis_dir = vis_dir if vis_dir else (self.vis_dir / basename) file_vis_dir.mkdir(exist_ok=True) total = len(self.viz_steps) # Phase 1: determine which visualizations need generation force_viz = self.force if force_images is None else force_images needs_generation = {} # name -> True/False for name, func in self.viz_steps: if force_viz: needs_generation[name] = True else: expected_webp = self._expected_output_path(name, basename, file_vis_dir, self.output_format) needs_generation[name] = not expected_webp.exists() to_generate = [n for n, needed in needs_generation.items() if needed] needs_shared = any(name not in ('ortho', 'topo') for name in to_generate) if not to_generate: logger.info(" Toutes les visualisations déjà existantes — ignorées") # Still need to return results dict for PDF check vis_results = {} for name, func in self.viz_steps: vis_results[name] = self._expected_output_path(name, basename, file_vis_dir, self.output_format) return vis_results # Phase 2: compute SharedDEM only if needed shared = None if needs_shared: logger.info(" Pré-calcul données partagées (gradient, LRM)...") t_shared = time.time() shared = SharedDEM(dtm_file, resolution) logger.info(f" ✓ Données partagées prêtes ({time.time()-t_shared:.1f}s)") # Phase 3: generate visualizations vis_results = {} for idx, (name, func) in enumerate(self.viz_steps, 1): if not needs_generation[name]: logger.info(f" [{idx}/{total}] {name}: déjà existant, ignoré") self._report(basename, "viz", "skip", name, res=resolution) vis_results[name] = self._expected_output_path(name, basename, file_vis_dir, self.output_format) continue # When regenerating, delete existing TIF to ensure clean regeneration if force_viz: for tif in file_vis_dir.glob(f"{basename}_{name}.tif"): tif.unlink(missing_ok=True) if name == 'pos_open': for tif in file_vis_dir.glob(f"{basename}_positive_openness.tif"): tif.unlink(missing_ok=True) elif name == 'neg_open': for tif in file_vis_dir.glob(f"{basename}_negative_openness.tif"): tif.unlink(missing_ok=True) elif name == 'hillshade': for tif in file_vis_dir.glob(f"{basename}_hillshade_multi.tif"): tif.unlink(missing_ok=True) logger.info(f" [{idx}/{total}] {name}...") self._report(basename, "viz", "start", name, res=resolution) t0 = time.time() try: # IGN overlays don't use SharedDEM (they download external data) # Non-IGN visualizations use safe_gpu_call for GPU→CPU fallback if name in ('ortho', 'topo'): result = func(dtm_file, basename, file_vis_dir, resolution) else: result = safe_gpu_call(func, dtm_file, basename, file_vis_dir, resolution, shared=shared) vis_results[name] = result elapsed = time.time() - t0 if result: logger.info(f" [{idx}/{total}] ✓ {name} ({elapsed:.1f}s)") self._report(basename, "viz", "ok", name, res=resolution) else: logger.warning(f" [{idx}/{total}] ✗ {name} — no output ({elapsed:.1f}s)") self._report(basename, "viz", "fail", name, res=resolution) except Exception as e: vis_results[name] = None logger.error(f" [{idx}/{total}] ✗ {name}: {e}", exc_info=True) self._report(basename, "viz", "fail", name, res=resolution) # Free GPU memory between visualizations to prevent OOM gpu_cleanup() # Convert to output format (only newly generated TIFs, not skipped ones) fmt_label = self.output_format.upper() logger.info(f" Conversion images {fmt_label}:") for name, tif_file in vis_results.items(): if tif_file and isinstance(tif_file, Path) and tif_file.suffix == '.tif' and tif_file.exists(): img_file = tif_to_crop(tif_file, file_vis_dir, resolution, keep_tif=self.keep_tif, quality=self.quality, output_format=self.output_format, subtiles_dir=self.output_dir) if img_file: logger.info(f" ✓ {img_file.name}") # Clean up remaining TIF files unless --keep-tif if not self.keep_tif: for tif in file_vis_dir.glob("*.tif"): tif.unlink(missing_ok=True) return vis_results @staticmethod def _res_suffix(resolution): """Return suffix for additional resolutions (empty string for primary).""" if resolution == 0.5: return "" # Default resolution — no suffix res_str = f"{resolution}".replace('.', 'p') return f"_r{res_str}" def _dtm_method_path(self, basename, res_suffix): """Sidecar path storing which ground method produced a DTM.""" return self.dtm_dir / f"{basename}_dtm{res_suffix}_method.txt" def _dtm_method_name(self, basename, res_suffix): """Read the recorded ground method for a DTM, or None if unknown.""" p = self._dtm_method_path(basename, res_suffix) if p.exists(): try: return p.read_text(encoding="utf-8").strip() or None except Exception: return None return None def _effective_ground_method(self): """Méthode effective pour le suivi de cache, classes IGN incluses. La méthode 'ign' est étiquetée avec les classes choisies (ex. 'ign_1_2') pour qu'un changement de --ign-classes déclenche la reclassification. """ if self.ground_method == 'ign': from .dtm import parse_ign_classes, ign_method_label return ign_method_label(parse_ign_classes(self.ign_classes)) return self.ground_method def _dtm_method_matches(self, basename, res_suffix): """True if the recorded ground method matches the requested one. A DTM without a recorded method is treated as matching so the existing cache is preserved; the method is adopted on its next reclassification. """ recorded = self._dtm_method_name(basename, res_suffix) return recorded is None or recorded == self._effective_ground_method() def _strip_align_matches(self, basename, res_suffix): """True si le sidecar de calage des faisceaux correspond à la config. Un DTM sans sidecar (antérieur au calage) est régénéré pour mesurer et consigner ses offsets ; un sidecar de version, de seuil ou de paramètres de gigue intra-faisceau différents aussi. Calage désactivé : tout DTM porteur d'un sidecar (donc calé) est régénéré non calé. """ sidecar = self.dtm_dir / f"{basename}_dtm{res_suffix}_stripalign.json" if not self.strip_align: return not sidecar.exists() if not sidecar.exists(): return False try: import json data = json.loads(sidecar.read_text(encoding="utf-8")) return (data.get("version") == STRIP_ALIGN_VERSION and abs(float(data.get("threshold", -1)) - STRIP_ALIGN_THRESHOLD) < 1e-9 and abs(float(data.get("jitter_bin", -1)) - STRIP_JITTER_BIN) < 1e-9 and int(data.get("jitter_smooth", -1)) == STRIP_JITTER_SMOOTH and int(data.get("line_window", -1)) == STRIP_LINE_WINDOW and abs(float(data.get("line_cell", -1)) - STRIP_LINE_CELL) < 1e-9 and data.get("line_model") == STRIP_LINE_MODEL) except Exception: return False def _write_dtm_method(self, basename, res_suffix): """Record the ground classification method used to build a DTM.""" try: self._dtm_method_path(basename, res_suffix).write_text( self._effective_ground_method(), encoding="utf-8") except Exception: pass def _edge_buffer_matches(self, dtm_path): """True si le tampon de raccord du DTM correspond à la config. Le tampon est inscrit dans le tag GeoTIFF LIDAR_EDGE_BUFFER (dtm.py). Un DTM sans tag (antérieur au raccord) compte comme tampon 0 : activer --edge-buffer régénère donc les DTM en cache, le désactiver régénère les DTM raccordés. """ from .dtm import read_dtm_edge_buffer return abs(read_dtm_edge_buffer(dtm_path) - self.edge_buffer) < 1e-6 def _gap_fill_matches(self, dtm_path): """True si le DTM a été comblé par la version courante (tag LIDAR_GAP_FILL, dtm.py) ; absent = ancien comblement à distance fixe.""" from .dtm import read_dtm_gap_fill, GAP_FILL_VERSION return read_dtm_gap_fill(dtm_path) == GAP_FILL_VERSION def _fetch_edge_neighbors(self, files): """Télécharge les dalles LAZ voisines manquantes (raccord des bords). La bande de raccord lit les 8 LAZ adjacentes de chaque tuile ; celles absentes sont téléchargées depuis le catalogue IGN avant de lancer les workers, dans le sous-dossier input/edge_neighbors/ : elles ne doivent PAS rejoindre input/ à plat, sinon les passes globales (« tout input/ ») les comptent comme des tuiles à rendre et la zone grandit d'un anneau à chaque relance. La liste est dédupliquée sur tout le lot : la couronne d'un bloc contigu ne coûte qu'un passage. Une dalle introuvable (zone non publiée) ou en échec laisse simplement la bande vide — le rendu continue (best-effort). """ if self.edge_buffer <= 0: return from .dtm import _tile_coords, _NEIGHBOR_OFFSETS, EDGE_NEIGHBORS_DIRNAME from .fetch_ign import fetch_tiles, tile_filename edge_dir = self.input_dir / EDGE_NEIGHBORS_DIRNAME wanted = set() for laz_file in files: coords = _tile_coords(Path(laz_file).name) if coords is None: continue col, row = coords for dcol, drow in _NEIGHBOR_OFFSETS: wanted.add((col + dcol, row + drow)) missing = sorted( cr for cr in wanted if not (self.input_dir / tile_filename(cr[0], cr[1])).exists() and not (edge_dir / tile_filename(cr[0], cr[1])).exists()) if not missing: return edge_dir.mkdir(parents=True, exist_ok=True) logger.info(f"Raccord des bords : {len(missing)} dalle(s) voisine(s) " f"absente(s) — téléchargement dans {EDGE_NEIGHBORS_DIRNAME}/ " f"(catalogue IGN)") t0 = time.time() from concurrent.futures import ThreadPoolExecutor with ThreadPoolExecutor(max_workers=4) as pool: fetched = [path for batch in pool.map( lambda spec: fetch_tiles(edge_dir, [spec]), missing) for path in batch] logger.info(f"Raccord des bords : {len(fetched)}/{len(missing)} dalle(s) " f"voisine(s) téléchargée(s) ({time.time() - t0:.0f} s)") def _report(self, basename, phase, state, detail=None, res=None): """Émet un événement de progression (file de génération, best-effort).""" report_event(self.output_dir, basename, phase, state, detail=detail, res=res) def process_file(self, laz_file): """Process a single LAZ file through the full pipeline. If self.resolutions has multiple entries, processes each resolution: - Primary resolution uses current naming (no suffix) - Additional resolutions use _r0p2 suffix in directories/filenames - Ground classification is done once and shared across resolutions """ basename = _file_basename(laz_file) _file_filter.basename = basename t_start = time.time() logger.info("=" * 60) logger.info(f"FICHIER : {basename}") logger.info("=" * 60) # Validate file integrity before any processing from .dtm import validate_laz if not validate_laz(laz_file): self._report(basename, "tile", "fail", "fichier LAZ invalide") return False # Step 1: Ground classification (shared across all resolutions) las_file = None t_classif = 0 dtm_rebuilt = False # The ground method is shared across resolutions. It is recorded per DTM # in a sidecar so that changing --ground-classification invalidates the # cache; otherwise the cached DTM would be reused and the new method # never applied (nothing would change visually). primary_suffix = self._res_suffix(self.resolutions[0]) method_matches = self._dtm_method_matches(basename, primary_suffix) for i, res in enumerate(self.resolutions): res_suffix = self._res_suffix(res) dtm_path = self.dtm_dir / f"{basename}_dtm{res_suffix}.tif" if dtm_path.exists() and not self.force_classify: if not self._strip_align_matches(basename, res_suffix): logger.info(f" DTM{res_suffix} sans calage de faisceaux conforme — régénération (offsets verticaux mesurés et appliqués)") dtm_path.unlink() elif not self._gap_fill_matches(dtm_path): logger.info(f" DTM{res_suffix} au comblement d'une version antérieure — régénération " f"(vides comblés dans l'enveloppe des points, rayon selon la densité)") dtm_path.unlink() elif not self._edge_buffer_matches(dtm_path): from .dtm import read_dtm_edge_buffer recorded = read_dtm_edge_buffer(dtm_path) logger.info(f" DTM{res_suffix} avec raccord de {recorded:g} m ≠ {self.edge_buffer:g} m " f"demandé — régénération (bande de bord des tuiles voisines)") dtm_path.unlink() elif method_matches: import rasterio try: with rasterio.open(dtm_path) as src: existing_res = abs(src.transform.a) if abs(existing_res - res) > 0.01: logger.info(f" DTM{res_suffix} existant à {existing_res}m/px — résolution demandée {res}m/px → régénération") dtm_path.unlink() else: if i == 0: logger.info(f"[1/5] Classification du sol — sautée (DTM existant)") logger.info(f"[2/5] Génération DTM {res}m/px — sautée (DTM existant)") self._report(basename, "classif", "skip", "DTM existant") else: logger.info(f" DTM {res}m/px déjà existant — ignoré") self._report(basename, "dtm", "skip", res=res) continue except Exception: logger.warning(f"Impossible de lire le DTM existant — régénération") dtm_path.unlink() else: logger.info(f" DTM{res_suffix} produit par {self._dtm_method_name(basename, primary_suffix) or '?'} ≠ {self._effective_ground_method()} → reclassification") dtm_path.unlink() # Need to classify/generate DTM for this resolution if las_file is None: # First time: do ground classification logger.info("[1/5] Classification du sol...") self._report(basename, "classif", "start") t1 = time.time() las_file = classify_ground(laz_file, self.temp_dir, method=self.ground_method, force=self.force_classify, ign_classes=self.ign_classes) t_classif = time.time() - t1 if not las_file: logger.error(f" ✗ Échec classification ({t_classif:.1f}s)") self._report(basename, "classif", "fail") self._report(basename, "tile", "fail", "échec classification") return False logger.info(f" ✓ Classification terminée ({t_classif:.1f}s)") self._report(basename, "classif", "ok") # Generate DTM at this resolution logger.info(f"{'[2/5]' if i == 0 else ' '} Génération DTM {res}m/px...") self._report(basename, "dtm", "start", res=res) t2 = time.time() # Classification IGN → mode pur : DTM = rasterisation brute des # classes choisies, sans plancher ni comblement pure_ign = "_ground_ign" in Path(las_file).name # Classes des voisines pour la bande de raccord : mêmes classes IGN # que le MNT (les voisines sont lues dans leur pré-classification # fournisseur, quelle que soit la méthode de la tuile centrale). from .dtm import parse_ign_classes neighbor_codes = parse_ign_classes(self.ign_classes) dtm_file = create_dtm_fast(las_file, basename, self.dtm_dir, res, force=self.force or self.force_classify, output_suffix=res_suffix, source_laz=laz_file, pure=pure_ign, strip_align=self.strip_align, edge_buffer=self.edge_buffer, neighbor_classes=neighbor_codes) t_dtm = time.time() - t2 if not dtm_file: logger.error(f" ✗ Échec DTM {res}m/px ({t_dtm:.1f}s)") self._report(basename, "dtm", "fail", res=res) if i == 0: self._report(basename, "tile", "fail", "échec DTM") return False # Primary resolution failure is fatal continue # Additional resolution failure is non-fatal logger.info(f" ✓ DTM {res}m/px terminé ({t_dtm:.1f}s)") self._report(basename, "dtm", "ok", res=res) dtm_rebuilt = True self._write_dtm_method(basename, res_suffix) if i == 0: # Sidecar qualité (densité sol, dates de vol) pour l'encart de # l'export PDF : calculé une fois, même si le DTM vient du cache. from .quality import ensure_quality ensure_quality(las_file, basename, self.output_dir) # Process each resolution: visualizations + PDF # Option de calcul (surcharge le défaut du module) appliquée ICI car les # workers (spawn) réimportent les modules à froid : c'est le seul endroit # qui s'exécute dans le processus qui fait le calcul. if self.openness_downsample is not None: from . import visualizations as _viz_mod _viz_mod.OPENNESS_DOWNSAMPLE = max(1, int(self.openness_downsample)) all_vis_results = {} for res in self.resolutions: res_suffix = self._res_suffix(res) dtm_path = self.dtm_dir / f"{basename}_dtm{res_suffix}.tif" if not dtm_path.exists(): logger.warning(f" DTM {res}m/px manquant — visualisations ignorées") continue import rasterio with rasterio.open(dtm_path) as src: actual_res = abs(src.transform.a) if len(self.resolutions) > 1: logger.info(f" --- Résolution {res}m/px ---") # For additional resolutions, use suffixed subdirectory if res_suffix: vis_dir = self.vis_dir / f"{basename}{res_suffix}" else: vis_dir = self.vis_dir / basename vis_dir.mkdir(exist_ok=True) self.generate_all_visualizations( dtm_path, basename, actual_res, vis_dir=vis_dir, force_images=self.force or self.force_classify or dtm_rebuilt) t_total = time.time() - t_start logger.info(f"✓ {basename} terminé en {t_total:.1f}s") self._report(basename, "tile", "ok", f"{t_total:.0f}s") _file_filter.basename = None return True def _rebuild_index_incremental(self): """Régénère la carte juste après une tuile terminée (mode incrémental). Réécrit index_tiles.json (et rafraîchit vignettes/sous-tuiles) pour que la carte affiche la tuile sans attendre la fin du run. Anti-rebond : 3 s minimum entre deux passes — les tuiles terminées pendant l'intervalle sont couvertes par la passe suivante ou par la passe finale. Les logs de build_index sont masqués pour ne pas noyer le journal du run. """ if self.no_index: return wait = 3.0 - (time.time() - self._last_index_rebuild) if wait > 0: # Anti-rebond : la dalle n'est pas oubliée, une passe différée la # couvre dès la fin de l'intervalle (sans attendre la dalle suivante). with self._index_lock: if self._index_timer is None: self._index_timer = threading.Timer(wait, self._rebuild_index_now) self._index_timer.daemon = True self._index_timer.start() return self._rebuild_index_now() def _rebuild_index_now(self): with self._index_lock: self._index_timer = None self._last_index_rebuild = time.time() try: from .index import build_index saved_level = logger.getEffectiveLevel() logger.setLevel(logging.WARNING) try: build_index(self.output_dir, self.output_format) finally: logger.setLevel(saved_level) except Exception as e: logger.debug(f"Rebuild incrémental de l'index ignoré : {e}") def process_all(self, files=None): """Process all LAZ files in input directory (or an explicit list).""" files = files if files is not None else self.find_laz_files() if not files: logger.error("Aucun fichier LAZ/LAS trouvé !") return logger.info("=" * 60) logger.info("PIPELINE ARCHÉOLOGIQUE LiDAR") logger.info("=" * 60) logger.info("Vérification des outils...") if not self.check_tools(): logger.error("Outils manquants — abandon") return 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) # Raccord des bords : pré-télécharger les voisines manquantes avant # les workers (no-op si le raccord est désactivé). self._fetch_edge_neighbors(files) if self.workers > 1 and len(files) > 1: n_gpus = num_gpus() or 1 if n_gpus > 1: logger.info(f"Traitement parallèle avec {self.workers} workers sur {n_gpus} GPUs...") else: logger.info(f"Traitement parallèle avec {self.workers} workers...") logger.info(f"Fichiers: {len(files)}") # Une place fixe par processus du pool (GPU ou CPU), prise à sa # création : au plus « VRAM libre / pic d'un worker » processus par # GPU. L'affectation par numéro de fichier (round-robin) laissait # 6 workers sur chaque GPU de 8 Go avec LIDAR_WORKERS=auto : OOM. active_ids = self.gpu_ids if self.gpu_ids else available_gpu_ids() slots = gpu_worker_slots(active_ids, self.workers) if active_ids: per_gpu = ", ".join(f"GPU {g} : {slots.count(g)}" for g in active_ids) n_cpu = slots.count(-1) logger.info(f"Répartition des workers : {per_gpu}" + (f", CPU : {n_cpu} (VRAM insuffisante)" if n_cpu else "")) slot_queue = multiprocessing.Queue() for slot in slots: slot_queue.put(slot) with ProcessPoolExecutor(max_workers=self.workers, initializer=_init_worker_slot, initargs=(slot_queue,)) as executor: 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.ign_classes, self.force_classify, self.keep_tif, self.quality, self.only_viz, self.skip_viz, self.output_format, None, self.openness_downsample, self.edge_buffer): laz_file for laz_file in files } done = 0 t_deadline = time.time() + 7200 try: # 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: f.cancel() executor.shutdown(wait=False, cancel_futures=True) logger.info("Travaux annulés.") return else: total = len(files) if self.workers == 1 and len(files) > 1: n_gpus = num_gpus() or 1 if n_gpus > 1: logger.info(f"Conseil : utilisez -w {n_gpus} pour exploiter tous les GPUs") for idx, laz_file in enumerate(files, 1): logger.info(f"--- Fichier {idx}/{total} ---") try: results[laz_file.name] = self.process_file(laz_file) if results[laz_file.name] and self.incremental_index: self._rebuild_index_incremental() except KeyboardInterrupt: logger.info("Interruption — arrêt immédiat.") return except Exception as e: logger.error(f"✗ Erreur traitement {laz_file.name}: {e}") logger.debug("Traceback:", exc_info=True) report_event(self.output_dir, _file_basename(laz_file), "tile", "fail", detail=str(e)) results[laz_file.name] = False # Summary t_pipeline_total = time.time() - t_pipeline_start success_count = sum(1 for v in results.values() if v) fail_count = sum(1 for v in results.values() if not v) logger.info("=" * 60) logger.info("RÉSUMÉ") logger.info("=" * 60) for name, ok in results.items(): status = "✓" if ok else "✗" logger.info(f" {status} {name}") logger.info("-" * 60) logger.info(f" Succès : {success_count}/{len(results)}") if fail_count: logger.info(f" Échecs : {fail_count}/{len(results)}") logger.info(f" Durée totale : {t_pipeline_total:.1f}s ({t_pipeline_total/60:.1f}min)") logger.info(f"\nRésultats dans: {self.output_dir}") logger.info(f" • DTM : {self.dtm_dir}") logger.info(f" • Visualisations: {self.vis_dir}") # Génère le catalogue des tuiles traitées (vignettes + inventaire) if not self.no_index: try: from .index import build_index index_path = build_index(self.output_dir, self.output_format) if index_path: logger.info(f" • Catalogue : {index_path}") except Exception as e: logger.warning(f"Index global non généré: {e}") # Clean up temporary files logger.info("Nettoyage des fichiers temporaires...") try: if self.temp_dir.exists(): shutil.rmtree(self.temp_dir) logger.info(" ✓ Fichiers temporaires supprimés") except Exception as e: logger.warning(f" Note: Impossible de supprimer les fichiers temporaires: {e}") def _init_worker_slot(slot_queue): """Initialiseur du pool : le processus prend sa place une fois pour toutes. Indice GPU → set_active_gpu ; -1 → CPU forcé (VRAM insuffisante) ; None ou file vide → choix laissé au worker (pas de GPU détecté). """ import queue from . import gpu try: slot = slot_queue.get_nowait() except queue.Empty: return if slot is None: return if slot < 0: gpu.force_cpu() else: gpu.set_active_gpu(slot) def _process_file_standalone(laz_file_str, input_dir, output_dir, resolution, force=False, ground_method='auto', ign_classes="sol", force_classify=False, keep_tif=False, quality=60, only_viz=None, skip_viz=None, output_format='avif', gpu_id=None, openness_downsample=None, edge_buffer=0.0): """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 CUDA_VISIBLE_DEVICES to balance load across GPUs. """ if gpu_id is not None and gpu_id >= 0: from .gpu import set_active_gpu set_active_gpu(gpu_id) # Configure logging in worker process (spawn doesn't inherit parent config) import logging import sys # Ensure UTF-8 output — spawn workers may default to ASCII if hasattr(sys.stdout, 'reconfigure'): sys.stdout.reconfigure(encoding='utf-8', errors='replace') if hasattr(sys.stderr, 'reconfigure'): sys.stderr.reconfigure(encoding='utf-8', errors='replace') worker_logger = logging.getLogger("lidar") if not worker_logger.handlers: handler = logging.StreamHandler(sys.stdout) handler.setFormatter(logging.Formatter("%(message)s")) worker_logger.setLevel(logging.INFO) worker_logger.addHandler(handler) worker_logger.addFilter(_file_filter) pipeline = LidarArchaeoPipeline(input_dir, output_dir, resolution=resolution, workers=1, force=force, ground_method=ground_method, ign_classes=ign_classes, force_classify=force_classify, keep_tif=keep_tif, quality=quality, only_viz=only_viz, skip_viz=skip_viz, output_format=output_format, openness_downsample=openness_downsample, edge_buffer=edge_buffer) basename = _file_basename(laz_file_str) pipeline.temp_dir = pipeline.output_dir / "temp" / basename pipeline.temp_dir.mkdir(exist_ok=True) laz_file = Path(laz_file_str) result = pipeline.process_file(laz_file) # Clean up per-file temp directory try: if pipeline.temp_dir.exists(): shutil.rmtree(pipeline.temp_dir) except Exception: pass return result