The pool of tile workers used to cancel every remaining tile after a
hardcoded 2-hour wall clock, silently truncating large batches (a
670-tile completion run lost its last 348 tiles that way). The timeout
now defaults to unlimited and can be capped per deployment with the
LIDAR_BATCH_TIMEOUT environment variable (seconds); the local worker
compose sets it to 6 hours.
💘 Generated with Crush
Assisted-by: Crush:glm-5.2
1034 lines
49 KiB
Python
1034 lines
49 KiB
Python
"""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 _batch_timeout_s():
|
|
"""Wall-clock budget for one batch, from LIDAR_BATCH_TIMEOUT (seconds).
|
|
|
|
Unset, empty or 0 means unlimited (default). A negative or non-numeric
|
|
value is ignored with a warning.
|
|
"""
|
|
raw = (os.environ.get("LIDAR_BATCH_TIMEOUT") or "").strip()
|
|
if not raw:
|
|
return 0.0
|
|
try:
|
|
return max(0.0, float(raw))
|
|
except ValueError:
|
|
logger.warning(f"LIDAR_BATCH_TIMEOUT={raw!r} is not a number — ignored (unlimited)")
|
|
return 0.0
|
|
|
|
|
|
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
|
|
budget_s = _batch_timeout_s()
|
|
deadline = time.time() + budget_s if budget_s > 0 else None
|
|
try:
|
|
# timeout= lets the deadline be checked between two
|
|
# completions; without it a stuck worker would block
|
|
# as_completed forever. No deadline (LIDAR_BATCH_TIMEOUT
|
|
# unset or 0) → wait indefinitely.
|
|
as_completed_kwargs = ({"timeout": max(1.0, deadline - time.time())}
|
|
if deadline is not None else {})
|
|
try:
|
|
for future in as_completed(future_to_file, **as_completed_kwargs):
|
|
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(f"Batch exceeded LIDAR_BATCH_TIMEOUT "
|
|
f"({budget_s:g} s) — 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 |