Files
Aislo/B04_wf1_Surface/B04_wf1_Surface_Engine.py
T
2026-07-18 12:25:14 +09:00

340 lines
14 KiB
Python

"""B04 지표면 분석 엔진 오케스트레이터.
원본 LAS를 구조화하고 지면 필터를 실행한 뒤 지표면 5종 모델을 빌드하는
동기 계산 파이프라인. 라우터에서 asyncio.to_thread로 호출한다.
"""
import json
import logging
import time
from collections.abc import Callable
from pathlib import Path
from typing import Any
import numpy as np
from B04_wf1_Surface.B04_wf1_Surface_Engine_Ground import (
build_ground_masks,
run_ground_filter,
summarize_masks,
)
from B04_wf1_Surface.B04_wf1_Surface_Engine_Pipeline import build_all_terrain_models
from B04_wf1_Surface.B04_wf1_Surface_Engine_Structurize import structurize_las
from common_util.common_util_atomic import atomic_write_npz
from common_util.common_util_json import atomic_write_json
from config.config_system import build_surface_model_config
logger = logging.getLogger(__name__)
# 진행 콜백 시그니처: (진행률 0~100, 현재 단계 키, 메시지)
ProgressCallback = Callable[[int, str, str], None]
GROUND_POINT_SAMPLE_LIMIT = 500_000
GROUND_POINT_CACHE_VERSION = 2
def _source_identity(las_path: Path) -> dict[str, Any]:
"""입력 LAS의 정체성(이름·크기·수정시각)으로 캐시 세대를 식별한다 (PLAN B-2)."""
stat = las_path.stat()
return {
"filename": las_path.name,
"size_bytes": int(stat.st_size),
"mtime": float(stat.st_mtime),
}
def _relative_to_project(project_root: Path, path: Path) -> str:
"""프로젝트 루트 기준 posix 상대 경로 문자열."""
return path.relative_to(project_root).as_posix()
def cache_ground_points(
structured_path: Path,
filter_key: str,
mask: np.ndarray | None = None,
) -> Path:
"""필터링된 지면 포인트 미리보기 캐시를 생성하고 경로를 반환한다.
mask 미전달 시 영구 저장된 mask_{filter}.npy를 우선 재사용한다 (PLAN A-2).
"""
with np.load(structured_path) as structured:
xyz = np.asarray(structured["xyz"], dtype=np.float32)
mask_path = structured_path.parent / f"mask_{filter_key}.npy"
if mask is None and mask_path.is_file():
stored_mask = np.load(mask_path)
if len(stored_mask) == len(xyz):
mask = np.asarray(stored_mask, dtype=bool)
if mask is None:
mask = build_ground_masks(structured, [filter_key])[filter_key]
np.save(mask_path, np.asarray(mask, dtype=bool))
ground_indexes = np.flatnonzero(mask)
ground_points = xyz[ground_indexes]
ground_point_count = int(len(ground_points))
if ground_point_count:
data_bounds = np.column_stack((ground_points.min(axis=0), ground_points.max(axis=0)))
else:
data_bounds = np.zeros((3, 2), dtype=np.float64)
if ground_point_count > GROUND_POINT_SAMPLE_LIMIT:
rng = np.random.default_rng(20260717)
sample_indexes = rng.choice(
ground_point_count, GROUND_POINT_SAMPLE_LIMIT, replace=False
)
ground_indexes = ground_indexes[sample_indexes]
points = ground_points[sample_indexes]
else:
points = ground_points
arrays = {
"xyz": points,
"bounds": np.asarray(structured["bounds"], dtype=np.float64),
"data_bounds": np.asarray(data_bounds, dtype=np.float64),
"cache_version": np.asarray(GROUND_POINT_CACHE_VERSION, dtype=np.int16),
"point_count": np.asarray(ground_point_count, dtype=np.int64),
"sampled_count": np.asarray(len(points), dtype=np.int64),
}
if "rgb" in structured:
rgb = np.asarray(structured["rgb"])
if rgb.ndim > 0 and len(rgb) == len(xyz):
arrays["rgb"] = rgb[ground_indexes]
cache_path = structured_path.parent / f"ground_points_{filter_key}.npz"
atomic_write_npz(cache_path, **arrays)
return cache_path
def run_surface_analysis(
project_root: Path,
las_path: Path,
*,
source_filters: list[str],
methods: list[str],
force: bool = False,
on_progress: ProgressCallback | None = None,
) -> dict[str, Any]:
"""구조화→필터→모델 빌드를 수행하고 산출 메타데이터를 반환한다.
반환 dict:
- processed: {processed_file_path, converted_file_path, point_count, bounds, statistics}
- ground_summary: 필터별 지면 포인트 요약
- manifest: 지표면 모델 파이프라인 manifest
- models: [{model_type, model_file_path, resolution_m, generation_params, layers}]
"""
def _report(percent: int, stage: str, message: str) -> None:
if on_progress is not None:
on_progress(percent, stage, message)
total_started = time.monotonic()
stage_root = project_root / "B04_wf1_Surface"
processed_dir = stage_root / "processed"
models_dir = stage_root / "models"
processed_dir.mkdir(parents=True, exist_ok=True)
models_dir.mkdir(parents=True, exist_ok=True)
# 0. 입력 세대 검증: LAS가 바뀌었으면 모든 캐시를 재계산한다 (PLAN B-2)
identity_path = processed_dir / "source_identity.json"
current_identity = _source_identity(las_path)
stored_identity: dict[str, Any] | None = None
if identity_path.is_file():
try:
stored_identity = json.loads(identity_path.read_text(encoding="utf-8"))
except (OSError, ValueError):
stored_identity = None
rebuild = force or stored_identity != current_identity
# 1. LAS 구조화 (structured.npz) — 캐시가 유효하면 스킵 (PLAN B-1)
structured_path = processed_dir / "structured.npz"
if rebuild or not structured_path.is_file():
_report(10, "structurize", "LAS 구조화 중")
step_started = time.monotonic()
structured_path = structurize_las(las_path, processed_dir)
atomic_write_json(identity_path, current_identity)
logger.info(
"B04 LAS 구조화 완료: %s (%.1fs)", las_path.name, time.monotonic() - step_started
)
else:
_report(10, "structurize", "구조화 캐시 재사용")
logger.info("B04 구조화 캐시 재사용: %s", structured_path.name)
with np.load(structured_path) as structured:
xyz = structured["xyz"]
bounds = structured["bounds"]
total_points = int(len(xyz))
stats = {
"min_z": float(bounds[2, 0]),
"max_z": float(bounds[2, 1]),
"mean_z": float(np.mean(xyz[:, 2])) if total_points else None,
}
bounds_dict = {
"x_min": float(bounds[0, 0]),
"x_max": float(bounds[0, 1]),
"y_min": float(bounds[1, 0]),
"y_max": float(bounds[1, 1]),
}
data = {"xyz": xyz, "bounds": bounds}
# 2. 지면 필터 실행 — mask_{filter}.npy 영구 캐시 우선 재사용 (PLAN A-1)
_report(40, "ground_filter", "지면 필터 적용 중")
masks: dict[str, np.ndarray] = {}
for filter_key in source_filters:
mask_path = processed_dir / f"mask_{filter_key}.npy"
mask: np.ndarray | None = None
if not rebuild and mask_path.is_file():
stored_mask = np.load(mask_path)
if len(stored_mask) == total_points:
mask = np.asarray(stored_mask, dtype=bool)
logger.info("B04 지면 필터 캐시 재사용: %s", mask_path.name)
if mask is None:
step_started = time.monotonic()
mask = np.asarray(run_ground_filter(filter_key, data), dtype=bool)
np.save(mask_path, mask)
logger.info(
"B04 지면 필터 계산 완료: %s (%.1fs)",
filter_key,
time.monotonic() - step_started,
)
masks[filter_key] = mask
ground_summary = summarize_masks(data, masks)
for filter_key, mask in masks.items():
cache_path = processed_dir / f"ground_points_{filter_key}.npz"
if not rebuild and cache_path.is_file():
with np.load(cache_path) as cached:
if (
"cache_version" in cached
and int(cached["cache_version"]) == GROUND_POINT_CACHE_VERSION
):
continue
cache_ground_points(structured_path, filter_key, mask)
# 3. 지표면 5종 모델 빌드
_report(70, "surface_model", "지표면 모델 생성 중")
config = build_surface_model_config()
config["source_filters"] = list(source_filters)
config["precompute"] = list(methods)
step_started = time.monotonic()
def _model_progress(percent: int, detail: str) -> None:
_report(70 + int(percent * 0.2), "surface_model", detail)
manifest = build_all_terrain_models(
data, masks, models_dir, config, force=rebuild, progress=_model_progress
)
logger.info(
"B04 지표면 모델 빌드 완료: status=%s (%.1fs)",
manifest.get("status"),
time.monotonic() - step_started,
)
# 3-2. VWorld 지도 및 국가 GIS 벡터 다운로드 (기존 산출물이 있으면 스킵)
_report(90, "download_maps", "VWorld 지도 및 GIS 벡터 데이터 다운로드 중")
try:
# 입력 LAS와 같은 폴더의 .prj를 우선 사용 (PLAN E-1: 전체 재귀 glob 제거)
prj_candidates = sorted(las_path.parent.glob("*.prj")) or sorted(
project_root.glob("B03_FileInput/**/*.prj")
)
prj_path = prj_candidates[0] if prj_candidates else project_root / "result.prj"
bounds_dict_for_download = {
"x": [float(bounds[0, 0]), float(bounds[0, 1])],
"y": [float(bounds[1, 0]), float(bounds[1, 1])],
"z": [float(bounds[2, 0]), float(bounds[2, 1])],
}
from B04_wf1_Surface.B04_wf1_Surface_Engine_GisVector import download_all_gis_vectors
from B04_wf1_Surface.B04_wf1_Surface_Engine_VWorld import download_vworld_satellite_map
# VWorld 지도 및 GIS 데이터 저장 위치는 B04_wf1_Surface/processed에 보관.
layers = [
{"layer": "Satellite", "ext": "jpeg"},
{"layer": "Hybrid", "ext": "png"},
{"layer": "white", "ext": "png"},
]
for item in layers:
meta_path = processed_dir / f"vworld_{item['layer'].lower()}_meta.json"
if not rebuild and meta_path.is_file():
continue
try:
step_started = time.monotonic()
download_vworld_satellite_map(
prj_path,
bounds_dict_for_download,
processed_dir,
layer_name=item["layer"],
ext=item["ext"],
)
logger.info(
"B04 VWorld %s 지도 다운로드 완료 (%.1fs)",
item["layer"],
time.monotonic() - step_started,
)
except Exception as exc:
logger.warning("B04 VWorld %s 지도 다운로드 실패: %s", item["layer"], exc)
if rebuild or not any(processed_dir.glob("*_bounds.geojson")):
try:
step_started = time.monotonic()
download_all_gis_vectors(prj_path, bounds_dict_for_download, processed_dir)
logger.info(
"B04 국가 GIS 벡터 다운로드 완료 (%.1fs)", time.monotonic() - step_started
)
except Exception as exc:
logger.warning("B04 국가 GIS 벡터 다운로드 실패: %s", exc)
except Exception as exc:
logger.warning("B04 지도·GIS 다운로드 단계 실패: %s", exc)
_report(95, "saving", "결과 저장 중")
processed = {
"processed_file_path": _relative_to_project(project_root, structured_path),
"converted_file_path": None,
"point_count": total_points,
"bounds": bounds_dict,
"statistics": stats,
}
# manifest에서 저장된 모델별 정보 추출 (필터별 대표 모델)
models: list[dict[str, Any]] = []
for filter_key, filter_entry in manifest.get("source_filters", {}).items():
for method, meta in filter_entry.get("methods", {}).items():
if meta.get("status") != "completed":
continue
model_file = meta.get("model_file")
model_path = (models_dir / model_file) if model_file else None
layers: list[dict[str, Any]] = []
if meta.get("preview_file"):
layers.append(
{
"layer_name": f"{method}_{filter_key}_preview",
"geometry_type": "MESH" if method != "meshfree" else "POINTCLOUD",
"file_path": _relative_to_project(
project_root, models_dir / meta["preview_file"]
),
"file_format": "glb" if method != "meshfree" else "ply",
}
)
models.append(
{
"model_type": method,
"source_filter": filter_key,
"representation": meta.get("representation"),
"model_file_path": _relative_to_project(project_root, model_path)
if model_path
else None,
"resolution_m": meta.get("grid_resolution_meters"),
"generation_params": {
"source_filter": filter_key,
"representation": meta.get("representation"),
"footprint_area_m2": meta.get("footprint_area_m2"),
},
"layers": layers,
}
)
logger.info(
"B04 WF1 분석 완료: 모델 %d개, 총 %.1fs", len(models), time.monotonic() - total_started
)
return {
"processed": processed,
"ground_summary": ground_summary,
"manifest": manifest,
"models": models,
}