"""Pipeline orchestration for LiDAR archaeological analysis. LidarArchaeoPipeline coordinates the full processing chain: 1. Ground classification (IGN pre-classification by default; SMRF/CSF via PDAL) 2. DTM generation 3. Visualization generation (17 available products; a default run only produces the map layers, PANEL_VIZ in index.py) 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): """Turn the -w option into an effective worker count. Resolved when each run starts (not at server startup, nor per tile: the process pool lives for the duration of the run). "auto" = logical cores - 2 (one for the OS/server, one for the I/O and indexing phases), clamped to [2, 16] — each worker processes one tile and may spawn a streaming PDAL process, so the count stays reasonable even on a large 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='IGN Aerial Photograph', legend_label='Orthophoto\nAerial image', description='IGN aerial photograph (orthophoto)', out_suffix='ortho')), ('topo', lambda d, b, v, r: generate_ign_overlay( d, b, v, r, layer='GEOGRAPHICALGRIDSYSTEMS.PLANIGNV2', title='IGN Topographic Map', legend_label='IGN map\nTopographic map', description='IGN topographic map (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 # Inventory updated after every tile by default: the map (and its # pyramid maintenance) follows the ongoing render, whatever the # launcher (map, compose, run.sh). --no-index disables it. 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"Directory not found: {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"Unknown visualizations: {', '.join(invalid)}. Available: {', '.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"Unknown visualizations: {', '.join(invalid)}. Available: {', '.join(all_viz_names)}") self.viz_steps = [(n, f) for n, f in VIZ_STEPS if n not in skip_viz] else: # Without --only/--skip: only the layers shown on the map # (PANEL_VIZ, index.py) are produced. 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 initialized") logger.info(f" Input : {self.input_dir}") logger.info(f" Output : {self.output_dir}") if len(self.resolutions) > 1: logger.info(f" Resolutions : {', '.join(f'{r} m/px' for r in self.resolutions)}") else: logger.info(f" Resolution : {self.resolution} m/px") logger.info(f" Workers : {workers}") logger.info(f" Force : {'YES' if self.force else 'no (skip existing)'}") logger.info(f" Ground classification: {self.ground_method}") logger.info(f" Force classif.: {'YES' if self.force_classify else 'no'}") logger.info(f" Keep TIFF : {'YES' if self.keep_tif else 'no'}") if self.edge_buffer > 0: logger.info(f" Edge buffer : {self.edge_buffer:g} m (ground points from neighboring tiles)") logger.info(f" {self.output_format.upper()} quality: {self.quality if self.quality < 100 else 'lossless'}") if only_viz: logger.info(f" Visualizations: only {', '.join(only_viz)}") elif skip_viz: logger.info(f" Visualizations: all except {', '.join(skip_viz)}") logger.info(f" Visualizations: {len(self.viz_steps)}/{len(VIZ_STEPS)}") def find_laz_files(self): """Find all LAZ/LAS files in the input directory, sorted north to south. LHD_FXX_{col}_{row} tiles are ordered by decreasing row (row = northing in km), then increasing column: workers pick files up in submission order, so the map fills from north to south during full passes. Files that do not match the LHD pattern are sorted by name at the end of the list. """ 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)} LiDAR file(s) found — sorted north to south") 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} not available") 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' # flat level classes: lossless WebP (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. Optimization: SharedDEM is only computed if at least one visualization needs to be generated. When all output images (AVIF/WebP) exist, SharedDEM is skipped entirely (saves time on re-runs). """ if resolution is None: resolution = self.resolution logger.info(" Generating visualizations:") # 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_img = self._expected_output_path(name, basename, file_vis_dir, self.output_format) needs_generation[name] = not expected_img.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(" All visualizations already exist — skipped") # Still return the results dict (expected output paths) 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(" Precomputing shared data (gradient, LRM)...") t_shared = time.time() shared = SharedDEM(dtm_file, resolution) logger.info(f" ✓ Shared data ready ({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}: already exists, skipped") 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" Converting images to {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): """Effective method for cache tracking, IGN classes included. The 'ign' method is labeled with the chosen classes (e.g. 'ign_1_2') so that changing --ign-classes triggers reclassification (the default class list, ground only, keeps the plain 'ign' label). """ 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 if the strip-alignment sidecar matches the current config. A DTM without a sidecar (older than strip alignment) is regenerated to measure and record its offsets; so is one whose sidecar has a different version, threshold, intra-strip jitter parameters or scan-line parameters (window, cell, model). With alignment disabled, any DTM carrying a sidecar (hence aligned) is regenerated unaligned. """ 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 if the DTM edge buffer matches the current config. The buffer is stored in the LIDAR_EDGE_BUFFER GeoTIFF tag (dtm.py). A DTM without the tag (older than edge stitching) counts as buffer 0: enabling --edge-buffer therefore regenerates cached DTMs, and disabling it regenerates buffered ones. """ 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 if the DTM gaps were filled by the current version (the LIDAR_GAP_FILL tag, dtm.py); missing = old fixed-distance filling.""" 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): """Download missing neighboring LAZ tiles (edge stitching). The edge band reads the 8 LAZ tiles adjacent to each tile; missing ones are downloaded from the IGN catalog before the workers start, into the input/edge_neighbors/ subdirectory: they must NOT land flat in input/, otherwise full passes ("all of input/") would count them as tiles to render and the area would grow by one ring on every rerun. The list is deduplicated across the whole batch: the ring around a contiguous block costs a single pass. A tile that cannot be found (unpublished area) or fails to download simply leaves the band empty — rendering continues (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 self._cleanup_edge_neighbor_duplicates(edge_dir) 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"Edge stitching: {len(missing)} missing neighboring " f"tile(s) — downloading into {EDGE_NEIGHBORS_DIRNAME}/ " f"(IGN catalog)") 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"Edge stitching: {len(fetched)}/{len(missing)} neighboring " f"tile(s) downloaded ({time.time() - t0:.0f} s)") def _cleanup_edge_neighbor_duplicates(self, edge_dir): """Remove from edge_neighbors/ the tiles already present in input/. A tile downloaded as a neighbor before being requested as a main tile (or the reverse) can end up duplicated in both locations (before fetch_tiles promotes the new neighbors); the copy in input/ is authoritative (_neighbor_laz_files searches input/ before edge_neighbors/), so the duplicate is useless and can weigh several hundred MB. ".part" files (download in progress) are never touched, and a duplicate of a different size is kept (the copy in input/ could be truncated: never delete what might be the only sound copy). """ from .dtm import EDGE_NEIGHBORS_DIRNAME if not edge_dir.is_dir(): return removed = 0 freed = 0 for path in edge_dir.iterdir(): if not path.is_file() or path.suffix == ".part": continue twin = self.input_dir / path.name if not twin.exists(): continue try: size = path.stat().st_size if twin.stat().st_size != size: logger.warning(f"Edge stitching: {path.name} present in input/ " f"and {EDGE_NEIGHBORS_DIRNAME}/ with different " f"sizes — duplicate kept") continue path.unlink() freed += size removed += 1 except OSError: pass if removed: logger.info(f"Edge stitching: {removed} duplicate(s) removed from " f"{EDGE_NEIGHBORS_DIRNAME}/ (already in input/, " f"{freed / 1e6:.0f} MB freed)") def _report(self, basename, phase, state, detail=None, res=None): """Emit a progress event (generation queue, 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"FILE: {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", "invalid LAZ file") 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} lacks up-to-date strip alignment — regenerating (vertical offsets measured and applied)") dtm_path.unlink() elif not self._gap_fill_matches(dtm_path): logger.info(f" DTM{res_suffix} gap-filled by an older version — regenerating " f"(gaps filled within the point envelope, radius based on density)") 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} has a {recorded:g} m edge buffer ≠ {self.edge_buffer:g} m " f"requested — regenerating (edge band from neighboring tiles)") 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" Existing DTM{res_suffix} at {existing_res} m/px — requested resolution {res} m/px → regenerating") dtm_path.unlink() else: if i == 0: logger.info("[1/5] Ground classification — skipped (existing DTM)") logger.info(f"[2/5] DTM generation {res} m/px — skipped (existing DTM)") self._report(basename, "classif", "skip", "existing DTM") # Quality sidecar (ground density, flight dates) for # the PDF export inset: DTM reused from the cache — # ground classification was skipped, so measure # directly on the input LAZ with the run's IGN # classes (ensure_quality recomputes nothing if the # sidecar already exists: steady-state cost = one # JSON read). from .dtm import parse_ign_classes from .quality import ensure_quality ensure_quality(laz_file, basename, self.output_dir, codes=tuple(parse_ign_classes(self.ign_classes))) else: logger.info(f" DTM {res} m/px already exists — skipped") self._report(basename, "dtm", "skip", res=res) continue except Exception: logger.warning("Cannot read the existing DTM — regenerating") dtm_path.unlink() else: logger.info(f" DTM{res_suffix} produced by {self._dtm_method_name(basename, primary_suffix) or '?'} ≠ {self._effective_ground_method()} → reclassifying") 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] Ground classification...") 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" ✗ Classification failed ({t_classif:.1f}s)") self._report(basename, "classif", "fail") self._report(basename, "tile", "fail", "classification failed") return False logger.info(f" ✓ Classification done ({t_classif:.1f}s)") self._report(basename, "classif", "ok") # Generate DTM at this resolution logger.info(f"{'[2/5]' if i == 0 else ' '} DTM generation {res} m/px...") self._report(basename, "dtm", "start", res=res) t2 = time.time() # IGN classification → "pure" flag. create_dtm_fast now ignores it # (kept for compatibility): small gaps are filled and strip # alignment applies whatever the classification method. pure_ign = "_ground_ign" in Path(las_file).name # Neighbor classes for the edge band: same IGN classes as the DTM # (neighbors are read in their provider pre-classification, # whatever the method used for the central tile). 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" ✗ DTM {res} m/px failed ({t_dtm:.1f}s)") self._report(basename, "dtm", "fail", res=res) if i == 0: self._report(basename, "tile", "fail", "DTM failed") return False # Primary resolution failure is fatal continue # Additional resolution failure is non-fatal logger.info(f" ✓ DTM {res} m/px done ({t_dtm:.1f}s)") self._report(basename, "dtm", "ok", res=res) dtm_rebuilt = True self._write_dtm_method(basename, res_suffix) if i == 0: # Quality sidecar (ground density, flight dates) for the PDF # export inset: freshly generated DTM — measured on the already # filtered ground LAS (the cached-DTM case is handled above, # before the `continue`). from .quality import ensure_quality ensure_quality(las_file, basename, self.output_dir) # Process each resolution: visualizations + image conversion # Computation option (overrides the module default) applied HERE because # workers (spawn) re-import modules from scratch: this is the only place # that runs inside the process doing the computation. 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 missing — visualizations skipped") continue import rasterio with rasterio.open(dtm_path) as src: actual_res = abs(src.transform.a) if len(self.resolutions) > 1: logger.info(f" --- Resolution {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} done in {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): """Rebuild the map inventory right after a tile finishes (incremental mode). Rewrites index_tiles.json (and refreshes thumbnails/subtiles) so the map shows the tile without waiting for the end of the run. Debounce: at least 3 s between two passes — tiles finished during the interval are covered by a deferred pass scheduled at the end of the interval (or by the final pass). build_index logs are silenced so they do not flood the run log. """ if self.no_index: return wait = 3.0 - (time.time() - self._last_index_rebuild) if wait > 0: # Debounce: the tile is not forgotten, a deferred pass covers it as # soon as the interval ends (without waiting for the next tile). 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"Incremental index rebuild skipped: {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("No LAZ/LAS file found!") return logger.info("=" * 60) logger.info("LiDAR ARCHAEOLOGICAL PIPELINE") logger.info("=" * 60) logger.info("Checking tools...") if not self.check_tools(): logger.error("Missing tools — aborting") 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) # Edge stitching: pre-download missing neighbors before starting the # workers (no-op when edge stitching is disabled). 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"Parallel processing with {self.workers} workers on {n_gpus} GPUs...") else: logger.info(f"Parallel processing with {self.workers} workers...") logger.info(f"Files: {len(files)}") # One fixed slot per pool process (GPU or CPU), taken when the # process starts: at most "free VRAM / peak per worker" processes # per GPU. Assignment by file number (round-robin) used to put # 6 workers on each 8 GB GPU with 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"Worker distribution: {per_gpu}" + (f", CPU: {n_cpu} (insufficient VRAM)" 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=: without it, as_completed blocks between two # completions and the 2 h deadline is never evaluated if # no worker returns (run stuck forever). 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(" 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("Timeout exceeded (2 h) — cancelling remaining workers") for f in future_to_file: f.cancel() except KeyboardInterrupt: logger.info("Interrupted — cancelling running jobs...") for f in future_to_file: f.cancel() executor.shutdown(wait=False, cancel_futures=True) logger.info("Jobs cancelled.") 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"Tip: use -w {n_gpus} to take advantage of all GPUs") for idx, laz_file in enumerate(files, 1): logger.info(f"--- File {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("Interrupted — stopping immediately.") return except Exception as e: logger.error(f"✗ Error processing {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("SUMMARY") logger.info("=" * 60) for name, ok in results.items(): status = "✓" if ok else "✗" logger.info(f" {status} {name}") logger.info("-" * 60) logger.info(f" Succeeded: {success_count}/{len(results)}") if fail_count: logger.info(f" Failed: {fail_count}/{len(results)}") logger.info(f" Total time: {t_pipeline_total:.1f}s ({t_pipeline_total/60:.1f} min)") logger.info(f"\nResults in: {self.output_dir}") logger.info(f" • DTM : {self.dtm_dir}") logger.info(f" • Visualizations: {self.vis_dir}") # Build the catalog of processed tiles (thumbnails + inventory) 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" • Catalog: {index_path}") except Exception as e: logger.warning(f"Global index not generated: {e}") # Clean up temporary files logger.info("Cleaning up temporary files...") try: if self.temp_dir.exists(): shutil.rmtree(self.temp_dir) logger.info(" ✓ Temporary files removed") except Exception as e: logger.warning(f" Note: could not remove temporary files: {e}") def _init_worker_slot(slot_queue): """Pool initializer: the process takes its slot once and for all. GPU index → set_active_gpu; -1 → forced CPU (insufficient VRAM); None or empty queue → left to the worker (no GPU detected). """ 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. The GPU is normally assigned once per pool process by _init_worker_slot (VRAM-bounded slots); gpu_id (process_all passes None) only forces a GPU for a direct call. """ 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