Files
lidar_rendu/lidar_pipeline/pipeline.py

867 lines
41 KiB
Python

"""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 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)
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,
)
from .gpu import gpu_cleanup, num_gpus, available_gpu_ids, restrict_gpus, safe_gpu_call
from .ign import generate_ign_overlay
from .rendering import tif_to_crop
# 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),
('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
self.incremental_index = incremental_index
self.strip_align = strip_align
self.openness_downsample = openness_downsample
self.edge_buffer = float(edge_buffer)
self._last_index_rebuild = 0.0
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:
self.viz_steps = VIZ_STEPS
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 == '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)
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)
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 _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 de input/ sont téléchargées depuis le catalogue IGN avant de
lancer les workers, pour que la bande soit remplie jusqu'au bord de la
zone même pour des tuiles jamais rendues. 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
from .fetch_ign import fetch_tiles, tile_filename
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())
if not missing:
return
logger.info(f"Raccord des bords : {len(missing)} dalle(s) voisine(s) "
f"absente(s) de input/ — téléchargement depuis le 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(self.input_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._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)
# 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.html + index_tiles.json pour que la webapp 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
now = time.time()
if now - self._last_index_rebuild < 3.0:
return
self._last_index_rebuild = now
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)}")
with ProcessPoolExecutor(max_workers=self.workers) as executor:
# Round-robin assign each file to a real GPU host index
active_ids = self.gpu_ids if self.gpu_ids else available_gpu_ids()
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, active_ids[file_idx % len(active_ids)] if active_ids else None, self.openness_downsample, self.edge_buffer): laz_file
for file_idx, laz_file in enumerate(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 la carte globale interactive des tuiles traitées
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" • Carte globale : {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 _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